使用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
相关产品推荐
相关产品推荐

