跳过导航

数据库与消息队列如何保证最终一致性?本地消息表、事务消息与 Outbox

约 11 分钟...次浏览
专栏MySQL 与 ORM第 12 篇

业务数据库写入成功后再发送消息,是微服务中最常见的双写:订单创建后通知库存、积分和搜索服务。如果把两步直接顺序执行,无论先写数据库还是先发消息,都存在进程崩溃、网络超时和局部成功窗口。最终一致性的核心不是让两套系统“同时提交”,而是保证业务事实一旦提交,相关事件最终可被可靠发现、投递和幂等消费。

1. 为什么普通双写不可靠

先写数据库再发消息:

@Transactional
public void createOrder(CreateOrder command) {
    orderRepository.save(toOrder(command));
    kafkaTemplate.send("order-created", event(command));
}

常见误解是 @Transactional 会同时管理数据库和 Kafka。普通本地事务只能保证数据库连接上的操作;消息发送可能发生在事务提交前、提交后或异步线程中,除非专门配置协议,两者没有共同原子边界。

失败矩阵如下:

顺序失败窗口后果
数据库提交后进程崩溃,消息未发业务有、事件无下游永久不知道
消息成功后数据库回滚事件有、业务无下游处理不存在的事实
发送超时但 Broker 已收到结果未知重试可能产生重复消息

网络超时只能说明发送方没有得到确定响应,不能证明对端没有执行。因此可靠方案必须天然接受“至少一次”和重复处理。

2. 本地消息表:把业务和待发送记录放进同一事务

在业务数据库中增加消息表:

CREATE TABLE local_message (
  id            BIGINT PRIMARY KEY,
  aggregate_type VARCHAR(64) NOT NULL,
  aggregate_id   VARCHAR(64) NOT NULL,
  event_type     VARCHAR(128) NOT NULL,
  payload        JSON NOT NULL,
  status         VARCHAR(16) NOT NULL DEFAULT 'NEW',
  retry_count    INT NOT NULL DEFAULT 0,
  next_retry_at  DATETIME(3) NOT NULL,
  created_at     DATETIME(3) NOT NULL,
  sent_at        DATETIME(3),
  KEY idx_dispatch(status, next_retry_at, id)
) ENGINE=InnoDB;

业务数据与消息记录在同一个本地事务内写入:

@Transactional
public Long createOrder(CreateOrder command) {
    Order order = orderRepository.save(Order.create(command));
    localMessageRepository.save(LocalMessage.of(
        eventIdGenerator.nextId(),
        "Order", order.getId().toString(),
        "OrderCreated", serialize(order.toEvent())
    ));
    return order.getId();
}

数据库保证订单和待发送记录同时提交或同时回滚。后台投递器轮询 NEW 消息,发送成功后标记 SENT

SELECT * FROM local_message
WHERE status='NEW' AND next_retry_at <= NOW(3)
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;

SKIP LOCKED 允许多个投递实例并行领取不同批次。事务中不应长时间等待 Broker;更常见的做法是短事务领取并设置租约,然后在事务外发送,成功后更新状态。租约超时的消息可被重新领取。

3. “发送成功再标记”仍然会重复

投递器可能完成以下步骤:

  1. Broker 已接收消息。
  2. 进程在更新 SENT 前崩溃。
  3. 重启后再次发送同一事件。

这意味着本地消息表通常只能做到至少一次投递。不要试图通过更复杂的状态更新消灭所有重复窗口,而应给事件稳定的 eventId,由消费者幂等处理。

消费者可建立去重表:

CREATE TABLE consumed_event (
  consumer_name VARCHAR(64) NOT NULL,
  event_id      BIGINT NOT NULL,
  consumed_at   DATETIME(3) NOT NULL,
  PRIMARY KEY (consumer_name, event_id)
);

在同一数据库事务中插入去重记录并执行本地业务:

@Transactional
public void handle(OrderCreated event) {
    if (!consumedEventRepository.tryInsert("inventory", event.eventId())) {
        return;
    }
    inventoryRepository.reserve(event.orderId(), event.items());
}

唯一键把并发重复消费变为确定的冲突。只有“去重记录”和业务修改处于同一事务,才能避免已记录消费但业务尚未提交的窗口。

