Java发送的ActiveMQ消息无法接收,Web控制台发送则正常
问题分析与解决方案
核心问题原因
- 事务会话未提交:
connect()方法创建的是事务性会话(createSession(true, ...)的第一个参数true表示开启事务),但send()方法发送消息后没有提交事务,导致消息一直停留在客户端本地,ActiveMQ Broker根本没收到这条消息,因此receive()会一直阻塞等待。而Web控制台发送的消息是直接提交到Broker的,消费者能正常接收。 - 无超时的阻塞接收:
receive()方法使用consumer.receive()(无参数),这是无限期阻塞调用,没有消息时会一直挂起,直到手动终止程序。 - 资源未正确释放:
send()和receive()方法中创建的MessageProducer、MessageConsumer使用后未关闭,可能引发连接资源泄漏。
具体修复方案
1. 修正事务会话配置
如果不需要事务,直接创建非事务性会话:
session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
如果确实需要事务,在send()方法末尾添加事务提交:
session.commit();
2. 给接收方法添加超时
将consumer.receive()改为带超时的调用,比如等待5秒:
Message message = consumer.receive(5000); // 超时时间单位:毫秒
超时后会返回null,避免无限阻塞。
3. 释放资源
使用try-with-resources自动关闭MessageProducer和MessageConsumer,确保资源及时释放。
修改后的完整代码
import javax.jms.*; import org.apache.activemq.ActiveMQConnectionFactory; import arc.ipc.IService; public class ActiveMQService<K, V> implements IService<K, V> { private String brokerAddress; private ConnectionFactory connectionFactory; private Connection connection; private Session session; public ActiveMQService(String brokerAddress) { this.brokerAddress = brokerAddress; } @Override public void send(String topic, V value) throws JMSException { // 使用try-with-resources自动关闭producer try (Destination destination = session.createTopic(topic); MessageProducer producer = session.createProducer(destination)) { producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT); TextMessage message = session.createTextMessage(value.toString()); producer.send(message); // 如果是事务性会话,必须提交 // session.commit(); } } @Override public String receive(String topic) throws JMSException { // 使用try-with-resources自动关闭consumer try (Destination destination = session.createTopic(topic); MessageConsumer consumer = session.createConsumer(destination)) { // 设置5秒超时,避免无限阻塞 Message message = consumer.receive(5000); if (message instanceof TextMessage) { TextMessage textMessage = (TextMessage) message; return textMessage.getText(); } return null; } } @Override public void connect() throws JMSException { connectionFactory = new ActiveMQConnectionFactory(brokerAddress); connection = connectionFactory.createConnection(); connection.start(); // 改为非事务性会话(如果不需要事务) session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE); // 如果需要事务,保留下面的配置,并在send时commit // session = connection.createSession(true, Session.AUTO_ACKNOWLEDGE); } @Override public void disconnect() throws JMSException { if (session != null) { session.close(); } if (connection != null) { connection.close(); } } }
额外说明
- 后续接入Kafka时,注意Kafka的客户端API设计与JMS差异较大,建议基于你已有的
IService接口做统一封装,降低服务切换成本。 - 生产环境中,建议为ActiveMQ连接添加用户名密码认证,避免匿名访问风险。
内容的提问来源于stack exchange,提问作者Remixt
相关产品推荐
相关产品推荐

