You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SpringBoot集成RabbitMQ Stream集群消费报错:预期server7:5557,实际server5:5555

SpringBoot客户端消费RabbitMQ Stream集群消息失败

报错日志

17:29:33.878 [main] DEBUG com.rabbitmq.stream.impl.Client - Trying to create stream connection to localhost:5555
17:29:33.899 [main] DEBUG com.rabbitmq.stream.impl.Client - Connection tuned with max frame size 1048576 and heartbeat 60
17:29:33.900 [main] DEBUG com.rabbitmq.stream.impl.Utils - Expected client server7:5557, got server5:5555: failure
17:29:33.900 [main] DEBUG com.rabbitmq.stream.impl.Client - Closing client
17:29:33.901 [main] DEBUG com.rabbitmq.stream.impl.Client - Closing Netty channel

消费者代码

import java.util.concurrent.atomic.AtomicInteger;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.ApplicationContext;

import com.rabbitmq.stream.Address;
import com.rabbitmq.stream.Consumer;
import com.rabbitmq.stream.Environment;
import com.rabbitmq.stream.OffsetSpecification;

@SpringBootApplication
public class StreamConsumerApplication {

static void log(String format, Object... arguments) {
    System.out.println(String.format(format, arguments));
}

public static void main(String[] args) throws Exception {
    log("Connecting...");
    Address entryPoint = new Address("localhost", 5555);
    try (Environment environment = Environment.builder().host(entryPoint.host()).port(entryPoint.port())
            .username("rabbit_admin").password(".123-321.").addressResolver(address -> entryPoint).build()) {

        log("Connected");

        AtomicInteger messageConsumed = new AtomicInteger(0);
        long start = System.currentTimeMillis();
        log("Start consumer...");
        Consumer consumer = environment.consumerBuilder().stream("finance.eletronics")
                .offset(OffsetSpecification.offset(0))
                .messageHandler((context, message) -> {
                    messageConsumed.incrementAndGet();
                    System.out.println("Received: "+new String(message.getBodyAsBinary()));
                })
                .build();
        Utils.waitAtMost(60, () -> messageConsumed.get() >= 1_000_000);
        log("Consumed %,d messages in %s ms", messageConsumed.get(), (System.currentTimeMillis() - start));
        log("Closing environment...");
    }
    log("Environment closed");
}

}

生产者代码

public void send() throws InterruptedException {
    log("Connecting...");
    Address entryPoint = new Address("localhost", 5555);
    try (Environment environment = Environment.builder().host(entryPoint.host()).port(entryPoint.port())
            .username("rabbit_admin").password(".123-321.").addressResolver(address -> entryPoint).build()) {

        log("Connected");

        log("Creating stream...");
        environment.streamCreator().stream("finance.eletronics").create();
        log("Stream created");

        log("Creating producer...");
        Producer producer = environment.producerBuilder().stream("finance.eletronics").build();
        log("Producer created");

        long start = System.currentTimeMillis();
        int messageCount = 3;
        CountDownLatch confirmLatch = new CountDownLatch(messageCount);
        log("Sending %,d messages", messageCount);
        IntStream.range(0, messageCount).forEach(i -> {
            Message message = producer.messageBuilder().properties().creationTime(System.currentTimeMillis())
                    .messageId(i).messageBuilder().addData("hello world".getBytes(StandardCharsets.UTF_8)).build();
            producer.send(message, confirmationStatus -> confirmLatch.countDown());
        });
        log("Messages sent, waiting for confirmation...");
        boolean done = confirmLatch.await(1, TimeUnit.MINUTES);
        log("All messages confirmed? %s (%d ms)", done ? "yes" : "no", (System.currentTimeMillis() - start));
        log("Closing environment...");
    }
    log("Environment closed");
}

RabbitMQ节点配置(以server2为例)

loopback_users.guest = true
stream.listeners.tcp.1 = 5552
stream.advertised_host = server2
stream.advertised_port = 5552
management.tcp.port = 15672
prometheus.tcp.port = 15692
listeners.tcp.default = 5672

Docker-Compose配置片段(server2)

version: "3.2"
services:
server2:
image: rabbitmq:3.10.9-management
hostname: server2
container_name: 'server2'
ports:
- "5672:5672"
- "15672:15672"
- "5552:5552"
- "15692:15692"
volumes:
- ./rabbitmq.conf:/etc/rabbitmq/rabbitmq.conf
- type: bind
  source: $PWD/.erlang.cookie
  target: /var/lib/rabbitmq/.erlang.cookie
networks:
- rabbitmq_s_net
environment:
- RABBITMQ_DEFAULT_USER=rabbit_admin
- RABBITMQ_DEFAULT_PASS=.123-321.
- RABBITMQ_CONFIG_FILES=/etc/rabbitmq/rabbitmq.conf
- RABBITMQ_ADVANCED_CONFIG_FILE=/etc/rabbitmq/advanced.config
- RABBITMQ_NODENAME=rabbit@server2
networks:
  rabbitmq_s_net:
    external: true
    name: sales_net

集群状态与Stream信息

RabbitMQ集群状态截图
RabbitMQ Stream信息截图


问题原因分析

报错信息Expected client server7:5557, got server5:5555本质是地址匹配失败:

  1. 客户端连接入口节点localhost:5555后,集群返回了实际存储Stream数据的节点地址(server7:5557)
  2. 但代码中硬编码的addressResolver(address -> entryPoint)强制将所有节点地址解析为localhost:5555,导致客户端尝试连接server7时被错误转向到server5的地址,触发校验失败
  3. 生产者能正常工作是因为仅需连接入口节点完成消息发布,而消费者需要直接连接存储Stream数据的节点进行消费

解决方案

方案1:移除硬编码的地址解析器,配置本地hostname映射

修改客户端Environment构建逻辑,去掉强制地址解析的代码,让客户端自动处理集群返回的节点地址:

// 修改后的Environment构建代码
try (Environment environment = Environment.builder()
        .host(entryPoint.host())
        .port(entryPoint.port())
        .username("rabbit_admin")
        .password(".123-321.")
        // 移除硬编码的addressResolver配置
        .build()) {
    // 剩余业务代码不变
}

同时在本地hosts文件中添加节点hostname与宿主机IP的映射(替换为实际宿主机IP):

192.168.1.100 server2
192.168.1.100 server5
192.168.1.100 server7

方案2:统一配置节点的advertised地址为宿主机IP

修改每个RabbitMQ节点的rabbitmq.conf,将stream.advertised_host改为Docker宿主机的IP,而非容器hostname:

# 以server2为例,替换为实际宿主机IP
stream.advertised_host = 192.168.1.100
stream.advertised_port = 5552

确保Docker-Compose中暴露了对应节点的Stream端口(如server5映射5555,server7映射5557),客户端无需额外配置地址解析器即可直接连接到正确的节点地址。


内容的提问来源于stack exchange,提问作者danilo

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 08:05:29