Java 基础体系 · 第 98/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。

Java RabbitMQ:Exchange、确认、重试、死信、顺序和幂等

RabbitMQ 生产架构中最容易被混淆的几个概念是:

  • Exchange 接收消息,但通常不保存消息;
  • Publisher Confirm 只能证明 RabbitMQ 接受了发布,不代表消费者已经处理;
  • Consumer Ack 只能证明消费者确认了一次投递,不代表业务事务一定成功;
  • 重试解决“暂时失败”,死信解决“无法继续正常处理”;
  • RabbitMQ 可以提供局部顺序,但不能自动提供全局顺序;
  • 只要存在重试、连接断开、消费者崩溃或发布端超时,就必须按“消息可能重复”设计幂等。

下面从消息流转开始,逐层建立这些机制之间的因果关系。


一、先明确 RabbitMQ 中的消息路径

RabbitMQ 中,一条消息通常经过以下路径:

flowchart LR
    P[Producer] -->|publish| E[Exchange]
    E -->|binding + routing key| Q[Queue]
    Q -->|deliver| C[Consumer]
    C -->|basic.ack| Q
    Q -->|remove after ack| D[已确认消息]

    E -->|unroutable + mandatory| R[Publisher Return]
    Q -->|reject/nack requeue=false| X[Dead Letter Exchange]
    Q -->|TTL expired / length limit| X
    X --> Q2[Retry Queue or Dead Queue]

这里有三个不同的“成功”:

  1. 发布成功:Exchange 接受了发布请求,通常通过 Publisher Confirm 判断。
  2. 路由成功:Exchange 找到至少一个匹配的 Queue;找不到时可能触发 Return。
  3. 业务成功:Consumer 完成数据库、外部服务等业务操作,并发送 Ack。

它们不是同一个状态。一个消息可以:

  • 发布成功,但因为没有绑定而未路由;
  • 路由成功,但消费者处理失败;
  • 消费者业务已提交,但 Ack 前进程崩溃,导致消息再次投递;
  • 消费者发送了 Ack,但业务事务其实没有提交,造成消息丢失。

因此,不能用 Publisher Confirm 替代 Consumer Ack,也不能把 Consumer Ack 当作数据库事务提交的自动证明。


二、Exchange:消息为什么能到达某个 Queue

2.1 Exchange 的职责

Producer 不直接把消息写入 Queue,而是把消息发布到 Exchange。Exchange 根据:

  • Exchange 类型;
  • routing key;
  • binding;
  • 可选的 headers;

决定消息应该复制到哪些 Queue。

Queue 才负责保存消息和向消费者投递。Exchange 一般是路由组件,不应理解为消息存储区。

RabbitMQ 中常见的 Exchange 类型如下。

2.2 Direct Exchange

Direct Exchange 按精确 routing key 匹配。

例如存在以下绑定:

order.exchange --[routing key = order.created]--> order.queue

只有 routing key 为 order.created 的消息才会进入 order.queue

rabbitTemplate.convertAndSend(
        "order.exchange",
        "order.created",
        event
);

Direct Exchange 适合:

  • 命令型消息;
  • 明确的事件类型;
  • 一个 routing key 对应一个或多个确定队列。

如果多个 Queue 使用同一个 routing key,它们都会收到消息。Exchange 的路由不是竞争消费,而是复制到每个匹配 Queue。

2.3 Topic Exchange

Topic Exchange 按点分隔的模式匹配 routing key。

例如:

order.created.eu
order.created.cn
order.cancelled.cn

绑定模式:

order.*.cn
order.#

其中:

  • * 匹配一个单词;
  • # 匹配零个或多个单词。

因此:

order.*.cn

可以匹配:

order.created.cn
order.cancelled.cn

但不能匹配:

order.created.eu

Topic Exchange 适合事件订阅,但 routing key 应保持稳定。如果把用户输入直接拼成 routing key,可能造成绑定失控或路由规则难以审计。

2.4 Fanout Exchange

Fanout Exchange 忽略 routing key,把消息广播给所有绑定 Queue。

event.exchange
 ├── audit.queue
 ├── analytics.queue
 └── notification.queue

每个 Queue 都会拥有一份独立消息。一个消费者 Ack 不会影响其他 Queue 中的副本。

2.5 Headers Exchange

Headers Exchange 根据消息 headers 路由,而不是主要依赖 routing key。它能表达更复杂的匹配条件,但使用和排查成本更高。除非路由条件确实不是简单的层级 key,否则 Topic 或 Direct 通常更容易维护。

2.6 Default Exchange

RabbitMQ 内置一个名称为空字符串的 Direct Exchange。Queue 声明后,RabbitMQ 会自动建立:

