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

ActiveMQ Artemis消费者的MessageListener无法接收消息且应用退出

问题:ActiveMQ Artemis 消费者绑定MessageListener后自动退出,无法持续运行

我正从ActiveMQ Classic 5.18.1迁移至ActiveMQ Artemis,已按文档更新生产者和消费者代码,生产者可正常发送消息。但绑定了MessageListener的消费者应用仍会退出,我期望MessageListener能持续运行,请问遗漏了什么?

import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.Session;
import javax.jms.Topic;
import javax.naming.Context;
import javax.naming.InitialContext;
import javax.naming.NamingException;
import java.util.Properties;


public class Consumer1 {
    public static void main(String[] args) throws Exception {
        String topicName = "jms/topic/default_EventTopic";
        final Properties properties = new Properties();
        properties.setProperty(Context.INITIAL_CONTEXT_FACTORY, "org.apache.activemq.artemis.jndi.ActiveMQInitialContextFactory");
        properties.setProperty(Context.PROVIDER_URL, "tcp://localhost:61616");

        Context context = new InitialContext(properties);
        ConnectionFactory factory = (ConnectionFactory) context.lookup("ConnectionFactory");
        Connection connection = factory.createConnection();
        try {
            connection.setClientID("SampleConnection");
            Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            final Topic topic = (Topic) context.lookup(topicName);
            final MessageConsumer consumer = session.createDurableSubscriber(topic, "EventConsumer");

            connection.start();
            consumer.setMessageListener(message -> {
                try {
                    System.out.println("Type " + message.getStringProperty("type"));
                } catch (JMSException e) {
                    e.printStackTrace();
                }
            });
        } catch (final JMSException | NamingException cause) {
            System.out.println("Error while creating Message Consumer");
            cause.printStackTrace();
        } catch (Exception e ){
            e.printStackTrace();
        }
    }
}

核心问题

你的消费者应用启动后直接退出,是因为主线程(main方法)执行完毕后没有阻塞,导致JVM直接终止。ActiveMQ Artemis客户端的消息监听线程默认是守护线程,当主线程结束后,JVM不会等待守护线程,直接退出,所以MessageListener没有机会持续接收消息。

解决方案

在main方法末尾添加阻塞逻辑,让主线程保持运行,避免JVM退出。常用的方式有两种:

方式1:使用CountDownLatch(推荐,更优雅)

通过CountDownLatch让主线程等待,直到收到退出信号(比如手动中断),同时添加关闭钩子确保资源正常释放:

import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.Session;
import javax.jms.Topic;
import javax.naming.Context;
import javax.naming.InitialContext;
import javax.naming.NamingException;
import java.util.Properties;
import java.util.concurrent.CountDownLatch;

public class Consumer1 {
    public static void main(String[] args) throws Exception {
        CountDownLatch latch = new CountDownLatch(1);
        
        String topicName = "jms/topic/default_EventTopic";
        final Properties properties = new Properties();
        properties.setProperty(Context.INITIAL_CONTEXT_FACTORY, "org.apache.activemq.artemis.jndi.ActiveMQInitialContextFactory");
        properties.setProperty(Context.PROVIDER_URL, "tcp://localhost:61616");

        Context context = new InitialContext(properties);
        ConnectionFactory factory = (ConnectionFactory) context.lookup("ConnectionFactory");
        Connection connection = factory.createConnection();
        
        // 添加关闭钩子,程序退出时关闭连接
        Runtime.getRuntime().addShutdownHook(new Thread(() -> {
            try {
                connection.close();
                latch.countDown();
            } catch (JMSException e) {
                e.printStackTrace();
            }
        }));
        
        try {
            connection.setClientID("SampleConnection");
            Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
            final Topic topic = (Topic) context.lookup(topicName);
            final MessageConsumer consumer = session.createDurableSubscriber(topic, "EventConsumer");

            connection.start();
            consumer.setMessageListener(message -> {
                try {
                    System.out.println("Type " + message.getStringProperty("type"));
                } catch (JMSException e) {
                    e.printStackTrace();
                }
            });
            
            // 阻塞主线程
            latch.await();
        } catch (final JMSException | NamingException cause) {
            System.out.println("Error while creating Message Consumer");
            cause.printStackTrace();
            latch.countDown();
        } catch (Exception e ){
            e.printStackTrace();
            latch.countDown();
        }
    }
}

方式2:使用System.in.read()(简单测试用)

如果只是快速测试,可以直接在main方法末尾添加读取控制台输入的代码,让主线程一直等待:

// 在try-catch块结束后添加
try {
    // 阻塞主线程,直到用户输入任意字符
    System.in.read();
} catch (Exception e) {
    e.printStackTrace();
}

补充说明

ActiveMQ Classic中客户端线程默认是非守护线程,主线程结束后程序不会退出,但Artemis调整了这个默认行为,使用守护线程减少资源占用,这是迁移时容易忽略的差异点。

内容的提问来源于stack exchange,提问作者Sohan Soni

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 02:10:35