提问者:小点点

RabbitMQ/AMQP:单个队列,同一消息的多个使用者?


我刚刚开始使用RabbitMQ和AMQP。

  • 我有一个消息队列
  • 我有多个消费者,我想用相同的消息做不同的事情。

RabbitMQ的大部分文档似乎都集中在循环(round-robin)上,即单个消息由单个消费者使用,负载在每个消费者之间分散。这的确是我目击的行为。

例如:生产者只有一个队列,每2秒发送一次消息:

var amqp = require('amqp');
var connection = amqp.createConnection({ host: "localhost", port: 5672 });
var count = 1;

connection.on('ready', function () {
  var sendMessage = function(connection, queue_name, payload) {
    var encoded_payload = JSON.stringify(payload);  
    connection.publish(queue_name, encoded_payload);
  }

  setInterval( function() {    
    var test_message = 'TEST '+count
    sendMessage(connection, "my_queue_name", test_message)  
    count += 1;
  }, 2000) 


})

这里有一个消费者:

var amqp = require('amqp');
var connection = amqp.createConnection({ host: "localhost", port: 5672 });
connection.on('ready', function () {
  connection.queue("my_queue_name", function(queue){
    queue.bind('#'); 
    queue.subscribe(function (message) {
      var encoded_payload = unescape(message.data)
      var payload = JSON.parse(encoded_payload)
      console.log('Recieved a message:')
      console.log(payload)
    })
  })
})

如果我启动消费者两次,我可以看到每个消费者都在以循环行为消费交替消息。我将在一个终端机看到1、3、5号信息,在另一个终端机看到2、4、6号信息。

我的问题是:

>

  • 我可以让每个消费者接收相同的消息吗?也就是说,两个消费者都得到消息1、2、3、4、5、6?这在AMQP/RabbitMQ Speak中叫什么?它通常是如何配置的?

    通常这样做吗?我是否应该让交换将消息路由到两个单独的队列中,并使用一个使用者?


  • 共1个答案

    匿名用户

    我能让每个消费者收到相同的消息吗?也就是说,两个消费者都得到消息1、2、3、4、5、6?这在AMQP/RabbitMQ Speak中叫什么?它通常是如何配置的?

    不,如果消费者在同一队列中,则不是。摘自RabbitMQ的AMQP概念指南:

    在AMQP0-9-1中,消息在使用者之间是负载平衡的,这一点很重要。

    这似乎意味着队列中的循环行为是给定的,而不是可配置的。即,为了使同一消息ID由多个使用者处理,需要单独的队列。

    通常这样做吗?我是否应该让交换将消息路由到两个单独的队列中,并使用一个使用者?

    不,它不是,单个队列/多个消费者,每个消费者处理相同的消息ID是不可能的。让exchange将消息路由到两个单独的队列中确实更好。

    由于我不需要太复杂的路由,扇出交换可以很好地处理这一点。我之前没有过多地关注交换,因为node-amqp有一个“默认交换”的概念,允许您直接将消息发布到连接,然而大多数AMQP消息都发布到特定的交换。

    这是我的扇出交换,发送和接收:

    var amqp = require('amqp');
    var connection = amqp.createConnection({ host: "localhost", port: 5672 });
    var count = 1;
    
    connection.on('ready', function () {
      connection.exchange("my_exchange", options={type:'fanout'}, function(exchange) {   
    
        var sendMessage = function(exchange, payload) {
          console.log('about to publish')
          var encoded_payload = JSON.stringify(payload);
          exchange.publish('', encoded_payload, {})
        }
    
        // Recieve messages
        connection.queue("my_queue_name", function(queue){
          console.log('Created queue')
          queue.bind(exchange, ''); 
          queue.subscribe(function (message) {
            console.log('subscribed to queue')
            var encoded_payload = unescape(message.data)
            var payload = JSON.parse(encoded_payload)
            console.log('Recieved a message:')
            console.log(payload)
          })
        })
    
        setInterval( function() {    
          var test_message = 'TEST '+count
          sendMessage(exchange, test_message)  
          count += 1;
        }, 2000) 
     })
    })