Skip to content

可靠消息与Outbox:双写、投递、消费和重放全过程

可靠消息通常以至少一次语义运行。消费者不能假设同一事件只来一次,完整幂等模型见:MQ消费幂等全过程

一、学完本页要真正会什么

可靠消息不是“调用send()没报错”。一条业务事件要跨越数据库本地事务、应用进程、Broker存储、消费确认和消费者数据库,每个边界都可能在“动作成功、确认丢失”时产生未知结果和重复。

学完本页,你应该能够:

  1. 列出数据库和MQ双写的所有错误顺序与崩溃窗口。
  2. 解释Outbox关闭了哪个窗口、为什么仍然至少一次。
  3. 设计Outbox状态、租约、owner token、重试和归档字段。
  4. 说明多个投递器怎样安全抢占,不重复并发处理同一行。
  5. 处理Broker已收但发送端状态回写失败。
  6. 比较轮询Outbox、CDC Outbox和MQ事务消息。
  7. 设计消费者去重记录与业务写入同一本地事务。
  8. 处理外部HTTP、短信等无法与消费日志同事务的副作用。
  9. 解释顺序、事件版本、死信和安全重放。
  10. 按eventId、businessKey、Outbox状态、Broker位置和消费日志排查。

二、先理解At-most-once、At-least-once与Exactly-once

语义含义代价/风险
At-most-once最多处理一次,不重试或先确认失败时可能丢失
At-least-once失败或确认未知时重试可能重复,消费者必须幂等
Exactly-once effect业务效果恰好一次依赖稳定业务键、原子去重和副作用边界,不是一个Broker开关

Kafka生产者幂等、事务或RocketMQ事务消息能增强特定边界,但不会自动覆盖消费者的MySQL、HTTP、短信和第三方支付。端到端“恰好一次效果”最终仍要由业务唯一键和状态机证明。

三、数据库和MQ为什么不能直接双写

3.1 先写数据库再发消息

mermaid
flowchart TD
    A["数据库事务提交订单"] --> B["应用准备发送消息"]
    B --> C["进程宕机或网络失败"]
    C --> D["订单存在,但事件永久没有发送"]

3.2 先发消息再写数据库

mermaid
flowchart TD
    A["消息先进入Broker并可消费"] --> B["消费者开始处理"]
    A --> C["生产者数据库事务失败或回滚"]
    C --> D["下游处理了一个不存在的业务事实"]

3.3 在本地事务里直接send

mqClient.send()写进@Transactional也不能让普通MQ加入数据库事务:

text
send成功 → 数据库后来回滚:幽灵消息
数据库准备提交 → send失败:整个长事务回滚且锁持有更久
send响应超时:不知道Broker是否已经收到

四、Outbox关闭双写窗口的原理

业务数据和事件意图写到同一个数据库、同一个本地事务:

mermaid
flowchart TD
    A["开始本地事务"] --> B["写订单等业务事实"]
    B --> C["写Outbox事件PENDING"]
    C --> D{"本地事务结果"}
    D -- "提交" --> E["业务事实和待发事件同时存在"]
    D -- "回滚" --> F["业务事实和事件同时不存在"]
    E --> G["独立投递器随后可靠发送"]

Outbox保证:只要订单事务提交,数据库里就有一条可恢复的待发事件。它没有保证投递器只发送一次,因为“Broker接收”和“Outbox标SENT”仍属于两个系统。

五、Outbox表设计

sql
create table t_outbox_event (
    event_id varchar(64) not null,
    aggregate_type varchar(64) not null,
    aggregate_id varchar(128) not null,
    aggregate_version bigint not null,
    event_type varchar(128) not null,
    event_version int not null,
    destination varchar(128) not null,
    message_key varchar(128) not null,
    payload text not null,
    headers text null,
    status varchar(32) not null,
    owner_token varchar(64) null,
    lease_until datetime null,
    retry_count int not null,
    next_retry_at datetime null,
    last_error_code varchar(64) null,
    last_error_message varchar(512) null,
    created_at datetime not null,
    sent_at datetime null,
    updated_at datetime not null,
    primary key (event_id),
    key idx_outbox_dispatch (status, next_retry_at, created_at),
    key idx_outbox_aggregate (aggregate_type, aggregate_id, aggregate_version)
);