routing key = queue name

的隐式绑定。因此:

rabbitTemplate.convertAndSend("", "order.queue", event);

可以直接把消息发送到名为 order.queue 的 Queue。

Default Exchange 适合简单场景,但在生产架构中显式命名 Exchange 通常更清楚,因为业务路由、权限和监控都更容易表达。


三、声明 Exchange、Queue 和 Binding

下面使用 Spring Boot 和 Spring AMQP 声明一个订单消息拓扑:

@Configuration
class RabbitTopology {

    static final String ORDER_EXCHANGE = "order.x";
    static final String ORDER_QUEUE = "order.q";
    static final String ORDER_ROUTING_KEY = "order.created";

    static final String RETRY_EXCHANGE = "order.retry.x";
    static final String RETRY_QUEUE = "order.retry.10s.q";
    static final String DEAD_QUEUE = "order.dead.q";

    @Bean
    DirectExchange orderExchange() {
        return new DirectExchange(ORDER_EXCHANGE, true, false);
    }

    @Bean
    DirectExchange retryExchange() {
        return new DirectExchange(RETRY_EXCHANGE, true, false);
    }

    @Bean
    Queue orderQueue() {
        return QueueBuilder.durable(ORDER_QUEUE)
                .withArgument("x-dead-letter-exchange", RETRY_EXCHANGE)
                .withArgument("x-dead-letter-routing-key", "order.retry")
                .build();
    }

    @Bean
    Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) {
        return BindingBuilder.bind(orderQueue)
                .to(orderExchange)
                .with(ORDER_ROUTING_KEY);
    }

    @Bean
    Queue retryQueue() {
        return QueueBuilder.durable(RETRY_QUEUE)
                .withArgument("x-message-ttl", 10_000)
                .withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
                .withArgument("x-dead-letter-routing-key", ORDER_ROUTING_KEY)
                .build();
    }

    @Bean
    Binding retryBinding(Queue retryQueue, DirectExchange retryExchange) {
        return BindingBuilder.bind(retryQueue)
                .to(retryExchange)
                .with("order.retry");
    }

    @Bean
    Queue deadQueue() {
        return QueueBuilder.durable(DEAD_QUEUE).build();
    }
}

这段配置表达的路径是:

order.x
  └── order.created
      └── order.q

order.retry.x
  └── order.retry
      └── order.retry.10s.q
          └── TTL 到期后回到 order.x

最终无法处理的消息
  └── order.dead.q

需要注意几个事实:

  1. Queue 或 Exchange 已经存在时,重新声明必须使用兼容属性,否则 RabbitMQ 会返回 PRECONDITION_FAILED 并关闭通道。
  2. durable=true 只表示拓扑在 Broker 重启后保留,不表示消息一定持久化。
  3. 消息还需要设置持久化属性,通常使用 MessageDeliveryMode.PERSISTENT
  4. 生产环境中,Queue 参数经常通过 RabbitMQ Policy 管理。硬编码 x-arguments 后再修改,可能因为声明不兼容而无法滚动升级。
  5. DLX 指向的 Exchange、routing key 和目标 Queue 必须实际存在并有正确权限,否则死信可能无法按照预期路由。

四、Publisher Confirm:发布端确认了什么

4.1 Confirm 的准确含义

Publisher Confirm 是 RabbitMQ 对发布者的异步确认机制。

当 Producer 发布消息后,RabbitMQ 会在处理到相应阶段后发送:

  • ack:Broker 接受了该消息;
  • nack:Broker 没有接受该消息。

Confirm 的确认对象是“发布操作”,不是业务处理结果。

可以把发布过程抽象为:

Producer --publish--> Broker
Producer <--confirm-- Broker

当 Broker 发送 ack 时,通常表示消息已经被 RabbitMQ 接受。对于持久化消息和持久化 Queue,确认时机还受到 RabbitMQ 存储和队列类型等因素影响,但它仍不等于“消费者已经处理”。

4.2 Return 和 Confirm 解决不同问题

如果消息发布到一个不存在的 Exchange,发布通道通常会出现协议错误,不能简单等待正常 Confirm。

如果 Exchange 存在,但没有任何 Queue 能匹配 routing key,则属于“不可路由”:

Producer -> Exchange -> 没有匹配 Queue

此时只有在发布设置 mandatory=true 时,RabbitMQ 才会把消息退回 Producer,Spring AMQP 中表现为 Return。

因此,可靠发布通常同时需要:

  • Publisher Confirm:Exchange 是否接受了发布;
  • Publisher Return:消息是否没有路由到任何 Queue;
  • 对 Exchange 不存在、通道关闭等同步异常进行处理。

4.3 Spring Boot 配置

spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true
spring.rabbitmq.template.mandatory=true

correlated 表示 Confirm 能关联到具体消息。mandatory 让不可路由消息返回给发布者。

示例发布代码:

@Service
class OrderPublisher {

    private final RabbitTemplate rabbitTemplate;

    OrderPublisher(RabbitTemplate rabbitTemplate) {
        this.rabbitTemplate = rabbitTemplate;
    }

    CompletableFuture<CorrelationData.Confirm> publish(OrderCreated event) {
        CorrelationData correlationData =
                new CorrelationData(event.eventId());

        rabbitTemplate.convertAndSend(
                RabbitTopology.ORDER_EXCHANGE,
                RabbitTopology.ORDER_ROUTING_KEY,
                event,
                message -> {
                    message.getMessageProperties()
                            .setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                    message.getMessageProperties()
                            .setMessageId(event.eventId());
                    return message;
                },
                correlationData
        );

        return correlationData.getFuture();
    }
}

调用方可以等待或异步观察 Confirm:

publisher.publish(event)
        .whenComplete((confirm, error) -> {
            if (error != null) {
                // 发布调用本身失败,例如连接或通道异常
                recordPublishFailure(event.eventId(), error);
                return;
            }

            if (confirm == null || !confirm.isAck()) {
                // Broker nack,或者实现返回了无效确认
                recordPublishFailure(
                        event.eventId(),
                        new IllegalStateException(
                                confirm == null ? "no confirm" : confirm.getReason()
                        )
                );
                return;
            }

            recordPublishAccepted(event.eventId());
        });

这里的 CompletableFuture 只表示发布确认结果。它不能告诉你:

  • 是否有消费者在线;
  • 消费者是否已经收到;
  • 消费者是否提交了数据库事务;
  • 业务是否调用外部系统成功。

4.4 Confirm 丢失时怎么办

假设 Producer 发布后连接断开:

1. Producer 发送消息
2. Broker 可能已经接受
3. Confirm 尚未返回
4. Producer 连接断开
5. Producer 不知道消息究竟是否存在

如果 Producer 直接重发,可能产生重复消息;如果不重发,可能丢失消息。

这不是通过网络协议简单消除的二选一问题。工程上通常选择:

  1. 给每个事件分配稳定的 eventId
  2. 发送端记录待确认状态;
  3. Confirm ack 后标记成功;
  4. 超时或 nack 后重发;
  5. 消费端使用幂等处理消除重复影响。

这就是“可靠发布通常是至少一次发布 + 消费端幂等”,而不是试图依靠一次网络调用实现绝对的一次且仅一次。


五、Consumer Ack:RabbitMQ 什么时候删除消息

5.1 自动 Ack 和手动 Ack

消费者获取消息时,可以使用自动确认或手动确认。

自动确认的逻辑大致是:

Broker 投递消息 -> 立即视为已确认

如果消费者刚收到消息就崩溃,消息可能已经从队列中移除,业务尚未执行完成。

手动确认的逻辑是:

Broker 投递消息
    -> 消费者执行业务
    -> 业务成功后 basic.ack
    -> Broker 删除该投递对应的消息

Spring Boot 配置:

spring.rabbitmq.listener.simple.acknowledge-mode=manual
spring.rabbitmq.listener.simple.prefetch=1

消费者示例:

@Component
class OrderConsumer {

    private final OrderService orderService;

    OrderConsumer(OrderService orderService) {
        this.orderService = orderService;
    }

    @RabbitListener(queues = RabbitTopology.ORDER_QUEUE)
    public void consume(Message message, Channel channel) throws IOException {
        long deliveryTag =
                message.getMessageProperties().getDeliveryTag();

        OrderCreated event = parse(message);

        try {
            orderService.handle(event);

            // 业务成功后确认当前一条消息
            channel.basicAck(deliveryTag, false);
        } catch (TransientBusinessException e) {
            // 让消息重新入队;如果没有退避,会快速重复投递
            channel.basicNack(deliveryTag, false, true);
        } catch (PermanentBusinessException e) {
            // 不重新入队,进入 DLX 流程
            channel.basicNack(deliveryTag, false, false);
        } catch (Exception e) {
            // 默认按暂时失败处理,实际项目应根据异常分类
            channel.basicNack(deliveryTag, false, true);
        }
    }

    private OrderCreated parse(Message message) {
        // 这里可以使用 Jackson 等转换器
        throw new UnsupportedOperationException("示例省略反序列化实现");
    }
}

5.2 Ack 的作用域

RabbitMQ 的 delivery tag 是按 Channel 分配的。Ack 必须在产生该 delivery tag 的同一 Channel 上发送,否则会出现协议错误。

