【问题标题】:RabbitMQ. Java client. Is it possible to acknowledge message not on the same thread it was received?兔MQ。 Java 客户端。是否可以确认消息不在收到的同一线程上?
【发布时间】:2016-10-09 23:50:42
【问题描述】:

我想获取几条消息,处理它们并在此之后将它们全部确认。所以基本上我收到一条消息,将其放入某个队列并继续接收来自兔子的消息。不同的线程将使用收到的消息监视此队列,并在数量足够时对其进行处理。我所能找到的关于 ack 的所有内容仅包含在同一线程上处理的一条消息的示例。像这样(来自官方文档):

channel.basicQos(1);

final 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");

    System.out.println(" [x] Received '" + message + "'");
    try {
      doWork(message);
    } finally {
      System.out.println(" [x] Done");
      channel.basicAck(envelope.getDeliveryTag(), false);
    }
  }
};

文档也这样说:

通道实例不能在线程之间共享。应用 应该更喜欢每个线程使用一个 Channel 而不是共享同一个 跨多个线程的通道。虽然通道上的一些操作是 安全地同时调用,有些不是并且会导致不正确 帧在电线上交错。

所以我在这里很困惑。如果我正在确认一些消息,同时频道正在接收来自 rabbit 的另一条消息,那么它当时是否被认为是两个操作?在我看来是的。

我尝试从不同线程确认同一通道上的消息,它似乎有效,但文档说我不应该在线程之间共享通道。所以我尝试用不同的频道在不同的线程上做确认,但是失败了,因为这个频道的传递标签是未知的。

是否可以确认消息不在收到的同一线程上?

UPD 我想要的示例代码。它在 scala 中,但我认为它很简单。

 case class AmqpMessage(envelope: Envelope, msgBody: String)

    val queue = new ArrayBlockingQueue[AmqpMessage](100)

    val consumeChannel = connection.createChannel()
    consumeChannel.queueDeclare(queueName, true, false, true, null)
    consumeChannel.basicConsume(queueName, false, new DefaultConsumer(consumeChannel) {
      override def handleDelivery(consumerTag: String,
                                  envelope: Envelope,
                                  properties: BasicProperties,
                                  body: Array[Byte]): Unit = {
        queue.put(new AmqpMessage(envelope, new String(body)))
      }
    })

    Future {
      // this is different thread
      val channel = connection.createChannel()
      while (true) {
        try {
          val amqpMessage = queue.take()
          channel.basicAck(amqpMessage.envelope.getDeliveryTag, false) // doesn't work
          consumeChannel.basicAck(amqpMessage.envelope.getDeliveryTag, false) // works, but seems like not thread safe
        } catch {
          case e: Exception => e.printStackTrace()
        }
      }
    }

【问题讨论】:

  • 能否请您详细说明这部分I want to fetch several messages, handle them and ack them all together after that. So basically I receive a message, put it in some **queue** and continue receiving messages from rabbit.**之间的队列是什么?另一个 RMQ 队列还是其他什么?
  • @cantSleep现在只是简单的java内存阻塞队列。我已经发布了示例来澄清。抱歉误导。
  • "是否可以确认消息不在收到的同一线程上?"答案是“是”

标签: java multithreading rabbitmq


【解决方案1】:

尽管文档非常严格,通道上的某些操作可以安全地同时调用。 只要消费acking是您在频道上执行的唯一操作,您就可以在不同的线程中确认消息。

查看这个 SO 问题,它处理相同的事情:

RabbitMQ and channels Java thread safety

【讨论】:

    【解决方案2】:

    对我来说,您的解决方案是正确的。您没有跨线程共享频道。 您永远不会将通道对象传递给另一个线程,而是在接收消息的同一线程上使用它。

    你是不可能的

    '正在确认一些消息,同时频道正在接收来自 rabbit 的另一条消息'

    如果您在 handleDelivery 方法中,则该线程被您的代码阻塞,并且没有机会接收另一条消息。

    如您所见,您无法使用用于接收消息的通道以外的通道确认消息。

    您必须使用相同的通道进行确认,并且必须在接收消息的同一线程上执行此操作。因此,您可以将通道对象传递给其他方法、类,但必须注意不要将其传递给另一个线程。

    我在我的项目中使用此解决方案它使用 RabbitMQ 侦听器和 Spring 集成。对于每条 AMQP 消息,都会创建一个 org.springframework.integration.Message。该消息将 AMPQ 消息正文作为有效负载,并将 AMQP 通道和传递标记作为我的 org.springframework.integration.Message 的标头。

    如果你想确认多条消息,并且它们是在同一个频道上传递的,你应该使用

    channel.basicAck(envelope.getDeliveryTag(), true);
    

    对于多通道,高效的算法是

    1. 假设您有 100 条消息,使用 10 个渠道传递
    2. 您需要找到每个频道的 max deliveryTag。
    3. 调用channel.basicAck(maxDeliveryTagForThatChannel, true);

    这样,您需要 10 个 basicAck(网络往返)而不是 100 个。

    【讨论】:

    • 这不是我的解决方案,它是文档中的示例,用于说明兔子团队建议如何处理 ack。但它不适合我的需要。可能我需要发布我想做的例子。
    • 似乎可以从不同的频道确认。不过还是谢谢你。
    【解决方案3】:

    正如文档所说,每个线程一个通道,其余通道没有限制。

    我只想就你的例子说几句话。您在这里尝试做的事情是错误的。只有在您从ArrayBlockingQueue 获取消息后才需要确认消息,因为一旦您将它放在那里,它就会一直留在那里。将其 ACK 到 RMQ 与其他 ArrayBlockingQueue 队列无关。

    【讨论】:

    • 是的,它在那里,但它没有被处理。我想在处理完几条消息后确认(从 java 队列中取出它们,做一些事情,确认),而不是在我将它们放入 java 队列之后。
    • 因此,一旦收到就处理它们,而不是执行多个 ACK​​。从消费线程
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-02-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-01-12
    • 2013-10-27
    相关资源
    最近更新 更多