Quarkus Artemis JMS库无法将消息投递到ActiveMQ Artemis死信队列
问题描述
我在Quarkus项目中实现了一个JMS消费者,监听ActiveMQ Artemis(版本2.32.0)上的anycast队列my.edu.queue。为测试业务错误时消息投递死信队列的逻辑,我设置每处理5条消息就模拟一次业务异常(抛出异常)。
参考ActiveMQ Artemis的死信队列示例实现,我基于quarkus-artemis-jms库完成了类似开发,但目前遇到问题:调用Session.rollback()后,消息并未进入死信队列。我已在broker.xml中为目标队列设置max-delivery-attempts为0,以下是我的JmsMessageConsumer代码:
@ApplicationScoped public class JmsMessageConsumer { private static final Logger log = Logger.getLogger(JmsMessageConsumer.class.getName()); @Inject ConnectionFactory connectionFactory; private JMSContext context; private Connection connection; private Session session; private MessageConsumer consumer; private MessageProducer dlqProducer; @ConfigProperty(name = "my.edu.queue.name", defaultValue = "my.edu.queue") String queueName; private AtomicInteger counter = new AtomicInteger(0); void onStart(@Observes StartupEvent ev) throws JMSException { connection = connectionFactory.createConnection(); session = connection.createSession(true, Session.SESSION_TRANSACTED); Queue queue = session.createQueue(queueName); consumer = session.createConsumer(queue); Queue dlq = session.createQueue(queueName + ".dlq"); dlqProducer = session.createProducer(dlq); connection.start(); receiveMessages(); } void onStop(@Observes ShutdownEvent ev) throws JMSException { connection.close(); } private void receiveMessages() { try { consumer.setMessageListener(message -> { try { processMessage(message.getBody(String.class)); session.commit(); } catch (Exception e) { log.severe("Error processing message: %s".formatted(e.getMessage())); try { // sendToDLQ(message); session.rollback(); } catch (JMSException ex) { throw new RuntimeException(ex); } } }); } catch (JMSException e) { log.severe("Error setting message listener: %s".formatted(e.getMessage())); } } private void processMessage(String text) { counter.incrementAndGet(); if (counter.get() % 5 == 0) { throw new RuntimeException("Error in business logic"); } log.info("Processed message: " + text); } }
解决方案
核心问题:max-delivery-attempts配置错误
ActiveMQ Artemis中,max-delivery-attempts设为0代表无限重试,Broker会不断重新投递失败的消息,永远不会将其转入死信队列。你需要将该值设为大于0的整数(比如1),指定消息最多重试的次数,达到次数后Broker会自动将消息路由到死信队列。
代码优化建议
- 移除手动DLQ生产者:ActiveMQ Artemis会自动处理死信路由,无需手动创建Producer发送消息到DLQ,配置正确后Broker会自动完成转移。
- 使用JMSContext简化资源管理:Quarkus推荐使用
JMSContext替代手动管理Connection和Session,它会自动处理资源的创建与释放,更符合Quarkus的轻量、自动管理特性。
优化后的代码示例:
@ApplicationScoped public class JmsMessageConsumer { private static final Logger log = Logger.getLogger(JmsMessageConsumer.class.getName()); @Inject JMSContext jmsContext; @ConfigProperty(name = "my.edu.queue.name", defaultValue = "my.edu.queue") String queueName; private AtomicInteger counter = new AtomicInteger(0); void onStart(@Observes StartupEvent ev) { Queue queue = jmsContext.createQueue(queueName); MessageConsumer consumer = jmsContext.createConsumer(queue); consumer.setMessageListener(message -> { try { processMessage(message.getBody(String.class)); jmsContext.commit(); } catch (Exception e) { log.severe("Error processing message: %s".formatted(e.getMessage())); jmsContext.rollback(); } }); } private void processMessage(String text) { counter.incrementAndGet(); if (counter.get() % 5 == 0) { throw new RuntimeException("Error in business logic"); } log.info("Processed message: " + text); } }
修正后的broker.xml队列配置示例
<address name="my.edu.queue"> <anycast> <queue name="my.edu.queue"> <!-- 设置为1代表最多尝试投递1次,失败后转入死信队列 --> <max-delivery-attempts>1</max-delivery-attempts> <!-- 指定死信队列地址,若不指定则使用默认DLQ --> <dead-letter-address>my.edu.queue.dlq</dead-letter-address> </queue> </anycast> </address>
完成上述配置和代码调整后,当第5条消息处理失败触发rollback时,Broker会在完成指定次数的重试后,自动将消息转入对应的死信队列。
内容的提问来源于stack exchange,提问作者joe_specimen
相关产品推荐
相关产品推荐