4. Transactional Outbox 与本地消息表有什么区别

两者都把事件记录与业务变更写进同一数据库事务,工程实现经常高度相似。语义上的区别通常是:

  • 本地消息表偏向“待执行的投递任务”,字段包含重试次数、状态和下次执行时间。
  • Outbox 偏向“已经发生的领域事件日志”,记录聚合类型、聚合 ID、事件类型、版本和负载,再由发布器转发。

Outbox 表例子:

CREATE TABLE outbox_event (
  event_id       CHAR(36) PRIMARY KEY,
  aggregate_type VARCHAR(64) NOT NULL,
  aggregate_id   VARCHAR(64) NOT NULL,
  aggregate_version BIGINT NOT NULL,
  event_type     VARCHAR(128) NOT NULL,
  payload        JSON NOT NULL,
  occurred_at    DATETIME(3) NOT NULL,
  published_at   DATETIME(3),
  UNIQUE KEY uk_aggregate_version(
    aggregate_type, aggregate_id, aggregate_version
  ),
  KEY idx_unpublished(published_at, occurred_at)
);

Outbox 强调事件是业务事务的结果,而不是调用消息 SDK 的副作用。它还能通过聚合版本帮助消费者识别乱序和缺失。

5. Outbox 的两种发布方式

5.1 应用轮询

应用定时扫描未发布记录并发送。优点是技术简单、与数据库日志格式解耦;缺点是轮询延迟、扫描压力和抢占逻辑需要自行治理。

不要每次执行无索引的:

SELECT * FROM outbox_event WHERE published_at IS NULL;

应使用覆盖筛选的索引、批次上限和游标。已发布数据要归档或按时间分区,否则活跃索引会持续膨胀。

5.2 CDC 捕获

Debezium 等 CDC 工具读取数据库变更日志,将 outbox 的插入转换为消息。业务应用不再轮询,也无需写回 published_at;数据库提交顺序可以更自然地映射到事件流。

CDC 不是零成本方案:

  • 需要维护 Binlog/WAL 权限、保留期和连接器状态。
  • Schema 演进、快照和断点恢复需要规范。
  • 连接器到 Broker 仍可能至少一次投递。
  • 数据库日志顺序不等于跨分区的全局业务顺序。

如果组织已经有成熟 CDC 平台,Outbox + CDC 通常是很强的组合;小规模系统用轮询更容易落地和排障。

6. 消息队列事务消息

部分消息队列提供事务消息能力。典型流程是:

生产者发送半消息
  -> Broker 持久化但暂不投递
  -> 生产者执行本地数据库事务
  -> 成功:提交消息
  -> 失败:回滚消息
  -> 状态未知:Broker 回查生产者

本地事务状态必须可回查。例如建立事务记录表,或根据订单状态判断事件是否应提交。回查逻辑要幂等,且不能只依赖生产进程内存。

事务消息的优点是低延迟、无需业务轮询消息表;缺点是应用与特定 Broker 协议耦合,回查状态机和运维要求更高。它通常保证的是“消息提交状态与本地事务结果最终协调”,不是让所有消费者的数据库事务与生产者原子提交。

还要区分 Kafka 的 transactional producer:它可以原子写入多个 Kafka 分区,并配合消费位点实现 Kafka 内部的 consume-transform-produce exactly-once;它不会自动把普通 MySQL 事务纳入同一个原子提交。跨 MySQL 与 Kafka 仍需 Outbox、连接器或额外协调方案。

7. 三种方案怎样选择

维度本地消息表Transactional OutboxBroker 事务消息
原子边界业务行 + 消息任务行业务行 + 领域事件行半消息状态 + 本地事务结果协调
发布方式应用轮询/任务调度轮询或 CDCBroker 协议与回查
Broker 耦合低到中
延迟取决于轮询周期轮询中等,CDC 较低通常较低
运维重点积压、重试、表膨胀CDC/轮询、事件 Schema回查、半消息积压
典型语义可靠任务投递领域事件发布Broker 原生事务流程

没有一种方案能免除消费者幂等、失败重试、死信处理和对账。选择依据应包括现有基础设施、团队运维能力、延迟目标与 Broker 锁定成本。