字段意义:

字段为什么需要
event_id生产、投递和消费端稳定去重键
aggregate_id/version同一业务对象排序和旧事件保护
event_type/version事件契约演进
destination/message_key路由与分区键
statusPENDING、SENDING、SENT、RETRY、DEAD
owner_token/lease_until多投递器抢占与宕机接管
retry_count/next_retry_at有界退避和可索引扫描
last_error错误分类和排查

payload中不要直接放密码、Token、密钥和不必要的完整敏感信息。事件是长期复制的数据资产,脱敏和最小化应在产生时完成。

六、生产者本地事务Demo

java
@Transactional(rollbackFor = Exception.class)
public void createOrder(CreateOrderCommand command) {
    Order order = Order.create(command.getOrderNo(), command.getUserId());
    orderRepository.insert(order);

    OutboxEvent event = OutboxEvent.pending(
            command.getEventId(),
            "Order",
            command.getOrderNo(),
            order.getVersion(),
            "OrderCreated",
            1,
            JsonUtils.toJson(OrderCreatedEvent.from(order)));
    outboxRepository.insert(event);
}

关键约束:

  • eventId在业务操作开始时生成并跨重试稳定。
  • 订单号和eventId有唯一约束。
  • 订单写入和Outbox插入使用同一数据源、同一本地事务。
  • 不在提交前把事件变成可消费消息。

七、多投递器怎样安全抢占

多个实例同时扫描PENDING时,不能都读取同一批后再发送。常见方案:

7.1 数据库锁定领取

在数据库支持且版本语义明确时,可以在短事务中使用锁定读和SKIP LOCKED领取不同批次,然后快速更新为SENDING并提交。不同数据库、MySQL版本和隔离级别支持不同,必须验证。

7.2 条件更新领取

sql
update t_outbox_event
set status = 'SENDING',
    owner_token = :token,
    lease_until = :leaseUntil,
    updated_at = now()
where event_id = :eventId
  and status in ('PENDING', 'RETRY')
  and next_retry_at <= now();

影响行数为1才获得发送权。领取事务必须短,不能持数据库行锁等待MQ网络发送。

mermaid
flowchart TD
    A["扫描可发送eventId"] --> B["用唯一ownerToken条件抢占"]
    B --> C{"更新是否为1行"}
    C -- "是" --> D["提交领取事务后发送MQ"]
    C -- "否" --> E["其他投递器已领取,跳过"]

八、SENDING投递器宕机怎样接管

投递器领取后可能在发送前或发送后宕机。使用lease_until

  1. 正常投递器只更新自己ownerToken对应的记录。
  2. SENDING租约未到期,不允许其他实例并发接管。
  3. 租约过期后,新实例先判断Broker/业务是否可查询,再以新ownerToken条件接管。
  4. 旧实例恢复后,更新状态时ownerToken不匹配,不能覆盖新处理者。

租约只能防止大部分并发发送,不能消除“旧投递器长时间暂停后恢复继续send”的风险。由于消息发送本身允许至少一次,消费者幂等仍是最终防线。

九、Broker已收但Outbox未标SENT为什么必然重复

mermaid
flowchart TD
    A["投递器发送eventId=E1"] --> B["Broker持久化E1"]
    B --> C["Broker响应发送成功"]
    C --> D["投递器在更新Outbox前宕机"]
    D --> E["Outbox仍为SENDING/RETRY"]
    E --> F["恢复任务再次发送E1"]
    F --> G["消费者必须按E1复用第一次结果"]

即使Producer SDK自身有幂等能力,也必须理解其作用域、会话、序列和Broker版本。业务Outbox不能依赖一个不明确的传输层假设来取消消费者幂等。

