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

使用JNDI或ActiveMQConnectionFactory连接EmbeddedActiveMQ时消息未被消费

嵌入式ActiveMQ Artemis消息无法消费问题

代码能够正常初始化连接、会话等组件,但消息始终无法被消费,且日志中未出现任何异常。以下测试代码可复现该问题:

import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
import javax.jms.Queue;
import javax.jms.Session;
import javax.naming.Context;

import java.io.File;
import java.util.Hashtable;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

import org.apache.activemq.artemis.api.core.QueueConfiguration;
import org.apache.activemq.artemis.api.core.SimpleString;
import org.apache.activemq.artemis.core.config.Configuration;
import org.apache.activemq.artemis.core.config.impl.ConfigurationImpl;
import org.apache.activemq.artemis.core.server.embedded.EmbeddedActiveMQ;
import org.apache.activemq.artemis.core.settings.impl.AddressSettings;
import org.apache.activemq.artemis.jndi.ActiveMQInitialContextFactory;
import org.junit.After;
import org.junit.Before;

public class Test {

    EmbeddedActiveMQ jmsServer;

    final String QUEUE_NAME = "myQueue";

    @Before
    public void setUp() throws Exception {
        final String baseDir = File.separator + "tmp";
        final EmbeddedActiveMQ embeddedActiveMQ = new EmbeddedActiveMQ();
        final Configuration config = new ConfigurationImpl();
        config.setPersistenceEnabled(true);
        config.setBindingsDirectory(baseDir + File.separator + "bindings");
        config.setJournalDirectory(baseDir + File.separator + "journal");
        config.setPagingDirectory(baseDir + File.separator + "paging");
        config.setLargeMessagesDirectory(baseDir + File.separator + "largemessages");
        config.setSecurityEnabled(false);

        AddressSettings adr = new AddressSettings();
        adr.setDeadLetterAddress(new SimpleString("DLQ"));
        adr.setExpiryAddress(new SimpleString("ExpiryQueue"));
        config.addAddressSetting("#", adr);

        config.addAcceptorConfiguration("invmConnectionFactory", "vm://0");
        embeddedActiveMQ.setConfiguration(config);
        this.jmsServer = embeddedActiveMQ;
        this.jmsServer.start();

        System.out.println("creating queue");
        final boolean isSuccess = jmsServer.getActiveMQServer().createQueue(new QueueConfiguration(QUEUE_NAME)) != null;
        if(isSuccess) {
            System.out.println(QUEUE_NAME + "queue created");
        }

    }

    @After
    public void tearDown() {
        try {
            this.jmsServer.stop();
        } catch(Exception e) {
            // ignore
        }
    }

    @org.junit.Test
    public void simpleTest() throws Exception {

        Hashtable d = new Hashtable();
        d.put(Context.INITIAL_CONTEXT_FACTORY, "org.apache.activemq.artemis.jndi.ActiveMQInitialContextFactory");
        d.put("connectionFactory.invmConnectionFactory", "vm://0");
        final ActiveMQInitialContextFactory activeMQInitialContextFactory = new ActiveMQInitialContextFactory();
        Context initialContext = activeMQInitialContextFactory.getInitialContext(d);
        ConnectionFactory connectionFactory = (ConnectionFactory) initialContext.lookup("invmConnectionFactory");
        Connection connection = connectionFactory.createConnection();
        Session session = connection.createSession(true, Session.SESSION_TRANSACTED);
        Queue queue = session.createQueue(QUEUE_NAME);

        MessageProducer producer = session.createProducer(queue);
        MessageConsumer consumer = session.createConsumer(queue);
        CountDownLatch latch = new CountDownLatch(1);
        consumer.setMessageListener(message -> {
            System.out.println("=== " + message);
            try {
                message.acknowledge();
                session.commit();
                latch.countDown();
            } catch(JMSException e) {
                e.printStackTrace();
            }
        });

        connection.start();
        producer.send(session.createMessage());
        session.commit();
        if(!latch.await(2, TimeUnit.SECONDS)) {
            throw new IllegalStateException();
        }
        connection.close();
    }
}

内容的提问来源于stack exchange,提问作者Wojciech Piotrowiak

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 02:55:23