ActiveMQ Artemis 2.32.0消费者偶发无法接收JMS Topic消息问题
ActiveMQ Artemis 消费者偶尔无法接收Topic消息的问题排查
在使用ActiveMQ Artemis 2.32.0时,将两个消费者(Tommy、Harry)放在独立线程中运行,生产者向JMS Topic发布消息,偶尔会出现其中一个消费者因未收到消息阻塞的情况,并非每次复现;但如果直接通过各自main方法启动消费者(不使用线程),两个消费者都能正常接收消息。
相关代码
Tommy.java
public class Tommy implements Runnable { private static final Logger logger = LoggerFactory.getLogger(Tommy.class); public void receiveMessage() throws NamingException { // Create a new initial context, which loads from jndi.properties file Context context = new InitialContext(); // Lookup an existing Destination which is a topic in our example Topic topic = (Topic)context.lookup("jms/test/topic"); //Object in a try-with-resources block the close method will be called automatically at the end of the block. try(ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); JMSContext jmsContext = connectionFactory.createContext()) { //Create consumer and receive String message on the fly (i.e. without need to type caste to Message etc.) String messageReceived = jmsContext.createConsumer(topic).receiveBody(String.class); logger.info("Message received by Tommy >>> {}", messageReceived); } } @Override public void run() { try { receiveMessage(); } catch (NamingException e) { e.printStackTrace(); } } }
Harry.java
public class Harry implements Runnable { private static final Logger logger = LoggerFactory.getLogger(Harry.class); public static void receiveMessage() throws NamingException { // Create a new initial context, which loads from jndi.properties file Context context = new InitialContext(); // Lookup an existing Destination which is a topic in our example Topic topic = (Topic)context.lookup("jms/test/topic"); //Object in a try-with-resources block the close method will be called automatically at the end of the block. try(ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); JMSContext jmsContext = connectionFactory.createContext()) { //Create consumer and receive String message on the fly (i.e. without need to type caste to Message etc.) String messageReceived = jmsContext.createConsumer(topic).receiveBody(String.class); logger.info("Message received by Harry >>> {}", messageReceived); } } @Override public void run() { try { receiveMessage(); } catch (NamingException e) { e.printStackTrace(); } } }
WeatherChannel.java
public class WeatherChannel { private static final Logger logger = LoggerFactory.getLogger(WeatherChannel.class); public void broadcastMessage() throws NamingException, JMSException, InterruptedException { // Create a new initial context, which loads from jndi.properties file Context context = new InitialContext(); // Lookup an existing Destination which is a topic in our example Topic topic = (Topic)context.lookup("jms/test/topic"); //Object in a try-with-resources block the close method will be called automatically at the end of the block. try(ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); JMSContext jmsContext = connectionFactory.createContext()) { //Create producer and send message on the fly jmsContext.createProducer().send(topic, "Today's weather at Kolkata is pleasant with max temp 27 C"); logger.info("Message sent successfully by producer"); } } }
调用逻辑
new Thread(new Tommy()).start(); new Thread(new Harry()).start(); new WeatherChannel().broadcastMessage();
问题原因
核心问题是线程启动的异步性导致生产者可能在消费者完成订阅前发送消息。
JMS Topic的默认非持久订阅模式下,只有当消费者成功建立连接、完成订阅流程后,才能收到后续发布的消息。而线程启动需要时间,包括初始化JNDI上下文、创建连接工厂、建立JMS连接、创建消费者这些步骤都存在耗时。如果生产者在这些步骤完成前就发送了消息,未完成订阅的消费者就无法接收到这条消息,进而因为调用无超时的receiveBody()方法陷入永久等待。
当直接通过main方法启动消费者时,消费者的初始化流程在生产者发送消息前已经完成,所以能稳定接收消息。
解决方案
可以通过同步机制确保所有消费者完成订阅后,生产者再发送消息,常用方式是使用CountDownLatch:
修改消费者代码(Tommy为例,Harry同理)
public class Tommy implements Runnable { private static final Logger logger = LoggerFactory.getLogger(Tommy.class); private final CountDownLatch latch; public Tommy(CountDownLatch latch) { this.latch = latch; } public void receiveMessage() throws NamingException { Context context = new InitialContext(); Topic topic = (Topic)context.lookup("jms/test/topic"); try(ActiveMQConnectionFactory connectionFactory = new ActiveMQConnectionFactory(); JMSContext jmsContext = connectionFactory.createContext()) { // 创建消费者后,标记订阅完成 jmsContext.createConsumer(topic); latch.countDown(); String messageReceived = jmsContext.createConsumer(topic).receiveBody(String.class); logger.info("Message received by Tommy >>> {}", messageReceived); } } @Override public void run() { try { receiveMessage(); } catch (NamingException e) { e.printStackTrace(); } } }
修改调用逻辑
// 初始化计数器,数量等于消费者个数 CountDownLatch latch = new CountDownLatch(2); new Thread(new Tommy(latch)).start(); new Thread(new Harry(latch)).start(); // 等待所有消费者完成订阅 latch.await(); // 再发送消息 new WeatherChannel().broadcastMessage();
另外,也可以给receiveBody()设置超时时间避免永久阻塞,比如receiveBody(String.class, 5000)(5秒超时),但这只是临时规避阻塞问题,无法解决消息接收不到的核心问题。
内容的提问来源于stack exchange,提问作者Subhajit Roy
相关产品推荐
相关产品推荐