十、投递状态机和错误分类

text
PENDING → SENDING → SENT

           RETRY → SENDING

            DEAD

发送结果分类:

类型处理
明确成功ownerToken条件更新SENT,记录Broker元数据
瞬时失败RETRY,指数退避和抖动
响应未知RETRY或查询Broker,允许重复发送
永久配置错误DEAD并告警,例如Topic不存在/权限拒绝
Payload不合法DEAD,修复数据或转换器后受控重放

不能把所有异常无限立即重试。权限错误、Schema错误不会靠第10000次重试自动恢复。

十一、CDC Outbox怎样工作

轮询器不断查询Outbox表;CDC则读取数据库提交日志:

mermaid
flowchart TD
    A["订单和Outbox同一事务提交"] --> B["数据库产生Redo/Binlog/WAL变更"]
    B --> C["CDC Connector按提交顺序读取"]
    C --> D["只筛选Outbox表事件"]
    D --> E["转换为Broker消息"]
    E --> F["保存Connector位点"]

优点:减少业务表轮询和延迟,按数据库提交日志捕获。仍需处理:

  • Connector位点持久化和恢复。
  • 位点提交前后崩溃导致重复。
  • DDL/Schema演进。
  • 大事务和大payload。
  • 事件路由和删除/墓碑语义。
  • Outbox归档与CDC读取位置关系。
  • 下游幂等。

直接CDC业务表会暴露存储Schema并难表达领域事件;CDC Outbox由业务显式写事件,语义通常更稳定。

十二、RocketMQ事务消息关闭哪个窗口

抽象流程:

mermaid
flowchart TD
    A["Producer发送半消息"] --> B["Broker持久化但暂不可消费"]
    B --> C["Producer执行本地数据库事务"]
    C --> D{"本地事务事实"}
    D -- "提交" --> E["Producer通知Broker Commit"]
    D -- "回滚" --> F["通知Broker Rollback"]
    D -- "通信未知" --> G["Broker回查本地事务"]
    G --> H["回查订单/事务流水,不读内存变量"]
    H --> E
    H --> F

事务消息把“本地事务结果与消息是否可见”的协调交给Broker回查机制。它不保证:

  • 消费者只收到一次。
  • 消费者数据库和消息确认原子。
  • 消费者调用外部HTTP只执行一次。
  • 所有MQ产品都有相同语义。

本地事务需要保存可回查状态;若回查永远返回UNKNOWN,Broker最终会按产品策略处理,不能把它当永久等待队列。具体次数、间隔和状态以RocketMQ版本配置为准。

十三、消费者至少一次处理全过程

mermaid
flowchart TD
    A["Consumer收到eventId"] --> B["开始本地数据库事务"]
    B --> C["原子插入消费日志或占用唯一键"]
    C --> D{"是否首次消费"}
    D -- "是" --> E["校验事件版本和业务状态"]
    E --> F["更新业务表"]
    F --> G["消费日志标SUCCESS"]
    G --> H["提交本地事务"]
    H --> I["再向Broker ACK/提交Offset"]
    D -- "否且已SUCCESS" --> J["复用结果并ACK"]

业务更新和消费日志在同一数据库时,应放在同一本地事务。若业务提交后ACK丢失,消息会重投;消费日志让第二次跳过副作用并重新ACK。

十四、消费日志表

sql
create table t_message_consume (
    consumer_group varchar(128) not null,
    event_id varchar(64) not null,
    event_type varchar(128) not null,
    request_hash varchar(128) not null,
    status varchar(32) not null,
    result_ref varchar(128) null,
    retry_count int not null,
    first_received_at datetime not null,
    updated_at datetime not null,
    primary key (consumer_group, event_id)
);

为什么主键包含consumer_group:同一事件需要被库存、积分、通知等不同逻辑各处理一次;不同消费者不能共用一条全局“已消费”记录互相跳过。

为什么保存request_hash:同一个eventId却携带不同payload是协议冲突,不能复用旧成功结果。

