You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Java发送的ActiveMQ消息无法接收,Web控制台发送则正常

问题分析与解决方案

核心问题原因

  1. 事务会话未提交:connect()方法创建的是事务性会话(createSession(true, ...)的第一个参数true表示开启事务),但send()方法发送消息后没有提交事务,导致消息一直停留在客户端本地,ActiveMQ Broker根本没收到这条消息,因此receive()会一直阻塞等待。而Web控制台发送的消息是直接提交到Broker的,消费者能正常接收。
  2. 无超时的阻塞接收:receive()方法使用consumer.receive()(无参数),这是无限期阻塞调用,没有消息时会一直挂起,直到手动终止程序。
  3. 资源未正确释放: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 15:50:16