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

Java 消息队列工程:Kafka、RabbitMQ、确认、幂等、重试和 Outbox

消息队列不是“把方法调用换成发送消息”。它改变了调用双方的时间关系、失败边界和数据一致性边界:

同步调用:

调用方 ──请求──> 服务方
调用方 <─结果── 服务方

消息调用:

调用方 ──消息──> Broker ──投递──> 消费方
调用方 <─发送结果
消费方                         └─处理结果

调用方通常只能确认“消息是否被 Broker 接收”,不能在发送时确认消费方是否已经完成业务处理。因此,消息工程的核心问题不是如何调用 send(),而是:

  1. 消息是否真的被持久化;
  2. 消息是否可能重复;
  3. 消费失败后是否重试;
  4. 重试是否会破坏业务;
  5. 数据库事务和消息发送如何保持可恢复的一致性;
  6. 系统如何在故障后继续推进,而不是静默丢失或无限循环。

本文以 Java 25 LTS 为语言运行时背景,使用 Spring Boot 与 Spring Framework 的常见编程模型说明这些机制。具体配置项会受 Spring Boot、Spring Kafka、Spring AMQP 和 Broker 版本影响,生产环境应以对应版本的官方 Reference 为准。


一、先建立消息队列的正确抽象

1. 消息、生产者、Broker、消费者

消息是一个具有业务含义的数据单元,通常包括:

消息信封:
  messageId      全局唯一消息 ID
  eventType      事件类型
  aggregateId    业务聚合 ID,例如订单 ID
  occurredAt     事件发生时间
  traceId        链路追踪 ID
  schemaVersion  消息结构版本
  payload        业务数据

其中:

  • messageId 用于识别同一条消息的重复投递;
  • aggregateId 常用于决定分区、路由键或业务顺序;
  • eventType 区分 OrderCreatedPaymentSucceeded 等事件;
  • schemaVersion 用于消息结构演进;
  • traceId 用于把 HTTP 请求、生产消息和消费消息串起来。

生产者生成消息并把它发送给 Broker。Broker负责暂存、路由、持久化和向消费者投递。消费者从 Broker 获取消息并执行本地业务。

需要特别区分三个动作:

发送请求成功
    ≠ Broker 已持久化
    ≠ 消费者已收到
    ≠ 消费者已成功处理

例如,Kafka Producer 的 send() 返回一个异步结果。这个结果通常只能在回调或 Future 完成后说明 Broker 对这次发送的确认状态。RabbitMQ 的 publisher confirm 也只确认 Broker 对发布过程的处理,不确认消费者业务已经完成。

2. 队列模型与日志模型

RabbitMQ 和 Kafka 都能实现异步消息传递,但底层模型不同。

RabbitMQ:路由到队列,再由消费者取走

RabbitMQ 常见路径是:

Producer
   │
   ▼
Exchange ──routing key──> Queue ──> Consumer

Exchange 根据类型和绑定规则路由消息:

  • direct:路由键精确匹配;
  • topic:使用通配符进行层级匹配;
  • fanout:广播到所有绑定队列;
  • headers:根据消息头匹配。

一个队列通常由多个消费者竞争消费。同一条消息在一个队列中通常只会被其中一个消费者处理;如果需要广播给多个独立业务方,应为每个业务方建立独立队列,而不是让多个业务方共享同一个队列。

Kafka:追加日志、分区和消费位点

Kafka 的逻辑结构是:

Topic
 ├── Partition 0: offset 0, 1, 2, 3 ...
 ├── Partition 1: offset 0, 1, 2, 3 ...
 └── Partition 2: offset 0, 1, 2, 3 ...

Kafka 消费者不是简单地“删除队列头部消息”,而是维护每个分区的消费位点(offset)。消息即使已经被一个消费者读取,也可以在保留期内被重新读取。

同一个 consumer group 中,一个分区在同一时刻只会分配给一个消费者实例,因此:

消费者数量 <= 分区数量

时,增加消费者通常可以提高并行度;当消费者数量超过分区数时,多出来的消费者会处于空闲状态。

Kafka 的顺序保证通常是单分区内有序,而不是整个 Topic 全局有序。若订单 order-100 的创建、支付、发货事件必须有序,应将 order-100 作为 Kafka key,使它们进入同一个分区:

kafkaTemplate.send("order-events", orderId, event);

这仍然不是绝对的业务顺序保证。消费者重试、多个下游系统的处理速度不同、数据库事务提交时间不同,都可能使“业务可见顺序”与 Kafka 的分区顺序不同。


二、确认机制:确认了什么,没有确认什么

“确认”不是单一概念。至少要区分生产确认、传输确认、消费确认和业务确认。

1. 生产者确认

生产者确认回答的问题是:

Broker 是否接受并按照相应级别处理了这条消息?

它不回答:

消费者是否已经成功执行了业务事务?

Kafka 的确认

Kafka Producer 的关键配置包括:

spring:
  kafka:
    producer:
      properties:
        acks: all
        enable.idempotence: true
        delivery.timeout.ms: 120000
        request.timeout.ms: 30000
        max.in.flight.requests.per.connection: 5

acks 的含义:

  • acks=0:生产者不等待 Broker 确认,延迟低,但可能静默丢失;
  • acks=1:Leader 确认写入,但 Leader 在复制前宕机时存在丢失风险;
  • acks=all:等待满足 ISR 条件的副本确认,可靠性更高,但延迟和可用性取舍更明显。

enable.idempotence=true 使用 Kafka Producer 的幂等发送机制,避免 Producer 因网络重试造成同一生产会话中的重复追加。它不能防止:

  • 应用层生成两条不同 messageId 的相同业务事件;
  • 应用重启后重新构造并发送业务事件;
  • 消费者重复处理;
  • Kafka 之外的数据库事务与 Kafka 发送不一致。

在 Kafka 中启用幂等生产通常要求使用兼容的 acks、重试和并发配置。不要只看到 enable.idempotence=true 就推导出“端到端 exactly-once”。

RabbitMQ 的 Publisher Confirm

RabbitMQ 发布流程至少涉及两个阶段:

Producer ──publish──> Exchange
                         │
                         ├──路由成功──> Queue
                         └──路由失败

Publisher confirm 主要确认 RabbitMQ 是否接受了发布操作;mandatory return 还可以发现消息没有路由到任何队列。

Spring Boot 中常见配置如下:

spring:
  rabbitmq:
    publisher-confirm-type: correlated
    publisher-returns: true
    template:
      mandatory: true