十五、消费者先查再处理为什么不安全

错误流程:

text
exists(eventId) == false
→ 执行业务
→ insert consume_log

两个并发投递都可能查到false并重复执行。应先通过唯一键原子占位,或让业务表本身以eventId/businessKey唯一;唯一冲突后读取已有状态。完整原理见分布式幂等

十六、消费者调用HTTP、短信等外部副作用怎么办

消费日志和本地业务表可以同事务,但外部HTTP无法加入普通本地事务:

text
外部短信发送成功
→ 应用在记录SUCCESS前宕机
→ 消息重投
→ 短信可能再次发送

可选策略:

  • 下游支持幂等键:使用eventId调用并查询结果。
  • 再写本地Command Outbox,由专门投递器调用外部系统。
  • 对天然不可幂等通知,业务上接受至少一次并去重模板/频率。
  • 高风险资金接口使用稳定业务流水和事实查询。

不能在调用外部接口前就把消费日志标SUCCESS,否则调用失败后Broker不再重投,副作用永久丢失。

十七、事件顺序和版本保护

同一订单可能产生:

text
OrderCreated version=1
OrderPaid version=2
OrderCanceled version=3

消息网络和重试可能让旧事件晚到。消费者应:

  • 使用aggregateId作为同一分区/队列路由键,尽量保持局部顺序。
  • 在派生表保存last_applied_version。
  • 只接受期望的下一版本,或按业务允许跳跃但拒绝旧版本。
  • 缺版本时进入等待/重查事实源,而不是无条件覆盖。
sql
update t_order_search_projection
set status = :status,
    last_applied_version = :version
where order_no = :orderNo
  and last_applied_version < :version;

仅靠Broker顺序不够:多生产者、失败重试、分区扩容和重放都可能影响观察顺序。

十八、事件契约与Schema演进

事件是已经发生的事实,发布后可能被长期存储和重放。演进原则:

  • 事件包含eventTypeeventVersion
  • 新增字段提供默认语义。
  • 不随意改变已有字段含义或类型。
  • 消费者支持新旧版本并在完成迁移后再淘汰旧版。
  • Payload过大时存引用要考虑源数据变化和保留期。
  • 敏感字段变更要考虑历史消息和归档中的删除/脱敏要求。

十九、重试、死信和停车场队列

错误分类:

类型示例处理
瞬时技术错误数据库短暂不可用、网络抖动指数退避重试
业务等待条件依赖事件尚未到达延迟重试或查事实源
永久数据错误缺字段、无法反序列化死信并修复转换/数据
永久业务拒绝状态不允许记录异常,不无限重试
结果未知外部调用响应丢失使用原幂等键查询事实

死信不是垃圾桶。必须有:

  • 告警和负责人。
  • 原Topic、分区/队列、Offset/消息标识。
  • eventId、businessKey、错误分类。
  • Payload安全存储和保留期。
  • 修复后的受控重放工具。
  • 重放审计和速率限制。

二十、安全重放流程

mermaid
flowchart TD
    A["选择死信/时间范围/业务键"] --> B["在隔离环境验证修复后的消费者"]
    B --> C["确认消费者幂等和版本保护"]
    C --> D["小批灰度重放并限速"]
    D --> E["监控成功、重复、下游负载和新死信"]
    E --> F["扩大批次或停止回滚"]
    F --> G["记录操作者、原因、范围和结果"]

不能把死信直接批量重新发回原Topic且不限制速度,否则可能再次打垮仍未恢复的下游。重放使用原eventId,不能生成新ID绕过幂等。

二十一、Outbox与消费日志怎样归档

永久保留所有记录会导致表膨胀和索引变慢,过早删除又失去去重和审计依据。

保留期至少考虑:

  • Broker最大重试/保留/可回放窗口。
  • 客户端和人工可能重放的最长时间。
  • 业务对账和审计要求。
  • 法规与敏感数据删除要求。