8. 事件载荷应该保存什么

Outbox 中可保存完整事件快照,也可只保存实体 ID。

只保存 ID,发布时再查业务表,看似节省空间,却可能读到后续状态:订单创建事件发布时,订单已经被取消,下游收到的就不再是“创建时事实”。完整事件快照能保持历史语义,但带来存储和 Schema 演进成本。

推荐事件信封包含:

{
  "eventId": "0190a4d4-...",
  "eventType": "OrderCreated",
  "eventVersion": 2,
  "aggregateId": "202607120001",
  "aggregateVersion": 5,
  "occurredAt": "2026-07-12T10:30:00Z",
  "traceId": "...",
  "data": {}
}

事件版本应描述消息契约版本,不要直接把 ORM 实体序列化成消息。实体字段重命名、懒加载代理和内部敏感字段都不应泄露到公共契约。

9. 顺序、乱序和并发消费

全局严格顺序通常代价过高,业务真正需要的往往是“同一个订单的事件有序”。把 aggregateId 作为消息分区键,可让同一聚合事件进入同一分区。

即便如此,重试、死信回放和多来源写入仍可能造成乱序。消费者可以利用 aggregateVersion

  • 收到版本 6,而本地只有版本 4:暂缓、重试或触发补偿查询。
  • 已处理版本 6,又收到版本 5:识别为旧事件并忽略。
  • 相同版本与 eventId 重复:按幂等规则直接成功。

不能只用消息产生时间排序,因为节点时钟和网络延迟都可能不同。

10. 重试与死信不能代替对账

投递失败应使用指数退避和随机抖动,避免 Broker 恢复时大量消息同时重试:

nextDelay = min(maxDelay, base * 2^retryCount) + randomJitter

超过阈值的消息进入死信或人工处理队列,但不能悄悄标记成功。至少监控:

  • 最老未投递消息年龄。
  • NEW、SENDING、FAILED 数量与增长速度。
  • 投递成功率和重试次数分布。
  • 消费延迟、重复率和死信量。
  • CDC 位点落后时间。

即使投递链路设计正确,也可能因程序缺陷、误操作和数据修复绕过正常流程。定期对账应比较业务事实与下游结果,例如已支付订单是否都生成积分记录。对账任务是最终一致性的最后一道防线。

11. 不要忽视数据库提交后的回调窗口

另一种常见写法是在事务提交后回调发送消息。它避免数据库回滚时发出消息,却仍有“提交完成、回调尚未执行时进程崩溃”的窗口。内存事件总线、Spring AFTER_COMMIT 监听器适合非关键通知或配合 Outbox 唤醒发布器,不应独自承担不可丢失的业务事件。

同理,定时扫描业务表推导“哪些事件没发”虽然可以补偿,但若没有明确状态或版本,难以判断事件是否已经发布以及发布的是哪个历史状态。

12. 生产落地检查清单

  • 业务数据与 outbox/message 行是否由同一个数据库事务提交?
  • 事件 ID 是否稳定,重试时是否保持不变?
  • 消费去重与本地业务修改是否在同一事务?
  • 消息载荷是历史快照还是查询引用?语义是否明确?
  • 是否按聚合键分区,并包含聚合版本?
  • 发布器崩溃、发送超时和重复发送时会发生什么?
  • 是否有退避、租约、死信、告警和人工恢复流程?
  • Outbox 表如何归档,索引是否能支撑持续扫描?
  • 事件 Schema 如何兼容升级和回放旧消息?
  • 是否有独立于消息链路的周期性对账?

最终一致性不是一句“失败就重试”。完整方案必须把可靠记录、至少一次投递、幂等消费、顺序控制、可观测性和对账闭环组合起来。本地消息表、Outbox 和事务消息只是解决“如何可靠产生并发现事件”的不同入口,系统能否经受故障,最终取决于整个状态机是否可恢复、可验证。

分享:
文章作者:狼码纪
版权声明:本博客所有文章除特别声明外,均采用 CC BY-NC-SA 4.0 许可协议。文章可能参考了其他优秀文章,如有侵权请联系删除。