basicAck(deliveryTag, false) 只确认当前消息。

basicAck(deliveryTag, true) 会确认当前 tag 以及此前同一 Channel 上尚未确认的投递。批量 Ack 可以减少协议交互,但如果中间某条消息业务并未成功,就不能盲目使用 multiple=true

5.3 Ack 与业务事务的顺序

正确的基本顺序通常是:

读取消息
  -> 执行业务事务
  -> 数据库提交成功
  -> basic.ack

不能反过来:

读取消息
  -> basic.ack
  -> 数据库提交

后者在数据库提交前进程崩溃时会造成消息丢失。

但即使按正确顺序执行,也仍然存在另一个窗口:

1. 数据库事务提交
2. 进程在发送 Ack 前崩溃
3. RabbitMQ 重新投递
4. 业务再次执行

所以手动 Ack 解决的是“未成功处理时不要过早删除”,不能单独解决重复消费。


六、重试:暂时失败和永久失败必须分开

6.1 什么是暂时失败

暂时失败是指稍后重试可能成功,例如:

  • 下游 HTTP 服务临时超时;
  • 数据库连接池暂时耗尽;
  • 外部服务返回限流;
  • 某个依赖正在滚动发布。

永久失败包括:

  • 消息格式不合法;
  • 业务数据违反不可变约束;
  • 必要字段缺失;
  • 订单状态不允许当前操作。

永久失败不应无限重试,否则会形成消息热循环。

6.2 直接 requeue=true 的问题

如果消费者这样处理:

channel.basicNack(tag, false, true);

RabbitMQ 会把消息重新入队。没有任何延迟控制时,路径可能变成:

投递 -> 失败 -> 立即入队 -> 立即投递 -> 失败 -> ...

后果包括:

  • CPU 持续消耗;
  • 日志快速刷屏;
  • 同一条坏消息阻塞有效消息;
  • 下游服务被更高频率地冲击;
  • 队列吞吐下降。

因此,“重新入队”不是完整的重试策略,只是重新安排一次投递。

6.3 TTL + DLX 实现延迟重试

RabbitMQ 没有把任意消息精确睡眠后再投递给消费者的通用语义。常见做法是使用一个带 TTL 的重试 Queue:

消费失败
  -> 发布到 retry exchange
  -> retry queue 暂存 10 秒
  -> TTL 到期
  -> retry queue 的 DLX 将消息送回主 exchange
  -> 主 queue 再次消费

前面的拓扑中:

QueueBuilder.durable(RETRY_QUEUE)
        .withArgument("x-message-ttl", 10_000)
        .withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
        .withArgument("x-dead-letter-routing-key", ORDER_ROUTING_KEY)
        .build();

TTL 的变量是毫秒数。10_000 表示消息在该 Queue 中至少等待约 10 秒后才具备过期条件,但它不是实时定时器,实际重新投递还受队列调度和 Broker 状态影响。

一个重要边界是 Queue 中消息的过期处理可能受到队首消息影响。使用单一 Queue 存放不同 TTL 的消息,不能把它当作精确的时间轮。需要多个固定延迟 Queue 时,通常按 10 秒、1 分钟、10 分钟等档位分别建立拓扑,或使用 RabbitMQ 延迟消息插件;插件属于额外组件能力,不应当当作所有 RabbitMQ 部署的内置保证。

6.4 “发布到重试 Queue 后再 Ack”与失败窗口

消费者失败后,可以:

1. 将消息发布到 retry exchange
2. 等待 Publisher Confirm
3. Confirm 成功后 Ack 原消息

这比直接 nack(requeue=true) 更容易控制延迟和次数:

void retry(Message message, Channel channel, long tag, int attempt)
        throws IOException {

    RetryEvent retryEvent = new RetryEvent(
            message.getBody(),
            attempt + 1,
            message.getMessageProperties().getMessageId()
    );

    CorrelationData correlationData =
            new CorrelationData(retryEvent.eventId());

    rabbitTemplate.convertAndSend(
            RabbitTopology.RETRY_EXCHANGE,
            "order.retry",
            retryEvent,
            correlationData
    );

    correlationData.getFuture().whenComplete((confirm, error) -> {
        if (error == null && confirm != null && confirm.isAck()) {
            try {
                channel.basicAck(tag, false);
            } catch (IOException ackError) {
                // Ack 失败时,原消息可能再次投递;
                // retryEvent 必须允许重复,不能依赖“恰好一次”
                recordAckFailure(retryEvent.eventId(), ackError);
            }
        } else {
            // 重试消息未被确认,不能确认原消息
            recordRetryPublishFailure(retryEvent.eventId(), error);
        }
    });
}

这个流程仍然不是原子操作:

重试消息 Confirm 成功
原消息 Ack 失败

会导致重试消息存在两份。因此重试消息和业务处理都必须幂等。

6.5 重试次数

重试次数不能只依赖客户端内存计数,因为消费者重启后内存状态会丢失。常见方式是:

  • 在消息 header 中保存 attempt
  • 根据 RabbitMQ 的死信历史头部 x-death 推断经历过的死信次数;
  • 使用单独的重试记录表;
  • 把重试次数放入业务事件体。

Header 方案简单,但要注意消息在多次 DLX 转发过程中可能增加或更新死信历史;不同 RabbitMQ 版本和路由链路下,应在目标环境验证实际 header 结构,不应只凭日志猜测。

当达到最大次数后,消费者应将消息转入最终死信 Queue,而不是继续回到主 Queue:

attempt < 5  -> retry queue
attempt >= 5 -> dead queue

七、死信:消息为什么离开正常处理路径

7.1 死信的触发条件

消息成为死信,通常有以下原因:

  1. 消费者 basic.reject(requeue=false)
  2. 消费者 basic.nack(requeue=false)
  3. 消息在 Queue 中过期;
  4. Queue 达到长度限制后,旧消息被淘汰;
  5. 某些队列类型和 RabbitMQ 版本支持的投递次数限制触发。

消息成为死信后,RabbitMQ 会把它重新发布到该 Queue 配置的 Dead Letter Exchange。

因此:

DLX 不是一个神奇的“错误存储区”

它仍然是一次 Exchange 路由。必须为 DLX 配置目标 Queue,否则消息可能没有可见的最终落点。

7.2 消费失败进入 DLX

channel.basicReject(deliveryTag, false);

等价于拒绝当前投递且不重新入队。若原 Queue 设置了:

x-dead-letter-exchange = order.retry.x

消息将被发送到该 Exchange,再按 routing key 路由到 retry Queue 或 dead Queue。

7.3 DLX 的可靠性边界

必须区分:

  • 正常 Producer 使用 Publisher Confirm 发布;
  • RabbitMQ 内部把死信转发到 DLX。

死信转发并不自动等价于应用 Producer 的端到端 Confirm。RabbitMQ 对不同队列类型、版本和死信配置存在不同可靠性行为;例如某些版本和队列类型支持更强的 at-least-once dead-lettering,但需要显式配置并承担额外资源成本。不能仅凭“配置了 DLX”就断言死信绝不会丢失。

生产上应验证:

  1. 原 Queue 是否确实配置了 DLX;
  2. DLX Exchange 是否存在;
  3. routing key 是否匹配;
  4. DLX 目标 Queue 是否绑定;
  5. Broker 日志和权限是否正常;
  6. 目标队列是否有消息积压;
  7. RabbitMQ 版本和队列类型对死信确认的具体支持。

最终死信 Queue 不应只是“丢进去就结束”。通常还需要:

  • 保存原始消息和 headers;
  • 保存失败原因、重试次数和时间;
  • 提供人工查看;
  • 支持修复后重新投递;
  • 防止人工重放绕过幂等检查。

八、完整的状态和故障路径

一个订单事件的正常路径可以表示为:

sequenceDiagram
    participant P as Producer
    participant E as order.x
    participant Q as order.q
    participant C as Consumer
    participant DB as Database

    P->>E: publish(eventId)
    E-->>P: confirm ack
    E->>Q: route by order.created
    Q->>C: deliver(tag=17)
    C->>DB: transaction
    DB-->>C: commit
    C->>Q: basic.ack(17)
    Q-->>Q: remove message

消费者崩溃发生在数据库提交之后、Ack 之前:

sequenceDiagram
    participant Q as order.q
    participant C1 as Consumer 1
    participant DB as Database
    participant C2 as Consumer 2

    Q->>C1: deliver(eventId=E1)
    C1->>DB: commit business effect
    C1--xQ: process crashes before ack
    Q->>C2: redeliver(eventId=E1)
    C2->>DB: detect duplicate eventId
    C2->>Q: basic.ack

这条路径说明:

数据库已提交 != RabbitMQ 已 Ack

也说明幂等检查必须发生在业务操作的事务边界内,而不能只在内存中维护一个已处理集合。


九、幂等:为什么至少一次系统必须接受重复

9.1 重复消息的来源

重复可能来自多个位置:

Producer 重试

publish -> 网络中断 -> Producer 未收到 Confirm -> 重发

第一次发布可能已经成功,所以 Broker 中可能存在两条逻辑相同的消息。

Consumer Ack 丢失

业务提交 -> Ack 网络失败 -> Broker 重新投递

消费者主动重试

消息被发送到 retry Queue 后,原消息 Ack 失败,也会形成重复。