只归档已确认终态且超过安全窗口的记录。进行中的SENDING、RETRY、PROCESSING不能按创建时间粗暴删除。大表可按时间分区、冷热归档,但索引和扫描策略要随之调整。

二十二、Outbox、CDC与事务消息对比

方案事实记录位置投递机制优点主要代价
轮询Outbox业务数据库定时扫描简单、可见、易补偿轮询延迟、表压力、抢占
CDC Outbox业务数据库+日志Connector读取提交日志低轮询压力、提交顺序CDC基础设施、位点和Schema治理
MQ事务消息Broker半消息+本地事实提交/回滚和回查发送链路自然依赖产品语义、回查和本地状态
afterCommit回调仅应用内存提交后立即send简单提交后宕机仍会丢,不可靠

afterCommit适合非关键提示或作为低延迟尝试,但不能替代持久Outbox。可以在提交后立即尝试唤醒投递器,可靠性仍由数据库Outbox兜底。

二十三、商业场景:支付成功传播

mermaid
flowchart TD
    A["支付回调验签、金额和状态校验"] --> B["本地事务条件更新支付/订单"]
    B --> C["同事务写PaymentSucceeded Outbox"]
    C --> D["投递器或CDC发布事件"]
    D --> E["积分消费者幂等入账"]
    D --> F["通知消费者幂等/限频发送"]
    D --> G["ES消费者按订单版本更新"]
    E --> H["对账验证所有派生结果"]
    F --> H
    G --> H

支付回调主事务不应同步等待积分、短信、ES。每个消费者使用自己的consumerGroup + eventId去重,失败互不阻塞,并可独立重试和对账。

二十四、生产排查Runbook

24.1 Outbox PENDING持续堆积

  1. 看待发数量、最老年龄和新增速率。
  2. 检查扫描索引是否命中、查询是否锁等待。
  3. 检查投递器实例、线程池、连接池和租约。
  4. 看Broker发送TPS、延迟、权限和Topic路由。
  5. 按错误码区分瞬时失败和永久配置错误。
  6. 扩容前确认数据库扫描和Broker容量,避免投递器争抢更严重。

24.2 大量SENDING租约过期

检查投递器GC/宕机、发送超时、lease是否小于正常P99、ownerToken条件更新以及旧实例是否恢复后继续发送。允许接管后重复,但消费者必须幂等。

24.3 消息已在Broker但Outbox仍RETRY

这是发送确认未知窗口。按同一eventId再次发送或使用产品可用的查询能力,不要手工直接标SENT;先确认Broker和业务事实。消费者去重会吸收重复。

24.4 消费者重复修改业务

检查是否先查后写、消费日志和业务是否同事务、eventId是否稳定、不同payload是否复用ID、ACK是否在本地事务提交前、外部副作用是否支持幂等。

24.5 某个聚合事件乱序

检查messageKey/分区、生产者数量、aggregateVersion、重试队列是否改变路由、重放方式和消费者版本条件更新。根据数据库事实源补缺失版本或重建投影,不要无条件按到达顺序覆盖。

24.6 死信重放后再次堆积

确认根因是否真正修复、重放速率是否超过下游容量、旧payload是否能被新代码解析、依赖数据是否已准备。停止全量重放,回到小批灰度。

二十五、监控指标

指标价值
outbox_pending_count/oldest_age生产传播延迟
outbox_claim_latency数据库扫描和争抢
send_success/retry/dead发送健康度
send_unknown_totalBroker确认不确定窗口
lease_expired_total投递器宕机或超时配置
consumer_lag/oldest_message_age消费能力和用户影响
dedup_hit_total重复投递实际规模
consume_retry/dead消费错误健康度
version_gap/out_of_order顺序和版本问题
replay_rate/result人工恢复风险

指标必须按destination、eventType、consumerGroup、错误分类和必要业务域拆分,但避免把高基数eventId直接作为Metrics标签。

二十六、JDK 8 Demo:重复投递与消费结果复用

java
import java.util.LinkedHashMap;
import java.util.Map;

