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信息


问题原因分析
报错信息Expected client server7:5557, got server5:5555本质是地址匹配失败:
- 客户端连接入口节点localhost:5555后,集群返回了实际存储Stream数据的节点地址(server7:5557)
- 但代码中硬编码的
addressResolver(address -> entryPoint)强制将所有节点地址解析为localhost:5555,导致客户端尝试连接server7时被错误转向到server5的地址,触发校验失败 - 生产者能正常工作是因为仅需连接入口节点完成消息发布,而消费者需要直接连接存储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
相关产品推荐
相关产品推荐