Broker 或客户端恢复

连接恢复、消费者重新注册和未确认消息重新投递,都可能让同一业务事件再次到达。

因此,所谓幂等不是“消息永远只投递一次”,而是:

同一个业务操作执行一次或多次,最终业务状态和执行一次相同。

9.2 使用唯一事件 ID

事件必须包含稳定且全局唯一的 ID:

public record OrderCreated(
        String eventId,
        String orderId,
        long occurredAt,
        BigDecimal amount
) {
}

重试和重发时必须保留原 eventId,不能每次重新生成,否则消费者无法判断它们是否属于同一个业务操作。

9.3 Inbox 表实现消费幂等

一种通用做法是建立消费记录表:

CREATE TABLE consumed_message (
    consumer_name VARCHAR(100) NOT NULL,
    event_id      VARCHAR(100) NOT NULL,
    consumed_at   TIMESTAMP NOT NULL,
    PRIMARY KEY (consumer_name, event_id)
);

处理流程:

开始数据库事务
  -> INSERT consumer_name + event_id
  -> 插入成功:继续执行业务
  -> 唯一键冲突:说明已处理,跳过业务
  -> 提交数据库事务
  -> basic.ack

伪代码如下:

@Transactional
public void handle(OrderCreated event) {
    int inserted = jdbc.update("""
        INSERT INTO consumed_message
            (consumer_name, event_id, consumed_at)
        VALUES (?, ?, CURRENT_TIMESTAMP)
        ON CONFLICT (consumer_name, event_id) DO NOTHING
        """,
        "order-consumer",
        event.eventId()
    );

    if (inserted == 0) {
        // 已经处理过
        return;
    }

    jdbc.update("""
        INSERT INTO order_projection(order_id, amount)
        VALUES (?, ?)
        ON CONFLICT (order_id) DO NOTHING
        """,
        event.orderId(),
        event.amount()
    );
}

不同数据库的冲突语法不同:

  • PostgreSQL 使用 ON CONFLICT
  • MySQL 可使用 INSERT IGNOREON DUPLICATE KEY UPDATE
  • SQL Server 可使用唯一索引配合异常处理或其他写法。

关键不是具体 SQL,而是“唯一约束和业务写入在同一个事务中”。

如果先检查、后写入:

SELECT event_id 是否存在
业务写入
INSERT event_id

两个并发消费者可能同时查询到“不存在”,然后都执行业务。唯一约束必须成为最终并发仲裁机制。

9.4 非数据库操作的幂等

如果消费者调用外部支付接口、邮件服务或 HTTP API,数据库 Inbox 不能自动撤销外部副作用。

应优先使用:

  • 外部接口支持的幂等键,例如 eventId
  • 本地操作记录表;
  • 可查询的状态机;
  • 对账和补偿机制。

例如支付请求应把 eventId 作为下游幂等键:

POST /payments
Idempotency-Key: event-8b9...

如果下游不支持幂等,而操作本身又不可撤销,就无法仅靠 RabbitMQ Ack 保证 exactly-once。此时需要重新评估接口协议和业务补偿方案。


十、顺序:RabbitMQ 能保证哪一种顺序

10.1 Queue 内的局部顺序

对于一个普通 Queue,RabbitMQ 通常按照消息进入队列的顺序投递消息。但这不是“任意情况下的全局顺序保证”。

下列因素会改变可观察顺序:

  • 一个 Queue 有多个消费者;
  • prefetch 大于 1;
  • 消费者处理时间不同;
  • 消息被 nack 后重新入队;
  • 使用消息优先级;
  • 消息经过重试 Queue 后重新返回;
  • Producer 通过多个并发 Channel 发布。

例如:

Q: A, B
Consumer 1 收到 A,处理 10 秒
Consumer 2 收到 B,处理 1 秒

业务完成顺序是:

B -> A

即使投递顺序是:

A -> B

所以“Queue 中按序”不等于“业务完成按序”。

10.2 为什么 prefetch=1 仍不等于全局顺序

单消费者、prefetch=1 能减少并发未确认消息,使投递和处理更接近串行:

spring.rabbitmq.listener.simple.concurrency=1
spring.rabbitmq.listener.simple.max-concurrency=1
spring.rabbitmq.listener.simple.prefetch=1

但它仍不能解决:

  • 失败消息重新入队后的顺序变化;
  • 重试消息晚于后续消息返回;
  • Producer 并发发布产生的实际到达顺序;
  • 消费者进程崩溃后的重新投递;
  • 多个 Queue 之间没有全局顺序。

因此它是降低并发的配置,不是分布式顺序协议。

10.3 按业务键分片