public class AtLeastOnceMessageDemo {
    static final class Consumer {
        private final Map<String, String> results =
                new LinkedHashMap<String, String>();
        private int businessExecutions;

        String consume(String eventId) {
            String old = results.get(eventId);
            if (old != null) {
                return old;
            }
            businessExecutions++;
            String result = "handled-" + businessExecutions;
            results.put(eventId, result);
            return result;
        }
    }

    public static void main(String[] args) {
        Consumer consumer = new Consumer();
        String first = consumer.consume("EVENT-9001");
        String duplicate = consumer.consume("EVENT-9001");
        System.out.println("first=" + first);
        System.out.println("duplicate=" + duplicate);
        System.out.println("businessExecutions=" + consumer.businessExecutions);
    }
}

输出:

text
first=handled-1
duplicate=handled-1
businessExecutions=1

真实消费者用数据库唯一键和本地事务,而不是内存Map。Demo只证明重复到达和重复业务效果是两件不同的事。

二十七、常见错误与后果

错误后果正确方向
DB提交后直接send无Outbox进程宕机造成永久漏消息同事务持久化事件意图
Outbox标SENT后再发标记后宕机造成漏发先发再标,允许重复并幂等消费
多投递器先查后发同一事件被并发大量发送ownerToken条件抢占和租约
领取事务持锁等待MQ数据库锁和连接长期占用短事务领取,提交后发送
消费者先查日志再执行业务并发重复执行唯一键原子占位,同事务写业务
ACK在业务提交前宕机后消息不再投递但业务没完成业务事务提交后ACK
外部HTTP成功前先标消费成功外部失败时永久漏执行幂等下游或Command Outbox
重放生成新eventId绕过去重导致重复副作用保留原eventId并记录replay元数据
死信无限自动回原Topic毒消息循环和下游雪崩修复后小批限速重放
只靠Broker顺序不带版本旧事件晚到覆盖新状态aggregateVersion条件更新
afterCommit当可靠消息提交后进程宕机仍丢持久Outbox或事务消息

二十八、面试标准回答

28.1 Outbox解决什么

Outbox把业务数据和待发事件写入同一个数据库本地事务,保证业务提交时一定留下可恢复事件意图,关闭“数据库成功但消息没记录”的双写窗口。独立投递器或CDC随后发送,发送确认未知时允许重复,因此消费者仍要幂等。

28.2 为什么Outbox仍然会重复消息

Broker已持久化消息后,投递器可能在把Outbox标SENT前宕机;恢复扫描看到未SENT会再次发送。同理ACK丢失也会让Broker重投。至少一次系统选择重复而不是静默丢失,业务通过稳定eventId和原子消费日志吸收重复。

28.3 消费日志怎样和业务保持一致

若消费日志和业务表同库,消费者在一个本地事务中先用consumerGroup + eventId唯一键原子占位,校验状态后更新业务并标记成功,提交后再ACK。业务提交后ACK丢失时,重复消息命中成功记录并重新ACK。

28.4 CDC Outbox和轮询Outbox怎么选

轮询实现直观、可控,但有扫描延迟、索引和多实例抢占成本;CDC读取数据库提交日志,延迟低且减少轮询,但需要Connector位点、Schema演进和运维能力。两者都不能取消重复和消费者幂等。

28.5 事务消息能保证消费者只执行一次吗

不能。事务消息主要协调生产者本地事务与消息是否可见;消费者仍可能因ACK丢失、重平衡和恢复而重复收到。消费者的数据库或外部副作用仍需稳定业务键、幂等、状态机和对账。

二十九、关联知识点

本章小结

可靠消息的本质是把每个不可靠确认窗口变成可恢复状态:生产端用Outbox或事务消息保存事件意图,投递端用租约和重试发送,消费端用唯一键和本地事务吸收重复,死信和对账负责长期失败。至少一次不是缺陷,而是网络不确定下“宁可重复、不能静默丢失”的工程选择。