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

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会自动将消息路由到死信队列。

代码优化建议

  1. 移除手动DLQ生产者:ActiveMQ Artemis会自动处理死信路由,无需手动创建Producer发送消息到DLQ,配置正确后Broker会自动完成转移。
  2. 使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 11:34:53