含义是:

  • publisher-confirm-type: correlated:让发布确认与具体消息关联;
  • publisher-returns: true:启用无法路由消息的返回通知;
  • mandatory: true:要求无法路由的消息返回生产者。

必须同时处理两种失败:

  1. confirm negative:Broker 没有确认发布成功;
  2. returned message:消息到达 Exchange,但没有匹配到队列。

只处理第一种而忽略 returned message,会出现“Exchange 接受了消息,但没有任何队列收到”的业务丢失。

Publisher confirm 仍然不等价于消费者确认。消费者可能随后因为数据库异常、进程崩溃或超时而失败。

2. 消费者确认

消费者确认回答的问题是:

Broker 是否可以认为消费者已经完成了当前消息的处理?

RabbitMQ ACK

RabbitMQ 常见消费流程:

Queue ──delivery──> Consumer
Consumer ──业务成功──> basic.ack
Consumer ──业务失败──> basic.nack/reject

Spring AMQP 手动确认配置示例:

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 20

代码示意:

@RabbitListener(queues = "order-payment-queue")
public void consume(
        OrderPaymentEvent event,
        Channel channel,
        @Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException {

    try {
        paymentService.handle(event);
        channel.basicAck(deliveryTag, false);
    } catch (RetryablePaymentException ex) {
        channel.basicNack(deliveryTag, false, true);
    } catch (NonRetryableMessageException ex) {
        channel.basicNack(deliveryTag, false, false);
    }
}

这里的三个结果不同:

  • basicAck:确认成功,消息从当前投递流程中完成;
  • basicNack(..., requeue=true):重新放回队列,可能立即再次投递;
  • basicNack(..., requeue=false):拒绝且不重新入队,是否进入死信交换机取决于队列和 DLX 配置。

如果应用在数据库提交成功后、发送 basicAck 前宕机,消息会再次投递。因此 RabbitMQ 的手动 ACK 天然要求消费者幂等。

如果应用先 ACK,再执行数据库操作,数据库失败时消息已经被确认,消息就可能丢失。正确顺序通常是:

处理业务事务成功
    ↓
事务提交成功
    ↓
ACK

但即使如此,应用仍可能在“事务提交成功”和“ACK 发出”之间崩溃,所以仍然需要幂等。

Kafka offset 提交

Kafka 的消费确认本质上是提交 offset:

处理 record
    ↓
数据库事务提交
    ↓
提交 offset

若先提交 offset,再提交数据库事务,数据库失败后 Kafka 不会再次投递该消息,形成业务丢失。

Spring Kafka 的手动确认示意:

@KafkaListener(
        topics = "order-events",
        groupId = "inventory-service"
)
public void consume(
        ConsumerRecord<String, OrderEvent> record,
        Acknowledgment acknowledgment) {

    inventoryService.handle(record.value());

    // 只有 handle 返回且事务已成功时才确认
    acknowledgment.acknowledge();
}

这里还要看容器的 AckMode、事务配置和异常处理器。acknowledgment.acknowledge() 并不一定意味着 offset 已经同步写入 Kafka;它通常是向 Spring Kafka 容器表达“该记录可以提交”,最终提交时机由容器配置决定。

因此,排查“Kafka 消息是否丢失”时不能只看 Listener 方法有没有执行,还要看:

  • offset 是否在业务前提交;
  • 异常是否被 Listener 方法吞掉;
  • 错误处理器是否跳过了记录;
  • consumer group 是否发生 rebalance;
  • 处理线程是否超过 max.poll.interval.ms
  • 生产者和消费者是否使用了不同的 key 或 group。

三、消息投递语义:至少一次、至多一次和所谓 Exactly Once

1. 至多一次

至多一次通常是:

先确认消息已消费
    ↓
再执行业务

如果业务处理失败,消息不会再投递,因此可能丢失,但通常不会重复处理。

适合少量可丢失、重复代价高的场景,例如某些非关键监控采样。但不适合作为订单、支付、库存等核心事务的默认语义。

2. 至少一次

至少一次通常是:

先执行业务
    ↓
业务成功后确认

如果进程在业务成功与确认之间崩溃,消息会再次投递。于是:

至少一次 = 可能重复,但尽量不丢

大多数 RabbitMQ 手动 ACK 和 Kafka 手动提交 offset 的业务消费,实际都需要按至少一次来设计。

3. Exactly Once 的边界

“Exactly Once”必须说明作用范围。

Kafka 事务可以在 Kafka 内部实现类似:

读取 Kafka 消息
    ↓
处理
    ↓
向 Kafka 写入结果 + 提交消费 offset

这使“Kafka 输入到 Kafka 输出”可以形成一个事务边界,消费者使用 read_committed 时不会看到未提交事务。但如果处理过程还更新了 MySQL:

Kafka 输入
   ├── MySQL UPDATE
   └── Kafka 输出与 offset

Kafka 事务不能自动回滚 MySQL 的 UPDATE。因此它不是 Kafka、MySQL、Redis、HTTP 外部调用之间的全局 exactly-once。

更准确的工程目标通常是:

消息至少一次投递
+ 业务处理幂等
+ 失败可重试
+ 不可恢复消息可隔离

这比笼统地宣称“Exactly Once”更可验证。


四、幂等:把重复执行变成同一个结果

1. 幂等的形式化定义

设业务操作为函数:

f(state, message)

如果同一条消息执行一次和执行多次,对最终业务状态的影响相同,则称该处理在该消息语义下幂等:

f(f(state, m), m) = f(state, m)

这里的 m 必须代表同一个业务事实,而不是“字段刚好相同的两条消息”。

例如:

设置订单状态为 PAID

通常可以设计为幂等:

UPDATE orders
SET status = 'PAID'
WHERE order_id = ?
  AND status = 'CREATED';

第一次执行影响 1 行,第二次执行影响 0 行,但最终状态仍是 PAID

相反:

库存数量 = 库存数量 - 1

直接重复执行两次会扣减两次,不是幂等操作。

2. 幂等键与业务键

消息幂等通常需要一个稳定的唯一键:

  • 事件 ID:表示同一条事件;
  • 命令 ID:表示同一次业务请求;
  • 业务操作 ID:例如 paymentId
  • HTTP Idempotency-Key:表示同一次客户端请求。

不要只使用消息的时间戳或随机消费实例 ID。消费者重启后必须仍能识别这是同一条消息。

3. Inbox 表:用数据库唯一约束做去重

一个常见方案是建立 Inbox 表:

CREATE TABLE inbox_message (
    message_id   VARCHAR(128) PRIMARY KEY,
    consumer     VARCHAR(128) NOT NULL,
    received_at  TIMESTAMP NOT NULL,
    processed_at TIMESTAMP NULL
);

消费者收到消息后,与业务更新放在同一个数据库事务中:

BEGIN;

INSERT INTO inbox_message(message_id, consumer, received_at)
VALUES (:messageId, :consumer, CURRENT_TIMESTAMP);

-- 如果 message_id 已存在,则说明该消费者已经处理过
UPDATE orders
SET status = 'PAID',
    paid_at = CURRENT_TIMESTAMP
WHERE order_id = :orderId
  AND status = 'CREATED';

UPDATE inventory
SET available = available - :quantity
WHERE sku = :sku
  AND available >= :quantity;

UPDATE inbox_message
SET processed_at = CURRENT_TIMESTAMP
WHERE message_id = :messageId
  AND consumer = :consumer;

COMMIT;

关键点是:

  1. message_id 由数据库唯一约束保护;
  2. Inbox 插入和业务更新必须在同一个事务;
  3. 业务失败时整个事务回滚,Inbox 记录也回滚,下一次仍可重试;
  4. 重复消息到达时,唯一键冲突表示该消息已经成功处理,可以安全确认;
  5. 业务更新还应使用状态条件或版本条件,防止并发覆盖。

如果 Inbox 插入成功后业务更新失败,但二者不在同一事务中,下一次消息会被误判为重复,造成业务丢失。这是实现 Inbox 时最危险的错误之一。

4. 用唯一业务约束代替专门去重表

某些业务可以直接使用唯一索引:

CREATE UNIQUE INDEX uq_payment_order
ON payments(order_id);

消费 PaymentSucceeded 时执行插入:

INSERT INTO payment_records(order_id, transaction_id, amount)
VALUES (:orderId, :transactionId, :amount);

第二次插入触发唯一约束,说明支付记录已存在。这里的唯一约束不仅是优化,而是幂等正确性的组成部分。

5. 幂等不是“忽略所有重复”

重复消息可能携带不同内容。例如同一订单先收到:

OrderStatusChanged(orderId=1, status=PAID, version=5)

又收到:

OrderStatusChanged(orderId=1, status=REFUNDED, version=6)

不能因为 orderId=1 已经处理过,就把第二条消息忽略。正确的去重键应是事件 ID,正确的顺序条件可能是版本号:

UPDATE order_projection
SET status = :status,
    version = :version
WHERE order_id = :orderId
  AND version < :version;

这同时处理了:

  • 同一事件重复投递;
  • 不同事件乱序到达;
  • 旧事件覆盖新状态。

五、重试:根据故障类型决定是否再次投递

重试不是“失败后无限重新消费”。重试机制必须回答三个问题:

  1. 这次失败是否可能自行恢复;
  2. 下一次什么时候重试;
  3. 重试仍失败后放到哪里。

1. 可重试错误与不可重试错误

常见可重试错误:

  • 临时数据库连接失败;
  • 下游服务超时;
  • Redis 短暂不可用;
  • 网络连接重置;
  • 对方返回限流或暂时不可用;
  • Broker 临时故障。

常见不可重试错误:

  • 消息 JSON 无法解析;
  • 必填字段缺失;
  • 签名校验失败;
  • 业务状态非法;
  • 商品不存在且不会自动创建;
  • 金额格式不合法;
  • 代码缺陷导致的确定性异常。

把不可重试错误无限重试,会形成毒丸消息(poison message):它持续占用消费者,拖慢或阻塞后续消息。

2. 指数退避

设第 n 次重试的等待时间为:

delay(n) = min(cap, base × 2^(n-1)) + jitter

变量含义:

  • base:第一次等待时间,例如 1 秒;
  • cap:最大等待时间,例如 5 分钟;
  • n:重试次数,从 1 开始;
  • jitter:随机抖动,用来避免大量消费者同时重试。

例如 base=1scap=60s,不计算抖动时:

第 1 次:1 秒
第 2 次:2 秒
第 3 次:4 秒
第 4 次:8 秒
第 5 次:16 秒
第 6 次:32 秒
第 7 次:60 秒

如果所有消费者在相同时间启动且没有 jitter,数据库恢复瞬间可能承受一轮同步重试流量,造成再次过载。

3. RabbitMQ 重试的两种方式

方式一:原队列重新入队

channel.basicNack(deliveryTag, false, true);

优点是简单;问题是消息可能立刻重新投递:

取出消息
  ↓
失败
  ↓
立即 requeue
  ↓
再次取出
  ↓
再次失败

这会造成 CPU 空转、日志爆炸和队列头部阻塞,不适合长时间退避。

方式二:重试队列加 TTL 和死信交换机

一种拓扑是:

业务队列
   │失败
   ▼
retry-5s 队列 --TTL 5s--> retry-exchange --> 业务队列
   │
   └─超过次数--> dead-letter 队列

概念性配置示例:

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: manual
    template:
      mandatory: true

RabbitMQ 的具体 DLX、TTL、队列参数可以通过声明 Bean 或基础设施配置完成。需要注意:

  • TTL 到期后消息不是“定时器精确触发”,而是由队列内部机制转入死信路径;
  • 多条消息共享一个 retry 队列时,前面的长延迟消息可能影响后续消息可见性;
  • 重新发布到 retry 队列时必须保留原始 messageId 和重试次数;
  • 重试消息重新发布也需要处理 publisher confirm,否则可能在失败转移时丢失。

4. Kafka 重试的实现取舍

Kafka 中不建议简单地在 Listener 线程里 sleep(60_000)

  • 消费线程长时间不调用 poll()
  • 可能超过 max.poll.interval.ms
  • Broker 认为消费者失活并触发 rebalance;
  • 一个失败分区可能阻塞同一消费者负责的其他分区。

常见方案有:

  1. 使用 Spring Kafka 的错误处理器和重试主题;
  2. 将失败消息发布到 topic.retry.1mtopic.retry.10m 等不同 Topic;
  3. 使用独立重试消费者;
  4. 最终写入 DLT(Dead Letter Topic)。

重试 Topic 通常要携带:

originalTopic
originalPartition
originalOffset
originalMessageId
retryCount
firstFailedAt
lastErrorType
lastErrorMessage

重新投递时必须考虑分区顺序。将同一业务 key 重新发送到原 Topic 时,应继续使用相同 key;否则重试消息可能进入其他分区,破坏单 key 的顺序假设。

5. 重试必须具备上限

设最大重试次数为 N

attempt <= N

当第 N 次仍然失败时,消息应转移到隔离区域:

业务队列/Topic
    ↓
有限次重试
    ↓
DLQ/DLT

死信消息不能只是一个“垃圾桶”。至少需要:

  • 原始消息;
  • 失败异常类型;
  • 最后一次错误信息;
  • 重试次数;
  • 首次失败时间;
  • 最近失败时间;
  • 消费者名称;
  • Trace ID;
  • 原始 Topic、分区、offset 或 RabbitMQ routing key。

修复代码或下游依赖后,可以执行人工重放。重放前必须确认幂等逻辑仍然有效,否则死信恢复本身会引入重复业务。


六、Outbox:解决“数据库提交成功但消息发送失败”

1. 双写问题

以下代码存在经典双写风险:

@Transactional
public void createOrder(CreateOrderCommand command) {
    orderRepository.insert(command.toOrder());

    kafkaTemplate.send("order-events",
            new OrderCreated(command.orderId()));
}

数据库事务与 Kafka 发送不属于同一个本地事务。可能出现:

情况 A:
数据库提交成功
Kafka 发送失败
=> 订单存在,但下游永远不知道

情况 B:
Kafka 发送成功
数据库事务回滚
=> 下游收到一个实际不存在的订单

情况 C:
Kafka 发送结果未知
应用重试发送
=> 可能产生重复事件

这不是 Spring 的 @Transactional 失效,而是事务资源不同:

MySQL 事务:保护 MySQL
Kafka 发送:保护 Kafka
两者默认没有共同提交点

2. Outbox 的定义

Outbox Pattern 的做法是:

不在业务事务中直接依赖消息发送,而是在同一个数据库事务中写入业务数据和待发送事件;随后由独立 Relay 将 Outbox 事件发布到 Broker。

数据流变成:

业务请求
   │
   ▼
数据库本地事务
   ├──业务表
   └──outbox_event
         │
         ▼
      Relay
         │
         ▼
      Kafka/RabbitMQ

业务表和 Outbox 表在同一个数据库中,所以可以获得本地事务原子性:

业务数据提交成功
    ⇔
Outbox 记录提交成功

这里的 只表示数据库内的两个写入一起提交或一起回滚,不表示消息已经发布到 Broker。

3. Outbox 表设计

一个可用的基础表结构:

CREATE TABLE outbox_event (
    id              BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    event_id        VARCHAR(128) NOT NULL UNIQUE,
    aggregate_type  VARCHAR(128) NOT NULL,
    aggregate_id    VARCHAR(128) NOT NULL,
    event_type      VARCHAR(128) NOT NULL,
    schema_version  INT NOT NULL,
    payload         TEXT NOT NULL,
    status          VARCHAR(20) NOT NULL,
    attempts        INT NOT NULL DEFAULT 0,
    available_at    TIMESTAMP NOT NULL,
    locked_until    TIMESTAMP NULL,
    last_error      TEXT NULL,
    created_at      TIMESTAMP NOT NULL,
    sent_at         TIMESTAMP NULL
);

CREATE INDEX idx_outbox_ready
ON outbox_event(status, available_at);

event_id 必须具有稳定唯一性。status 可以包含:

READY       待发送
PROCESSING  已被 Relay 租约占用
SENT        已得到发布确认
FAILED      暂时失败,等待下次调度
DEAD        超过重试上限,等待人工处理

如果使用 PostgreSQL,payload 可以使用 JSONB;如果使用 MySQL,可以使用 JSON。消息体存储格式和序列化策略应固定,避免升级 Java 类后无法反序列化旧事件。

4. 业务事务中的 Outbox 写入

使用 Spring JDBC 的示意代码:

@Service
public class OrderApplicationService {

    private final JdbcTemplate jdbc;

    public OrderApplicationService(JdbcTemplate jdbc) {
        this.jdbc = jdbc;
    }

    @Transactional
    public void createOrder(String orderId, String customerId) {
        jdbc.update("""
            INSERT INTO orders(order_id, customer_id, status, created_at)
            VALUES (?, ?, 'CREATED', CURRENT_TIMESTAMP)
            """,
            orderId, customerId);

        String eventId = UUID.randomUUID().toString();
        String payload = """
            {"orderId":"%s","customerId":"%s"}
            """.formatted(orderId, customerId);

        jdbc.update("""
            INSERT INTO outbox_event(
                event_id, aggregate_type, aggregate_id,
                event_type, schema_version, payload,
                status, available_at, created_at
            )
            VALUES (?, 'Order', ?, 'OrderCreated', 1, ?,
                    'READY', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
            """,
            eventId, orderId, payload);
    }
}

这段代码成立的前提是:

  • ordersoutbox_event 使用同一个数据库连接和事务管理器;
  • @Transactional 的调用经过 Spring 代理;
  • createOrder 不是同一个类内部的直接自调用;
  • 业务异常能传播出去,使事务回滚。

如果 orderRepository 使用 MySQL,而 Outbox 使用另一个数据库,仍然不是一个本地事务。若确实需要跨资源原子提交,就会进入 XA、分布式事务或其他协调机制的复杂取舍,不能仅靠 @Transactional 解决。

5. Relay 的状态变化

Relay 不能简单写成:

for (OutboxEvent event : findReady()) {
    kafkaTemplate.send(event.payload());
    markSent(event.id());
}

因为发布结果是异步的,且进程可能在任意位置崩溃。一个更可恢复的状态流程是:

READY
  │领取并设置租约
  ▼
PROCESSING
  │
  ├──发布确认成功──> SENT
  │
  ├──临时失败──────> READY,available_at 延后
  │
  └──租约超时──────> 可被其他 Relay 重新领取

领取时需要避免多个 Relay 重复处理同一批记录。PostgreSQL 示例:

BEGIN;

SELECT id, event_id, aggregate_id, payload, attempts
FROM outbox_event
WHERE status = 'READY'
  AND available_at <= CURRENT_TIMESTAMP
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 100;

UPDATE outbox_event
SET status = 'PROCESSING',
    locked_until = CURRENT_TIMESTAMP + INTERVAL '2 minutes',
    attempts = attempts + 1
WHERE id IN (...);

COMMIT;

SKIP LOCKED 的作用是:一个 Relay 已经锁住的行,其他 Relay 跳过它们,继续领取其他行。它不是所有数据库都以相同方式支持,SQL 语法和锁行为必须按实际数据库验证。

发布成功后:

UPDATE outbox_event
SET status = 'SENT',
    sent_at = CURRENT_TIMESTAMP,
    locked_until = NULL
WHERE id = :id
  AND status = 'PROCESSING';

发布失败后:

UPDATE outbox_event
SET status = CASE
        WHEN attempts >= :maxAttempts THEN 'DEAD'
        ELSE 'READY'
    END,
    available_at = :nextAvailableAt,
    locked_until = NULL,
    last_error = :error
WHERE id = :id
  AND status = 'PROCESSING';

6. Outbox 为什么仍然会重复发送

关键故障窗口如下:

1. Relay 从数据库领取事件
2. Relay 成功发布到 Kafka
3. Kafka 已确认
4. Relay 在标记 SENT 前宕机

恢复后,outbox_event 仍是 PROCESSING,租约超时后会再次发布。因此 Outbox 的典型语义是:

业务数据与待发送事件不丢失
但同一事件可能被发布多次

这不是 Outbox 的失败,而是其故障恢复模型。解决办法不是把 SENT 提前写入,因为:

先写 SENT
  ↓
再发送 Broker
  ↓
进程宕机
=> 数据库认为已发送,但 Broker 实际没有消息

正确取舍是允许重复发布,再由消费者以 event_id 幂等处理。Outbox 解决的是“数据库提交成功但消息没有可靠进入发送流程”,不单独解决端到端 exactly-once。


七、一个完整的 Java 消费者实现思路

下面以订单事件消费者为例,说明“数据库幂等 + 手动确认 + 异常分类”的组合。

1. 数据库表

CREATE TABLE inbox_message (
    message_id  VARCHAR(128) NOT NULL,
    consumer    VARCHAR(128) NOT NULL,
    received_at TIMESTAMP NOT NULL,
    processed_at TIMESTAMP NULL,
    PRIMARY KEY (message_id, consumer)
);

CREATE TABLE inventory (
    sku       VARCHAR(128) PRIMARY KEY,
    available INT NOT NULL
);

同一事件被库存服务处理过,不代表支付服务也处理过,因此主键包含 consumer

2. 业务服务

@Service
public class InventoryService {

    private final JdbcTemplate jdbc;

    public InventoryService(JdbcTemplate jdbc) {
        this.jdbc = jdbc;
    }

    @Transactional
    public void handle(OrderCreatedEvent event) {
        int inserted = jdbc.update("""
            INSERT INTO inbox_message(
                message_id, consumer, received_at
            )
            VALUES (?, ?, CURRENT_TIMESTAMP)
            ON CONFLICT (message_id, consumer) DO NOTHING
            """,
            event.messageId(), "inventory-service");

        if (inserted == 0) {
            // 已在同一消费者中成功处理过
            return;
        }

        int updated = jdbc.update("""
            UPDATE inventory
            SET available = available - ?
            WHERE sku = ?
              AND available >= ?
            """,
            event.quantity(),
            event.sku(),
            event.quantity());

        if (updated != 1) {
            throw new InsufficientInventoryException(event.sku());
        }

        jdbc.update("""
            UPDATE inbox_message
            SET processed_at = CURRENT_TIMESTAMP
            WHERE message_id = ?
              AND consumer = ?
            """,
            event.messageId(), "inventory-service");
    }
}

这里的 ON CONFLICT 是 PostgreSQL 语法。MySQL 可以使用 INSERT IGNOREINSERT ... ON DUPLICATE KEY UPDATE,但应明确检查影响行数语义,不能未经验证直接替换。

完整消费路径:

@Component
public class OrderEventListener {

    private final InventoryService inventoryService;

    public OrderEventListener(InventoryService inventoryService) {
        this.inventoryService = inventoryService;
    }

    @KafkaListener(
            topics = "order-events",
            groupId = "inventory-service"
    )
    public void onMessage(
            OrderCreatedEvent event,
            Acknowledgment acknowledgment) {

        try {
            inventoryService.handle(event);
            acknowledgment.acknowledge();
        } catch (InsufficientInventoryException ex) {
            // 业务上通常不可通过立即重试解决
            throw new NonRetryableMessageException(ex);
        } catch (TransientDatabaseException ex) {
            // 抛出给容器错误处理器,进入重试路径
            throw ex;
        }
    }
}

这里必须保证:

InventoryService.handle() 的事务提交成功
    ↓
Listener 确认 offset

如果 handle() 抛出异常,监听器不应捕获后直接返回,否则容器可能把消息视为成功处理。

3. 重复消息的完整执行结果

假设同一个 eventId=evt-1 到达两次:

第一次:

INSERT inbox:成功,影响 1 行
UPDATE inventory:available 10 -> 8
事务提交
ACK/提交 offset

第二次:

INSERT inbox:冲突,影响 0 行
直接返回
ACK/提交 offset

最终库存只扣减 2,而不是 4。

如果第一次执行过程中数据库事务回滚:

INSERT inbox:回滚
UPDATE inventory:回滚
不 ACK

下一次重新投递时,Inbox 仍然不存在,消息可以重新执行。这就是 Inbox 必须和业务变更处于同一事务的原因。


八、Kafka 工程要点:分区、提交、再均衡和事务

1. 分区决定并行度与局部顺序

Kafka 的并行处理单位是分区,不是消息。对于一个 Topic:

P = 分区数
C = consumer group 中消费者数
有效并行消费者数 = min(P, C)

C > P 时,多出来的消费者不会增加该 Topic 的并发处理能力。

选择 key 时,要同时考虑:

  • 同一业务实体是否需要顺序;
  • key 是否分布均匀;
  • 是否会出现热点 key;
  • 未来是否需要扩容分区。

一个错误例子是把所有消息都使用固定 key:

kafkaTemplate.send("order-events", "all-orders", event);

这样所有消息进入同一分区,Topic 虽然配置了很多分区,实际仍接近单线程处理。

2. 长事务与 max.poll.interval.ms

Kafka 消费者需要持续调用 poll()。如果单条消息处理时间超过 max.poll.interval.ms,消费者可能被认为失活并触发 rebalance。

例如:

max.poll.interval.ms = 5 分钟
单条消息业务处理 = 8 分钟

即使业务最终成功,消费者也可能在处理中被移出 group。其他消费者重新领取分区后,同一消息可能被再次处理。

解决方向不是盲目增大参数,而是先判断:

  • 是否可以拆分长任务;
  • 是否需要把任务状态写入数据库后异步推进;
  • 是否应限制 max.poll.records
  • 是否应使用独立的任务执行器;
  • 执行器中的业务是否仍然能正确控制 offset 提交。

如果把消息交给线程池后立即提交 offset:

poll
  ↓
submit 到线程池
  ↓
立即 commit offset
  ↓
线程池任务失败

就会出现消息丢失。异步处理必须建立“任务完成后再提交 offset”的明确协调机制,并控制线程池队列容量,否则会造成内存堆积。

3. Kafka 事务的适用范围

Kafka 事务适合:

Kafka Consumer
  ├──读取输入
  ├──处理
  ├──写入 Kafka 输出
  └──提交输入 offset

例如流式转换:

orders-topic ──消费──> payment-events-topic

若输出写入 Kafka 事务中,消费者可以设置:

spring:
  kafka:
    consumer:
      properties:
        isolation.level: read_committed

这样不会读取未提交事务中的消息。

但以下流程仍需额外设计:

Kafka 输入
   ├──更新 MySQL
   └──写 Kafka 输出

即使 Kafka 输出和 offset 在 Kafka 事务中,MySQL 更新也不会自动回滚。需要使用 Outbox、Inbox,或明确接受最终一致性。


九、RabbitMQ 工程要点:路由、预取和消息生命周期

1. Exchange、Queue 与 Routing Key

生产者通常不直接面向业务队列,而是发布到 Exchange:

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

RabbitMQ 根据绑定关系决定消息进入哪些队列。生产者如果把队列名和所有消费者拓扑都硬编码在业务代码中,会使路由变更成本很高。

一个广播场景:

order.exchange (fanout)
   ├── inventory.queue
   ├── notification.queue
   └── analytics.queue

库存、通知和分析都应各自拥有独立队列。若三者共享一个队列,则一条订单事件只会被其中一个消费者拿到,不是广播。

2. Prefetch 与未确认消息

prefetch 限制消费者同时持有但尚未确认的消息数。它影响:

  • 消费者吞吐;
  • 内存占用;
  • 消息分配公平性;
  • 单个消费者故障时的重投递批量。

例如:

spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 20

设置过大时,一个慢消费者可能预取大量消息,其他消费者空闲,故障时还会产生大量重新投递。设置过小时,吞吐可能下降。它不是越大越好,应结合单条处理时间、消息大小和消费者数量压测。

3. RabbitMQ 的消息持久化边界

要降低 Broker 重启造成的消息丢失,通常需要同时考虑:

  • Exchange 是否持久化;
  • Queue 是否持久化;
  • Message 是否设置为 persistent;
  • Broker 集群和存储配置;
  • Publisher confirm 是否等待成功。

只把消息标记为 persistent,而队列是临时队列,不能获得期望的持久性。只声明 durable queue,但发布没有等待 confirm,也不能知道 Broker 是否接受成功。


十、Outbox Relay 的并发、顺序与删除策略

1. 多 Relay 并发

多个 Relay 可以提高吞吐,但要避免同一行被重复领取。常见方案:

  • 数据库行锁加 SKIP LOCKED
  • 原子 UPDATE ... WHERE status='READY'
  • 分片领取,例如按 id % N
  • 使用独立任务调度系统。

租约机制必须处理 Relay 崩溃:

PROCESSING + locked_until < now()
    => 允许重新进入 READY

不要永久保留 PROCESSING 状态,否则一次进程崩溃就会让事件永远卡住。

2. 同一聚合的发布顺序

如果订单事件必须按创建、支付、发货顺序发布,单纯按 Outbox 自增主键排序并不足够:

  • 多个事务的提交顺序可能不同;
  • 多个 Relay 并行发布;
  • Kafka 不同分区之间没有全局顺序;
  • RabbitMQ 多消费者也可能并行处理。

可行方法包括:

  1. 使用 aggregate_id 作为 Kafka key;
  2. 在 Outbox 中保存聚合版本;
  3. 消费者拒绝或延迟处理版本跳跃的事件;
  4. 对同一聚合使用单线程或分区内顺序消费。

例如:

OrderCreated version=1
PaymentSucceeded version=2
OrderShipped version=3

消费者收到 version=3,但数据库当前版本是 1 时,不能简单地把 version=3 当作最新状态写入,除非业务明确允许跳过中间事实。

3. Outbox 清理

SENT 记录不能无限增长。清理前需要确认:

  • Broker 保留期是否已覆盖需要的重放窗口;
  • 是否有审计要求;
  • 是否有下游仍依赖 Outbox 重放;
  • 是否需要归档到对象存储或历史库。

一种策略是:

DELETE FROM outbox_event
WHERE status = 'SENT'
  AND sent_at < CURRENT_TIMESTAMP - INTERVAL '14 days';

清理任务应分批执行,避免长事务锁住大量行。若使用主从复制、CDC 或审计系统,还要确认清理不会破坏下游同步。


十一、消息 Schema、版本和兼容性

消息一旦进入 Broker,就可能在数小时、数天甚至更久后被消费。生产者和消费者不一定同时发布,因此消息结构必须考虑演进。

1. 兼容性规则

从旧消费者角度,新增字段通常比删除字段安全:

{
  "orderId": "o-1",
  "customerId": "c-1",
  "coupon": "NEW"
}

旧消费者忽略未知字段通常可以继续工作。

删除字段、修改字段类型、改变字段语义则可能破坏消费者:

amount: 100

改成:

amount: {"value":100,"currency":"CNY"}

这不是普通字段新增,而是结构不兼容。

2. 事件与命令的区别

  • 事件:某件事已经发生,例如 OrderCreated
  • 命令:要求对方执行某件事,例如 ReserveInventory

事件通常由事实产生者发布,多个消费者独立订阅。命令通常具有明确目标和执行责任。

如果把命令错误地广播给多个实例,可能导致多个实例重复执行同一动作。应通过队列、消费者组或业务幂等键明确命令的归属。


十二、失败路径:从发送到处理的完整时序

下面的时序展示了一个使用 Outbox、Kafka 和 Inbox 的典型流程:

sequenceDiagram
    participant C as HTTP Client
    participant A as Order Service
    participant DB as Order DB
    participant R as Outbox Relay
    participant K as Kafka
    participant I as Inventory Service
    participant IDB as Inventory DB

    C->>A: POST /orders + Idempotency-Key
    A->>DB: BEGIN
    A->>DB: INSERT orders
    A->>DB: INSERT outbox_event(event_id)
    A->>DB: COMMIT
    A-->>C: 201 Created

    R->>DB: 领取 READY 事件并设置租约
    R->>K: publish(event_id, orderId)
    K-->>R: producer ack
    R->>DB: 标记 SENT

    K->>I: deliver record
    I->>IDB: BEGIN
    I->>IDB: INSERT inbox(event_id)
    I->>IDB: 扣减库存
    I->>IDB: COMMIT
    I->>K: commit offset

关键故障点如下:

故障点一:业务事务提交前进程崩溃

orders 未提交
outbox 未提交

客户端可以重试 HTTP 请求,但必须使用相同的 Idempotency-Key,否则可能创建两个订单。

故障点二:业务事务提交后 Relay 崩溃

orders 已提交
outbox=READY 或 PROCESSING

Relay 恢复后继续领取,消息不会因为 HTTP 请求已经返回而丢失。

故障点三:Kafka 发布成功后 Relay 崩溃

Kafka 已有事件
outbox 尚未标记 SENT

Relay 可能重复发布。消费者的 Inbox 表保证库存服务只执行一次。

故障点四:库存事务提交后消费者崩溃

库存已扣减
offset 尚未提交

消息再次投递,Inbox 唯一键阻止第二次扣减,然后消费者确认 offset。

故障点五:消息结构损坏

反序列化失败

如果错误处理器继续无限重试,分区或队列可能持续阻塞。此类消息应进入 DLT/DLQ,并保留原始内容供诊断。


十三、HTTP 幂等与消息幂等要连起来

消息幂等不能替代入口请求幂等。

例如客户端调用:

POST /orders
Idempotency-Key: req-abc-123

服务端应保存请求键与响应结果:

CREATE TABLE http_idempotency (
    idempotency_key VARCHAR(128) PRIMARY KEY,
    request_hash    VARCHAR(128) NOT NULL,
    status           VARCHAR(20) NOT NULL,
    response_body   TEXT NULL,
    created_at      TIMESTAMP NOT NULL
);

同一个 key 再次到达时:

  • 请求内容相同:返回之前的结果;
  • 请求内容不同:返回冲突错误;
  • 第一次请求仍处于处理中:根据协议返回处理中或等待结果。

如果只在消息消费者端去重,HTTP 请求已经可能创建两个不同订单:

第一次请求 -> 订单 A -> 事件 A
客户端超时
第二次请求 -> 订单 B -> 事件 B

消费者无法判断 A 和 B 是否是同一次用户意图,除非业务模型中有更上层的请求幂等键约束。


十四、Redis、Elasticsearch 与消息的一致性边界

消息驱动系统常把 Redis 作为缓存、Elasticsearch 作为检索投影。这两者都不应直接被当作业务事实的唯一来源。

1. 数据库与缓存

如果订单数据库提交成功,缓存删除消息发布失败:

数据库:新状态
Redis:旧状态

这属于缓存最终一致性问题。可以通过 Outbox 发布缓存失效事件:

数据库事务
  ├──更新订单
  └──写 outbox(OrderUpdated)

Relay
  └──发布事件

Cache Consumer
  └──删除或刷新 Redis

Redis 中的“已处理 messageId”可以作为性能优化,但不应单独作为可靠幂等依据:

  • Redis key 可能过期;
  • Redis 可能故障;
  • 写 Redis 成功而数据库失败;
  • 多区域或故障恢复时可能丢失。

核心幂等记录仍应由具有持久唯一约束的数据库或可靠事件存储保护。

2. Elasticsearch 投影

Elasticsearch 通常是数据库事实的异步投影:

MySQL 事实表
    ↓ Outbox/Event
消息队列
    ↓
ES Projection Consumer
    ↓
Elasticsearch 文档

ES 消费者必须支持重复事件和乱序事件。常见做法是把数据库版本写入 ES,并使用版本条件或在消费者侧比较版本:

事件 version=10 到达 -> 写入
事件 version=9 到达  -> 丢弃
事件 version=10 重复 -> 幂等覆盖或忽略

如果搜索结果暂时落后,查询接口需要明确这是最终一致性,而不是在所有路径上假设 ES 立即反映数据库提交。


十五、常见错误与失败表现

错误一:发送成功就认为业务成功

失败表现:

Producer 收到 confirm
但消费者一直没有处理

诊断:

  • 检查消息是否进入正确 Topic、Exchange 和 Queue;
  • 检查 Kafka consumer group 是否订阅正确;
  • 检查 RabbitMQ binding 和 routing key;
  • 检查消费者是否启动、是否被限流或阻塞;
  • 检查消息是否进入 DLT/DLQ。

错误二:业务成功后没有幂等保护

失败表现:

同一支付消息重复到达
支付记录插入两次
库存扣减两次
通知发送两次

诊断:

  • 使用 messageId 查询消费日志;
  • 查询业务表是否有唯一约束;
  • 查询 Inbox 是否和业务事务同库同事务;
  • 检查应用是否在 ACK/offset 提交前崩溃;
  • 检查 Relay 是否因租约超时重复发布。

错误三:捕获异常后正常返回

@KafkaListener(topics = "orders")
public void consume(OrderEvent event) {
    try {
        service.handle(event);
    } catch (Exception ex) {
        log.error("handle failed", ex);
        // 没有继续抛出
    }
}

如果容器认为方法正常返回,可能提交 offset 或发送 ACK,消息随后不会再次投递。日志中虽然有错误,但业务数据已经无法自动恢复。

错误四:失败后永久 requeue

失败表现:

队列深度不降
同一 messageId 每秒出现大量日志
CPU 升高
后续消息延迟持续增加

应检查:

  • 是否对确定性异常使用了 requeue=true
  • 是否记录并限制重试次数;
  • 是否存在 retry queue 和 DLQ;
  • 是否需要暂停消费者处理某一类毒丸消息。

错误五:Outbox 写入和业务写入不在同一个事务

失败表现:

业务表有订单,但没有 outbox
或 outbox 有事件,但业务表没有订单

检查:

  • 两个 Repository 是否使用同一个数据源;
  • 是否存在多个事务管理器;
  • @Transactional 是否被代理;
  • 是否发生同类内部方法自调用;
  • 是否把消息写入放在了另一个异步线程。

错误六:把 Kafka 全局顺序当成业务顺序

失败表现:

订单已经 REFUNDED
但 ES 投影又被旧的 PAID 事件覆盖

检查:

  • 事件 key 是否始终使用同一个 aggregate ID;
  • 消费者是否并发处理同一聚合;
  • 是否有事件版本;
  • 重试 Topic 是否改变了分区;
  • 投影更新是否带版本条件。

十六、生产诊断需要观察哪些指标

消息系统的问题通常表现为延迟、堆积、重复和失败,而不是单个 HTTP 请求直接报错。

Kafka

重点指标包括:

  • consumer lag:消费位点与生产最新位点的差距;
  • 每个分区的 lag,识别热点分区;
  • rebalance 次数;
  • poll 间隔和处理耗时;
  • producer error、retry、record error;
  • DLT 消息数;
  • Topic 分区分布。

lag 增大不一定表示消费者挂了,也可能是:

生产速率 > 消费处理速率

应进一步比较:

生产吞吐
消费成功吞吐
单条处理延迟
重试吞吐

如果业务处理失败后进入重试 Topic,原 Topic lag 可能下降,但系统整体仍在积压,因此不能只监控主 Topic。

RabbitMQ

重点指标包括:

  • ready messages;
  • unacknowledged messages;
  • consumer 数量;
  • publish rate;
  • deliver rate;
  • ack rate;
  • redelivery rate;
  • retry queue 和 DLQ 深度;
  • connection/channel 数量;
  • publisher confirm 失败数。

unacknowledged 持续升高通常表示:

  • 消费者处理过慢;
  • prefetch 过大;
  • 消费者卡在外部调用;
  • ACK 没有发出;
  • 消费者线程池耗尽。

应用与业务

消息指标必须关联业务指标:

messageId
eventId
aggregateId
traceId
topic/queue
partition/offset
retryCount
数据库事务结果
ACK/offset 提交结果

仅有“消费成功日志”不够,还应能够回答:

这条订单事件是否写入 Outbox?
是否被 Relay 发布?
是否被哪个消费者领取?
是否进入 Inbox?
数据库事务是否提交?
是否最终 ACK?

十七、可靠性设计的验证方法

消息系统不能只靠代码阅读验证。应对关键故障窗口做测试。

1. 生产确认测试

模拟:

  • Broker 不可用;
  • Exchange 无绑定;
  • Kafka Leader 切换;
  • 网络在发送后断开;
  • Producer 请求超时但 Broker 实际已写入。

验证:

是否有明确错误?
是否会重试?
重试后是否可能重复?
消费者是否幂等?

2. Outbox 故障注入

在以下位置强制进程退出:

业务表写入后、Outbox 写入前
Outbox 写入后、事务提交前
事务提交后、Relay 发布前
Broker 确认后、标记 SENT 前
标记 SENT 后

预期结果应是:

崩溃位置 预期
业务写入前 业务和 Outbox 都不存在
业务与 Outbox 同一事务中 二者一起提交或一起回滚
Outbox 已提交、未发布 Relay 恢复后继续发布
发布成功、未标记 SENT 可能重复发布,但消费者不重复执行业务
已标记 SENT 不再重复扫描

3. 消费者故障注入

在消费者中分别在以下位置退出:

数据库事务开始前
Inbox 插入后
业务更新后、事务提交前
事务提交后、ACK 前
ACK 后

预期:

  • 事务未提交:消息重试,业务重新执行;
  • 事务已提交、ACK 前:消息重试,Inbox 识别重复;
  • ACK 后:不再重复处理。

这组测试可以直接验证“至少一次 + 幂等”的设计,而不是停留在配置层面。


十八、如何选择 Kafka、RabbitMQ 和 Outbox

Kafka 更适合

  • 高吞吐事件流;
  • 需要较长时间保留和重放;
  • 多个独立消费组;
  • 按 key 保持分区内顺序;
  • 流式处理和 Kafka 到 Kafka 的事务链路;
  • 事件日志、审计、数据管道。

Kafka 的代价是需要理解分区、消费位点、rebalance、lag 和 Topic 生命周期。

RabbitMQ 更适合

  • 复杂路由;
  • 工作队列;
  • 任务分发;
  • 请求处理完成后再 ACK;
  • 需要较细粒度的队列、交换机和死信拓扑;
  • 业务规模中等但路由语义复杂的场景。

RabbitMQ 的代价是需要管理 Exchange、Queue、Binding、ACK、prefetch、重投递和死信拓扑。

Outbox 适合

只要存在:

本地数据库事务 + 异步消息发布

就应认真考虑 Outbox。尤其是:

  • 创建订单后必须发布 OrderCreated
  • 修改支付状态后必须通知下游;
  • 更新数据库后必须刷新缓存或 ES 投影;
  • 业务操作不能接受“数据库成功但消息完全没有”的窗口。

如果业务本身允许丢失通知,或者消息只是非关键统计,也可以不用 Outbox,但应明确这是业务取舍,而不是误以为普通 send() 已经具备原子性。


十九、一套可落地的端到端原则

一个订单服务和库存服务的可靠链路可以归纳为:

HTTP Idempotency-Key
    ↓
订单数据库事务
    ├──订单表
    └──Outbox 表
          ↓
       Relay + Publisher Confirm
          ↓
       Kafka/RabbitMQ
          ↓
       有限重试 + DLT/DLQ
          ↓
       消费者 Inbox + 业务事务
          ↓
       ACK 或提交 offset
          ↓
       Redis/Elasticsearch 异步投影

这条链路的语义是:

  1. HTTP 重试不会无条件创建重复业务对象;
  2. 业务数据提交成功时,待发送事件一定同时存在;
  3. Relay 故障恢复后可以继续发送;
  4. Relay 可能重复发布,但消费者不会重复执行核心业务;
  5. 临时故障可以退避重试;
  6. 确定性错误不会无限阻塞;
  7. 无法处理的消息进入可诊断、可重放的隔离区域;
  8. Redis 和 Elasticsearch 的延迟不会被误认为主数据库事务失败;
  9. Kafka 和 RabbitMQ 的确认只在其各自边界内解释;
  10. 所谓 exactly-once 必须指明资源范围,不能跨数据库、Broker 和外部 HTTP 调用任意推导。

消息队列工程的核心不是追求一个抽象的“绝不重复、绝不丢失”,而是把每个故障窗口转换为明确的状态变化:

未提交
待发送
处理中
已发送
已确认
可重试
已隔离

当这些状态有持久记录、唯一约束、租约恢复、有限重试和可观测证据时,系统即使经历网络抖动、进程崩溃、Broker 重启和重复投递,也能够恢复到可解释、可验证的结果。


系列导航与关联阅读

官方资料

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