【发布时间】: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) ,但是当客户端重新启动并与服务器连接时,得到了相同的体验,即消息重新传递。是否需要任何服务器端配置?我不想在队列中设置消息过期,但是如果消费者读取消息然后从队列中删除它,我想这样做