【发布时间】:2015-10-05 11:52:39
【问题描述】:
我实际上正在寻找来自 ActiveMQ 的咨询或任何其他替代支持,以便在与 Consumer 关联的 MessageListener 完成处理消息时通知我。
MessageDelivered 建议似乎会在代理收到消息后立即通知。此外,MessageConsumed 建议声称在消费者收到消息时进行通知。
------------更新---------- --------
请在下面找到代码sn-p:
public class SampleListener implements MessageListener {
private Session session;
public SampleListener(Session session) {
this.session = session;
}
public void onMessage(Message message) {
try {
// do something
session.commit();
} catch (Exception e) {
e.printStackTrace();
}
}
}
public class SampleConsumer {
private boolean stopConnection = false;
public static void main(String[] args) {
new SampleConsumer().start();
}
public void start() {
ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory("tcp://localhost:61616");
Connection connection;
try {
connection = connectionFactory.createConnection();
connection.start();
Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
Destination destination = session.createTopic("test");
MessageConsumer messageConsumer = session.createConsumer(destination);
messageConsumer.setMessageListener(new SampleListener(session));
try {
synchronized (this) {
while (!stopConnection) {
wait();
}
}
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
session.close();
connection.close();
}
} catch (JMSException e) {
e.printStackTrace();
}
}
}
public void stop() {
synchronized (this) {
stopConnection = true;
notify();
}
}
}
public class SampleProducer implements MessageListener {
private boolean messageDelivered;
@Test
public void shouldTestSomething() throws JMSException, InterruptedException {
producerConnection = new ActiveMQConnectionFactory("tcp://localhost:61616").createConnection();
producerConnection.start();
Session session = producerConnection.createSession(true, SESSION_TRANSACTED);
Destination destination = session.createTopic("test");
MessageConsumer advisoryConsumer = session.createConsumer(AdvisorySupport.getMessageConsumedAdvisoryTopic(destination));
advisoryConsumer.setMessageListener(this);
Message message = session.createTextMessage("Hi");
Destination destination = session.createTopic("test");
MessageProducer producer = session.createProducer(destination);
producer.setDeliveryMode(DeliveryMode.PERSISTENT);
producer.send(message);
session.commit();
synchronized (this) {
while (!messageDelivered) {
wait();
}
}
session.close();
// some assertions
}
public void onMessage(Message message) {
// do something
synchronized (this) {
messageDelivered = true;
notify();
}
}
}
【问题讨论】:
-
如果您将 ack-mode 设置为事务处理,MessageConsumed 咨询仍会提前发生,会发生什么情况?根据issues.apache.org/jira/browse/AMQ-3361的关闭原因,它不应该。
-
Aksel,ack-mode 已经被交易。虽然 MessageDelivered 咨询发生得很早,但我无法通过 MessageConsumed 咨询获得任何通知。我有 PolicyEntry 参数来打开 MessageConsumed,但仍然没有任何消息排队等待 MessageConsumed 咨询。
-
我没有关注。问题是您认为消息消费会提前发生还是您无法为其启用建议?
-
后者。我确实在管理控制台的主题部分中看到了 MessageConsumed 咨询,但此主题没有任何反应,即配置为收听此咨询的侦听器永远不会被调用。然而,在收听 MessageDelivered 咨询时,通知会在消费者完成消费消息之前发生。
-
您在 onMessage() 中调用 session.commit() ?否则我无法猜测,并且需要查看一些代码才能更好地回答
标签: jms monitoring activemq