如果要求“同一订单内有序,不同订单可以并行”,应把顺序范围定义为业务聚合键,例如 orderId

目标是:

同一个 orderId -> 固定进入同一个分片 Queue
不同 orderId   -> 可以进入不同分片 Queue

可以采用:

partition = hash(orderId) mod N

Producer 根据分片结果选择 routing key 或 Queue。这样:

order-100 -> partition-2
order-101 -> partition-0
order-100 -> partition-2

同一个订单不会在多个分片之间漂移,多个分片可以并行消费。

但扩容时 N 改变会导致普通取模结果变化,使同一业务键可能进入不同分片。需要使用一致性哈希、固定分区数量,或维护稳定的分片映射。

10.4 用序列号检测乱序

仅靠 Queue 顺序不足以证明业务顺序。事件可以携带业务序列号:

public record OrderEvent(
        String eventId,
        String orderId,
        long sequence,
        String type
) {
}

消费者维护:

orderId = O1, lastSequence = 4
收到 sequence = 5 -> 应用
收到 sequence = 7 -> 暂存或拒绝,等待 6
收到 sequence = 3 -> 重复或过期,按幂等规则处理

这把“顺序”从传输层问题转化为业务状态机问题。若必须等待缺失事件,需要额外的暂存表、超时告警和人工恢复机制,否则一个丢失或永久失败的事件会阻塞整个聚合键。

10.5 重试会破坏简单顺序

考虑:

A: order.paid
B: order.shipped

如果 A 处理失败进入 10 秒重试 Queue,而 B 已经正常到达主 Queue,那么 B 可能先于 A 执行。

要保持同一订单的严格顺序,通常不能把 A 单独移出后让 B 继续执行。可选方案包括:

  1. 同一业务键串行消费,并在失败时暂停该键后续事件;
  2. 使用序列号和暂存机制;
  3. 将同一聚合键放入顺序队列,失败时阻塞后续;
  4. 让业务状态机拒绝非法状态转移并安排补偿。

“重试”和“严格顺序”存在天然张力。延迟重试越独立,越容易让后续消息越过失败消息。


十一、Outbox:解决数据库写入和消息发布之间的间隙

如果业务流程是:

数据库事务提交订单
发布 order.created

Producer 可能在两步之间崩溃:

订单已提交
消息尚未发布

此时 RabbitMQ Confirm 无法补救,因为消息甚至没有被发布。

Outbox 模式把业务数据和待发布事件写入同一个数据库事务:

CREATE TABLE outbox_event (
    event_id      VARCHAR(100) PRIMARY KEY,
    event_type    VARCHAR(100) NOT NULL,
    aggregate_id  VARCHAR(100) NOT NULL,
    payload       TEXT NOT NULL,
    status        VARCHAR(20) NOT NULL,
    created_at    TIMESTAMP NOT NULL
);

业务事务:

BEGIN
  INSERT order
  INSERT outbox_event(event_id, ...)
COMMIT

独立发布器扫描 status = 'NEW' 的事件:

读取 Outbox
  -> 发布到 RabbitMQ
  -> 等待 Publisher Confirm
  -> Confirm ack 后标记 SENT

如果发布器在 Confirm 后、更新 SENT 前崩溃,事件会再次发布。这里仍然是至少一次发布,因此消费者仍需使用 eventId 幂等。

Outbox 解决的是:

数据库提交成功,但消息没有被发布

它没有自动解决:

RabbitMQ 发布成功,但消费者处理失败

后者仍然依赖 Consumer Ack、重试、死信和幂等。


十二、不要把 RabbitMQ 事务误认为端到端事务

RabbitMQ 提供事务发布和 Confirm 等机制,但它们只能覆盖 RabbitMQ 协议范围,不能把以下操作自动绑定为一个原子事务:

数据库提交
RabbitMQ 发布
外部 HTTP 调用
RabbitMQ Ack

例如:

数据库事务提交成功
RabbitMQ 发布失败

或者:

RabbitMQ 发布成功
数据库事务回滚

都可能发生。

通常的架构取舍是:

  • 数据库写入 + 事件记录:使用 Outbox;
  • RabbitMQ 发布:使用 Publisher Confirm;
  • 消费业务 + 消费记录:使用 Inbox 和数据库事务;
  • 消费成功后:发送 Consumer Ack;
  • 外部副作用:使用下游幂等键和状态机。

这是一组组合协议,而不是 RabbitMQ 单独提供的 exactly-once 事务。


十三、一个可验证的 Spring Boot 处理骨架

依赖至少需要 Spring AMQP:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

使用 Java 25 运行时即可,但 Spring Boot 和 Spring AMQP 的具体版本必须选择官方声明支持 Java 25 的版本组合。Java 25 是运行时和编译目标,不能据此推断某个 Spring 版本自动兼容。

