【问题标题】:In RabbitMQ how to consume multiple message or read all messages in a queue or all messages in exchange using specific key?在 RabbitMQ 中,如何使用特定键消费多条消息或读取队列中的所有消息或交换的所有消息?
【发布时间】:2021-02-25 15:04:46
【问题描述】:

我想使用给定键从特定队列或特定交换中消费多条消息。

所以场景如下:

发布者通过队列 1 发布消息 1 发布者通过队列 1 发布消息 2 发布者通过队列 1 发布消息 3 发布者通过队列 2 发布消息 4 发布者通过队列 2 发布消息 5 .. 消费者消费来自队列 1 的消息 一次性获取 [消息 1, 消息 2, 消息 3] 并在一次回调中处理它们

listen_to(queue_name , num_of_msg_to_fetch or all, function(messages){
//do some stuff with the returned list
});

消息不会同时出现,就像事件一样,我想将它们收集到队列中,打包并将它们发送给第三方。

我也看过这篇文章:

http://rabbitmq.1065348.n5.nabble.com/Consuming-multiple-messages-at-a-time-td27195.html

谢谢

【问题讨论】:

  • 我不认为这是一个很好的用例,你必须知道你期望有多少条消息才能等到它们都被读取后再处理它们
  • 我想阅读队列中的任何内容,最多不超过最大数量。的消息。所以代码将是这样的: while(queue_has_messages || max_num_of_msgs == 100){ queue.consume; max_num_of_msgs++; }
  • 再次,我不了解用例。你有你想做的事,但我觉得重新分析这一点可能有助于找到比你刚才所说的更好的解决方案。消息是原子的,因此最好能一一处理。如果不是,也许队列不是您的最佳解决方案。
  • 如果我的应用程序要自己处理消息,我会同意你的看法。但是我将多条消息聚合并打包,然后将它们发送到远程机器,这台机器本身 - 目前 - 不能直接使用消息。这就是我为它消费和打包多条消息的原因。
  • 它看起来像是一个非常自定义的实现,带有移动的基于时间的窗口(至少在生产者端)。使用 RMQ 没有什么可以轻松实现的,因为它的 API 绑定到反复调用的单个消息流。

标签: rabbitmq message-queue


【解决方案1】:

不要直接从队列中消费,因为队列遵循循环算法(AMQP 要求) 使用 shovel 将队列内容传输到扇出交换并直接从该交换中消费消息。您将获得所有连接的消费者的所有消息。 :)

【讨论】:

  • 我不确定您是否理解 OP 的问题。他们希望能够一次读取多条消息,以便能够对这些多条消息执行批处理操作,而不是一个接一个地读取它们。
【解决方案2】:

如果你想消费来自特定队列的多条消息,你可以尝试如下。

channel.queueDeclare(QUEUE_NAME, false, false,false, null);
Consumer consumer = new DefaultConsumer(channel){
   @Override
   public void handleDelivery(String consumerTag,
                                       Envelope envelope,
                                       AMQP.BasicProperties properties,
                                       byte[] body)
                    throws IOException {
                
                    String message = new String(body, "UTF-8");
                    logger.info("Recieved Message --> " + message);
                
   }
};

【讨论】:

    【解决方案3】:

    您可能需要在概念上将域消息与 RMQ 消息分开。作为生产者,您可以将多个域消息捆绑到单个 RMQ 消息中,并将.produce() 捆绑到 RMQ。请记住,由于窗口的存在,这种设计会引入超时和延迟(您可能会从 Kafka 中获得一些印象,它会以延迟为代价进行捆绑以优化 I/O)。

    作为消费者,您将拥有一个消费者,具有典型的 .handleDelivery 实现,它将转换接收到的主体以进行处理:byte[] -> Set[DomainMessage] -> your listener

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2013-12-18
      • 2014-04-08
      • 2018-08-10
      • 2023-03-28
      • 2020-04-25
      • 1970-01-01
      • 2016-08-21
      相关资源
      最近更新 更多