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

如何配置RabbitMQ实现多消费者的Round-Robin轮询分发?

RabbitMQ轮询分发配置问题

项目架构

本项目包含三大核心组件:

  • 生产者(设备):通过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,导致消息被广播到所有临时队列。要实现轮询分发,需确保所有消费者绑定到同一个持久化队列,并调整监听容器配置:

  1. 确认队列配置有效性
    现有matrixQueue Bean默认创建的是持久化队列,无需修改,确保两个消费者实例都能连接到该队列。

  2. 修正消费者监听配置
    可以直接在@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);
    }
    
  3. 配置监听容器工厂,设置预取数
    添加容器工厂配置,设置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;
        }
    }
    
  4. 验证配置
    修改后重启两个消费者实例,通过RabbitMQ管理控制台确认两个消费者的队列名称均为matrix.queue。此时发送多条消息,日志会显示消息交替出现在8081和8082端口的实例中,实现轮询分发。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 17:02:04