消费者配置:

spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=app
spring.rabbitmq.password=app

spring.rabbitmq.listener.simple.acknowledge-mode=manual
spring.rabbitmq.listener.simple.prefetch=1
spring.rabbitmq.listener.simple.default-requeue-rejected=false

default-requeue-rejected=false 只能影响 Spring 容器对异常的默认处理,不能代替对异常分类。对于需要延迟重试的异常,应明确发布到 retry Queue 或配置 Spring Retry;对于永久失败,应明确进入死信路径。

一个更完整的消费决策可以写成:

try {
    service.handle(event);       // 数据库事务中完成 Inbox + 业务写入
    channel.basicAck(tag, false);
} catch (TransientBusinessException e) {
    publishToRetryAndConfirm(event);
    channel.basicAck(tag, false);
} catch (PermanentBusinessException e) {
    channel.basicReject(tag, false);
} catch (Exception e) {
    // 不应无条件无限 requeue
    channel.basicReject(tag, false);
}

其中 publishToRetryAndConfirm 必须等待重试消息的发布确认后才能确认原消息;否则可能出现原消息已删除、重试消息却未进入 Broker 的丢失窗口。


十四、失败表现和诊断方法

14.1 Producer 显示 Confirm ack,但 Queue 没有消息

优先检查:

  1. routing key 是否正确;
  2. Exchange 和 Queue 是否有 binding;
  3. 是否启用了 mandatory
  4. 是否收到 Return;
  5. 是否把消息发布到了错误的 vhost;
  6. Exchange 类型是否与预期一致。

Confirm ack 只说明 Exchange 接受发布,不说明一定路由成功。

14.2 消费者不断重复同一条消息

检查:

  1. 是否一直 nack(requeue=true)
  2. 业务异常是否在 Ack 前抛出;
  3. 数据库事务是否实际提交;
  4. Ack 是否使用了错误 Channel;
  5. 消费者是否频繁重启;
  6. prefetch 是否过大;
  7. 是否因为外部调用超时而重复执行。

查看 RabbitMQ 管理界面中的 messages_unacknowledged、消费者数和重新投递信息,并记录:

eventId
messageId
deliveryTag
redelivered
attempt
异常类型

redelivered 可以帮助识别 Broker 重新投递,但它不是可靠的业务重试次数,也不应替代持久化重试计数。

14.3 死信队列为空

检查:

  1. 原 Queue 是否真的设置了 x-dead-letter-exchange
  2. DLX Exchange 是否存在;
  3. DLX 的 routing key 是否正确;
  4. 目标 Queue 是否绑定;
  5. 消费者是否使用了 requeue=true,导致消息根本没有进入 DLX;
  6. 是否因权限不足而转发失败;
  7. TTL 或长度限制是否真的触发。

14.4 顺序偶尔颠倒

检查:

  1. 是否有多个消费者;
  2. prefetch 是否大于 1;
  3. 是否存在 retry Queue;
  4. 是否发生 nack 和重新入队;
  5. Producer 是否并发发布;
  6. 是否有多个 Queue 或多个分片;
  7. 业务是否观察的是“完成顺序”而非“投递顺序”。

如果业务确实要求顺序,应记录 aggregateIdsequence,只看日志时间通常无法准确判断消息因果关系。


十五、关键边界和取舍

Exchange

Exchange 解决路由,不解决持久化、不解决消费确认,也不保证同一业务键的全局顺序。

Publisher Confirm

Confirm 解决发布端与 Broker 之间的确认。它不确认路由到 Queue,也不确认消费者业务成功;不可路由消息还需要 Return。

Consumer Ack

Ack 解决 Broker 是否可以删除某次投递。Ack 必须晚于业务提交,但这会引入“业务已提交、Ack 未发送”的重复窗口。

重试

重试适合暂时失败。直接 requeue=true 没有退避,容易形成热循环;TTL + DLX 能提供固定延迟,但时间精度、队首阻塞和顺序都会受到影响。

死信

死信是离开正常处理路径的消息,不是自动可靠归档。DLX 仍然依赖 Exchange、binding、权限和目标 Queue;最终死信必须支持诊断和恢复。

顺序

RabbitMQ 更适合提供 Queue 内的局部顺序,而不是跨 Queue、跨消费者和跨重试链路的全局顺序。严格顺序需要业务分片、序列号和状态机共同实现。

幂等

在至少一次发布和至少一次消费模型下,幂等是业务正确性的基础。稳定事件 ID、唯一约束、Inbox、Outbox 和下游幂等键分别解决不同故障窗口,不能相互替代。


系列导航与关联阅读

官方资料

本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。