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

camel-quarkus-amqp异常消息未重入队列:如何实现失败消息重排队?

问题:消息处理失败后无法重入队列导致永久丢失

现象

当消息后续处理失败时,消息不会被重入队列而是永久丢失,其他功能均正常。

环境

  • 部署在AKS上的Quarkus应用,使用camel-quarkus-amqp组件
  • 队列采用Azure Service Bus
  • 同类基于JBoss服务器的应用,相同配置(队列属性一致、Camel路由结构类似)下无此问题

代码信息

应用最初基于Quarkus 1.11.3.Final构建,升级至2.7.6.Final后问题仍存在。

pom.xml 核心配置

<properties>
  ...
  <quarkus.platform.version>2.7.6.Final</quarkus.platform.version>
</properties>
<dependencies>
  <dependency>
    <groupId>org.apache.camel.quarkus</groupId>
    <artifactId>camel-quarkus-amqp</artifactId>
  </dependency>
  <dependency>
    <groupId>org.apache.camel.quarkus</groupId>
    <artifactId>camel-quarkus-bean</artifactId>
  </dependency>
</dependencies>

application.properties 配置

quarkus.qpid-jms.url=failover:(amqps://xyz.servicebus.windows.net)
quarkus.qpid-jms.username=xxx
quarkus.qpid-jms.password=xxx

消费者路由实现

public class MessageConsumer extends RouteBuilder {

    @Inject private UpdateService updateService;

    @Override
    public void configure() {
        from("amqp:queue:" + "queueName")
            .routeId(MessageConsumer.class.getName() + ".consumeQueue")
            .process(exchange -> exchange.getIn().setBody(UUID.fromString(exchange.getIn().getBody(String.class))))
            .bean(updateService);
    }
}

生产者实现

import javax.inject.Inject;
import javax.jms.ConnectionFactory;
import javax.jms.JMSContext;
import java.util.UUID;

public class MessageRepository {

    @Inject ConnectionFactory connectionFactory;

    public void sendMessage(UUID uuid) {
        try (JMSContext context = connectionFactory.createContext(JMSContext.AUTO_ACKNOWLEDGE)) {
            context.createProducer().send(context.createQueue("queueName"), uuid.toString());
        } catch (Exception e) {
            throw new ApplicationException("Invalid message");
        }
    }
}

需求

当Camel路由中updateService的bean处理抛出异常时,让消息重入队列(不确认消息)。


解决方法

1. 启用事务模式(推荐)

在Camel的AMQP消费端点添加transacted=true,让消息处理在事务中进行。一旦updateService抛出异常,事务会自动回滚,消息会被Azure Service Bus重新放回队列:

from("amqp:queue:queueName?transacted=true")
    .routeId(MessageConsumer.class.getName() + ".consumeQueue")
    .process(exchange -> exchange.getIn().setBody(UUID.fromString(exchange.getIn().getBody(String.class))))
    .bean(updateService);

这是最可靠的方式,和JBoss应用的行为对齐——传统Java EE应用通常默认使用事务处理JMS消息。

2. 调整ACK模式并配置错误处理

如果不想用事务,可以将端点的确认模式改为CLIENT_ACKNOWLEDGE,并通过Camel的错误处理器控制异常时的消息状态:

@Override
public void configure() {
    // 捕获所有异常,标记为未处理并触发重发
    onException(Exception.class)
        .handled(false)
        .maximumRedeliveries(-1); // 无限重试(可根据需求设置具体次数)

    from("amqp:queue:queueName?ackMode=CLIENT_ACKNOWLEDGE")
        .routeId(MessageConsumer.class.getName() + ".consumeQueue")
        .process(exchange -> exchange.getIn().setBody(UUID.fromString(exchange.getIn().getBody(String.class))))
        .bean(updateService);
}

这种模式下,只有当路由完全处理成功时,Camel才会确认消息;异常时会拒绝消息,触发Service Bus的重入逻辑。

3. 检查Azure Service Bus队列配置

登录Azure门户,检查目标队列的以下配置:

  • Max Delivery Count:设置为合理的重试次数(比如5),避免消息无限循环
  • Dead-Lettering On Message Expiration:启用该选项,超过重试次数的消息会进入死信队列(DLQ),不会永久丢失,方便后续排查问题

4. 验证Qpid JMS配置

确保application.properties中没有覆盖事务或ACK相关的默认配置,如果需要显式指定事务支持,可添加:

quarkus.qpid-jms.transacted=true

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 19:35:45