【问题标题】:Acknowledged messages are redelivered when crashed core client reconnects当崩溃的核心客户端重新连接时,重新传递确认的消息
【发布时间】:2019-04-21 22:31:34
【问题描述】:

我已经设置了一个在本地运行的独立 HornetQ 实例。出于测试目的,我使用 HornetQ core API 创建了一个消费者,它将每 500 毫秒接收一条消息。

当我的客户端连接并读取队列中的所有消息时,我在消费者端面临一个奇怪的行为,如果我强制关闭它(没有正确关闭会话/连接),那么下次我再次启动这个消费者时,它将读取队列中的旧消息。这是我的消费者示例:

// HornetQ 消费者代码

   public void readMessage() {
    ClientSession session = null;
    try {
        if (sf != null) {
            session = sf.createSession(true, true);

            ClientConsumer messageConsumer = session.createConsumer(JMS_QUEUE_NAME);
            session.start();

            while (true) {
                ClientMessage messageReceived = messageConsumer.receive(1000);
                if (messageReceived != null && messageReceived.getStringProperty(MESSAGE_PROPERTY_NAME) != null) {
                    System.out.println("Received JMS TextMessage:" + messageReceived.getStringProperty(MESSAGE_PROPERTY_NAME));
                    messageReceived.acknowledge();
                }

                Thread.sleep(500);
            }
        }
    } catch (Exception e) {
        LOGGER.error("Error while adding message by producer.", e);
    } finally {
        try {
            session.close();
        } catch (HornetQException e) {
            LOGGER.error("Error while closing producer session,", e);
        }
    }
}

有人能告诉我为什么它会这样工作吗?我应该在客户端/服务器端使用什么样的配置,这样如果消费者读取了一条消息,它就会从队列中删除它?

【问题讨论】:

  • 您似乎使用的是 HornetQ“核心”API 而不是 JMS(因为 ClientMessage 不是 JMS 对象)。你能确认一下吗?
  • 抱歉回复晚了。是的,这是真的 ClientMessage 来自 HornetQ 核心 API。但它与问题有什么关系呢?
  • 这与问题有关,因为这两个 API 的行为不同。确认完成后,您可能不会提交会话。您是在任何时候提交会话还是在创建会话时启用了自动提交以进行确认?
  • 我更新了有问题的代码。这就是我在 HornetQ 中测试消费者的方式
  • 我也尝试过使用 sf.createSession(true,false) ,但是当客户端重新启动并与服务器连接时,得到了相同的体验,即消息重新传递。是否需要任何服务器端配置?我不想在队列中设置消息过期,但是如果消费者读取消息然后从队列中删除它,我想这样做

标签: java hornetq


【解决方案1】:

确认完成后,您没有提交会话,并且您没有创建启用了自动提交确认的会话。因此,您应该执行以下操作之一:

  • 在一次或多次调用acknowledge() 后显式调用session.commit()
  • 或通过使用sf.createSession(true,true)sf.createSession(false,true) 创建会话来启用隐式自动提交确认(控制自动提交确认的布尔值是第二个)。

请记住,当您为确认启用自动提交时,有一个内部缓冲区需要在确认刷新到代理之前达到特定大小。像这样的批处理确认可以显着提高某些大容量用例的性能。默认情况下,您需要确认 1,048,576 字节的消息才能刷新缓冲区并将确认发送到代理。您可以通过在您的ServerLocator 实例上调用setAckBatchSize 或使用不同的createSession 方法(例如sf.createSession(true, true, myAckBatchSize))来更改此缓冲区的大小。

如果确认缓冲区未刷新并且您的客户端崩溃,则当客户端返回时,相应的消息仍将在队列中。如果缓冲区没有达到其阈值,当消费者正常关闭时,它仍然会被刷新。

【讨论】:

  • 它仍然在做同样的事情。我认为我在我的生产者或消费者代码中做了一些愚蠢的事情。我再次更新了一个问题,请您看看并建议我这是在这两个地方使用 autoCommit 会话的正确方法。
  • 我已从问题中删除了您的生产者代码,因为它与问题完全无关,并且我已对我的回答添加了一些说明。
  • 实际上,我正在关注 HornetQ 源代码和消费者实现中可用的示例,没有提及任何用于在确认后设置 MaxBatchSize 的内容。我现在正在查看文档以了解如何设置此信息。
  • 我已经概述了设置确认批量大小的最简单方法。至于 HornetQ 源代码中的示例,我的猜测是它依赖于 close() 调用来刷新确认缓冲区,或者它确认了足够的消息来自行刷新缓冲区。您能否提供相关 HornetQ 来源的链接?
  • 完美! ClientConsumer.Close() 做到了这一点。感谢您提供所有这些信息。
猜你喜欢
  • 1970-01-01
  • 2015-07-21
  • 2016-04-29
  • 2016-05-09
  • 1970-01-01
  • 2015-01-17
  • 1970-01-01
  • 2015-05-28
  • 2017-05-26
相关资源
最近更新 更多