如何配置RabbitMQ实现多消费者的Round-Robin轮询分发?
项目架构
本项目包含三大核心组件:
- 生产者(设备):通过MQTT协议向Broker发送数据
- RabbitMQ Broker
- 消费者(基于Spring AMQP):负责应用与Broker的通信
架构流程:(设备MQTT协议) 生产者 → RabbitMQ ← 消费者(Spring应用使用AMQP协议)
模拟生产者配置
使用MQTTX Client Toolbox模拟生产者:
quick.white.rabbit - 主题 ewogICAibWF0cml4IjoiZm9sbG93IHRoZSB3aGl0ZSByYWJiaXTigKYiLAogICAiZGF0ZSI6Ik1hcmNoIDMxLCAxOTk5Igp9 - Base64编码数据
RabbitMQ配置代码
@Configuration public class RabbitMQConfig { @Bean public TopicExchange exchange() { return new TopicExchange("amq.topic"); } @Bean public Queue matrixQueue(){ return new Queue("matrix.queue"); } @Bean public Binding binding(TopicExchange exchange, Queue matrixQueue) { return BindingBuilder.bind(matrixQueue) .to(exchange) .with("*.*.rabbit"); } }
消费者代码
@Service public class RabbitMQConsumer { static void dumb(Object payload, Map<String, Object> headers){ System.out.println("-----------------------------------"); headers.forEach((k,v) -> System.out.println(k + "=" + v)); System.out.println(payload); } @RabbitListener(queues = "#{matrixQueue.name}") public void consumePayload(@Payload String encodedMessage, @Headers Map<String, Object> headers){ String payload = new String(Base64.getDecoder().decode(encodedMessage)); dumb(payload, headers); } }
当前问题
启动两个相同的消费者实例(端口8081和8082):
java -jar target/consumer_rabbitmq-0.0.1-SNAPSHOT.jar --server.port=8081 java -jar target/consumer_rabbitmq-0.0.1-SNAPSHOT.jar --server.port=8082
发现同一条消息会同时投递到两个消费者实例中,日志示例如下:
java -jar target/deep_dive_rabbitmq-0.0.1-SNAPSHOT.jar --server.port=8081 ----------------------------------- amqp_receivedDeliveryMode=NON_PERSISTENT amqp_receivedRoutingKey=quick.white.rabbit amqp_receivedExchange=amq.topic x-mqtt-publish-qos=0 x-mqtt-dup=false amqp_deliveryTag=3 amqp_consumerQueue=spring.gen-X6e3uiz5S-K4_RA0tt89cg amqp_redelivered=false id=debc649f-9ae1-82a8-be5f-0a3f2c9e2f87 amqp_consumerTag=amq.ctag-brMleQwCDNmcCHyT_lbF8A amqp_lastInBatch=false timestamp=1708633236541 { "matrix":"follow the white rabbit…", "date":"March 31, 1999" }
java -jar target/deep_dive_rabbitmq-0.0.1-SNAPSHOT.jar --server.port=8082 ----------------------------------- amqp_receivedDeliveryMode=NON_PERSISTENT amqp_receivedRoutingKey=quick.white.rabbit amqp_receivedExchange=amq.topic x-mqtt-publish-qos=0 x-mqtt-dup=false amqp_deliveryTag=3 amqp_consumerQueue=spring.gen-Ym1vcVbUQ4aZqnG3Bu4OyA amqp_redelivered=false id=dd1a5d91-4972-dcca-0669-dc40f9d85a15 amqp_consumerTag=amq.ctag-zqy6gfLSLaZSf3H3W00M1A amqp_lastInBatch=false timestamp=1708633236541 { "matrix":"follow the white rabbit…", "date":"March 31, 1999" }
期望效果
配置RabbitMQ实现轮询(Round-Robin)分发,消息按顺序依次投递到不同消费者:
MESSAGE 1 // port 8081 MESSAGE 2 // port 8082 MESSAGE 3 // port 8081 MESSAGE 4 // port 8082
解决方案
当前问题的核心是两个消费者实际绑定的是RabbitMQ自动创建的临时队列(日志中amqp_consumerQueue为spring.gen-xxx),而非我们定义的matrix.queue,导致消息被广播到所有临时队列。要实现轮询分发,需确保所有消费者绑定到同一个持久化队列,并调整监听容器配置:
确认队列配置有效性
现有matrixQueueBean默认创建的是持久化队列,无需修改,确保两个消费者实例都能连接到该队列。修正消费者监听配置
可以直接在@RabbitListener中写死队列名称,避免SpEL表达式解析异常:@RabbitListener(queues = "matrix.queue") public void consumePayload(@Payload String encodedMessage, @Headers Map<String, Object> headers){ String payload = new String(Base64.getDecoder().decode(encodedMessage)); dumb(payload, headers); }配置监听容器工厂,设置预取数
添加容器工厂配置,设置prefetchCount=1,确保RabbitMQ每次只给一个消费者分发一条消息,处理完成后再分发下一条,保证轮询逻辑稳定生效:@Configuration public class RabbitMQConfig { // 原有Exchange、Queue、Binding Bean... @Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ConnectionFactory connectionFactory) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory); factory.setPrefetchCount(1); // 每次预取1条消息,避免批量分发 return factory; } }验证配置
修改后重启两个消费者实例,通过RabbitMQ管理控制台确认两个消费者的队列名称均为matrix.queue。此时发送多条消息,日志会显示消息交替出现在8081和8082端口的实例中,实现轮询分发。
内容的提问来源于stack exchange,提问作者Romillion

