消息驱动的分布式事务:本地消息表 + 四大模式

distributed-transaction 📚 learning local-message-table · distributed-transaction · local-message-table · mq · rocketmq · kafka-transaction · transactional-message · reliable-message · eventually-consistent · outbox

TL;DR(30 秒扫完)

  • 共同本质:MQ 是分布式事务最主流的落地手段(不是"有资源浪费"的替代品)
  • 共同前提:①消费端必须幂等(MQ 至少投递一次)②必须有对账/补偿兜底
  • 模式 1 本地消息表(Outbox):业务 SQL + 写消息表在同一本地事务,定时扫表投递;灵魂是"业务成功=消息存在"
  • 模式 2 事务消息(RocketMQ):half message + 回查机制,MQ 内置保证"业务与消息强绑定"
  • 模式 3 可靠消息 + 幂等消费:Producer acks=all + 重试 + Consumer 业务唯一键幂等;最简单但最弱(业务与消息不绑定)
  • 模式 4 Kafka 事务(2.5+):KafkaTransaction 跨 topic 原子写 + offset 一并提交;依赖幂等 Producer + 事务 Producer + read_committed
  • 实现的是最终一致:不是强一致,与 2PC/XA 的本质差异

关键结论

结论 A本地消息表的灵魂是"本地事务+异步消息",业务方在本地事务里把消息写进去
结论 B不丢失靠本地事务,不重复靠主键+消费端幂等,不丢单靠定时扫表+消费端重试
结论 C状态机每个状态转换必须幂等(重复执行结果一致)
结论 Dvs RocketMQ 事务消息,本地消息表依赖业务方运维能力,RocketMQ 依赖 MQ 实现

完整讲解(费曼四步)

STEP 1 · 概念

本地消息表:把"分布式事务"拆解为"本地事务+异步消息"的方案,业务数据和消息在同一本地事务里写入,保证"业务成功=消息存在"。

STEP 2 · 大白话

本地消息表类比:你要给朋友寄一封信(异步消息),但你怕自己忘了,所以先在笔记本上记一笔(消息表),然后每天检查笔记本(定时扫表),发现没寄的就寄出去(投递 MQ),朋友收到后让你打勾(消费确认)。 为什么能解决分布式事务:
  • 业务成功 → 消息一定存在(同一本地事务,要么都成功要么都失败)
  • 消息存在 → 一定能被扫表投递(定时任务保证)
  • 消息投递 → 消费端一定能处理(消费端幂等保证)

STEP 3 · 底层

核心思想

  • 本质:把"分布式事务"拆解为"本地事务 + 异步消息"
  • 核心思想:业务数据和消息在同一本地事务里写入,保证"业务成功 = 消息存在"
  • 为什么能解决分布式事务:
  • 业务成功 → 消息一定存在(同一事务,要么都成功要么都失败)
  • 消息存在 → 一定能被扫表投递(定时任务保证)
  • 消息投递 → 消费端一定能处理(消费端幂等保证)

完整状态机

INIT (本地事务已提交,消息待发送)
  ↓ 定时扫表 + MQ 投递
SENT (消息已发送,等待消费)
  ↓ 消费端 ACK
DELIVERED (消费端已确认)
  ↓ 业务处理完成
ACKED (业务已完成)
每个状态转换都是幂等的(重复执行结果一致)

核心组件

  • 消息表:存业务消息(msg_id、biz_id、payload、status、retry_count、next_retry_time)
  • 扫表任务:定时扫描 status=INIT 的消息,发 MQ,更新状态
  • 消费端:消费 MQ 消息,幂等处理业务,更新消息状态

可靠性保证(三大目标)

目标机制
不丢失本地事务保证"业务成功→消息必发"
不重复消息表主键去重 + 消费端幂等
不丢单定时扫表重试 + 消费端重试

与 RocketMQ 事务消息对比

维度本地消息表RocketMQ 事务消息
扫表逻辑业务方自建MQ 内置回查
业务侵入中低
可靠性依赖业务方依赖 MQ
部署复杂度低中
适合场景业务方运维能力强想用成熟 MQ、代码简洁

工程实现要点

  • 消息表字段:msg_id(唯一键)、status、retry_count、next_retry_time、create_time
  • 扫表 SQL:SELECT * FROM msg_table WHERE status = 'INIT' AND next_retry_time < NOW() LIMIT 100 FOR UPDATE SKIP LOCKED
  • 消费端幂等:基于 msg_id 建唯一键,重复消息直接返回成功
  • 死信队列:超过最大重试次数的消息转入死信队列,人工处理

消息表 DDL 示例

CREATE TABLE local_message (
    msg_id BIGINT PRIMARY KEY,        -- 消息唯一 ID(业务单号)
    biz_id VARCHAR(64) NOT NULL,      -- 业务 ID
    topic VARCHAR(128) NOT NULL,      -- MQ Topic
    payload TEXT,                      -- 消息内容(JSON)
    status ENUM('INIT','SENT','DELIVERED','ACKED') DEFAULT 'INIT',
    retry_count INT DEFAULT 0,
    next_retry_time DATETIME,
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP,
    update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    UNIQUE KEY uk_msg_id (msg_id),
    KEY idx_status_retry (status, next_retry_time)
);

扫表任务示例

@Scheduled(fixedDelay = 1000)
public void scanAndSend() {
    List<Message> messages = messageMapper.selectInitMessages(100);
    for (Message msg : messages) {
        try {
            // 发送 MQ 消息
            mqProducer.send(msg.getTopic(), msg.getPayload());
            // 更新状态为 SENT
            messageMapper.updateStatus(msg.getMsgId(), "SENT");
        } catch (Exception e) {
            // 更新重试次数和下次重试时间
            messageMapper.updateRetry(msg.getMsgId(), 
                msg.getRetryCount() + 1, 
                calculateNextRetryTime(msg.getRetryCount()));
        }
    }
}

消费端幂等示例

public void consume(Message msg) {
    // 基于 msg_id 检查是否已消费
    if (consumeRecordMapper.existsByMsgId(msg.getMsgId())) {
        return; // 已消费,直接返回成功
    }
    
    try {
        // 业务处理
        doBusiness(msg);
        // 记录消费记录(幂等键)
        consumeRecordMapper.insert(msg.getMsgId());
        // 更新消息状态为 DELIVERED
        messageMapper.updateStatus(msg.getMsgId(), "DELIVERED");
    } catch (Exception e) {
        throw new RetryException(e); // 消费端重试
    }
}

STEP 4 · 简化

一句话总结:本地消息表的灵魂是"本地事务+异步消息"——业务方在本地事务里把消息写进去,靠定时扫表和消费端幂等保证最终一致。

常见误区

"本地消息表不需要本地事务"
本地消息表的灵魂是"业务 SQL + 写消息表在同一本地事务",没有本地事务就无法保证"业务成功=消息存在"
"定时扫表是可选的"
定时扫表是本地消息表的关键组件,没有它消息会一直停留在 INIT 状态,永远不会被投递
"消费端不需要幂等"
消费端必须幂等(基于 msg_id 建唯一键),否则消息重复投递会导致重复消费
"本地消息表能保证 100% 可靠"
本地消息表依赖业务方运维能力,如果扫表任务挂了,消息会一直停留在 INIT 状态
"本地消息表和 RocketMQ 事务消息没区别"
本地消息表业务方自建扫表,RocketMQ 内置回查机制,业务方代码更简洁
"MQ 不是分布式事务的实现方式"
MQ 恰恰是分布式事务最主流的落地手段(本地消息表 / 事务消息 / Kafka 事务都基于 MQ),因为 MQ 提供"最终一致 + 异步解耦"的能力
"MQ 实现的是强一致"
MQ 实现的是最终一致,消费端可能出现"上游已提交、下游还没消费"的时间窗口
"消费端不需要幂等"
MQ 是"至少投递一次",可能重复消费,消费端必须幂等(基于业务唯一键去重)
"本地消息表和事务消息谁更好?"
本地消息表更通用(任何 MQ 都能用),事务消息更简洁(依赖 MQ 支持事务特性),没有绝对优劣

延伸追问

"如果定时扫表任务挂了,消息怎么办?"
要点:消息会一直停留在 INIT 状态,不会被投递。恢复方案:扫表任务集群化(多实例)+ 监控告警 + 人工触发兜底
"扫表任务多实例并发运行时,怎么避免重复投递?"
要点:用 FOR UPDATE SKIP LOCKED 锁住消息行,其他实例跳过;或加分布式锁
"本地消息表怎么和 RocketMQ 事务消息选型?"
要点:业务方运维能力强、想自建 → 本地消息表;想用成熟 MQ、代码简洁 → RocketMQ 事务消息
"消息表会无限增长吗?怎么清理?"
要点:ACKED 状态的消息可以定期归档/删除;保留策略根据业务需求(如审计要求保留 6 个月)
"消费端幂等怎么实现?"
要点:基于 msg_id 建唯一键,重复消息直接返回成功;或基于业务单号建幂等表

速查表

维度内容
核心思想本地事务 + 异步消息换分布式事务
灵魂业务 SQL + 写消息表在同一本地事务
流程 3 步本地事务 → 定时扫表投递 → 消费确认
状态机INIT → SENT → DELIVERED → ACKED
可靠性不丢失(本地事务)/ 不重复(主键+幂等)/ 不丢单(扫表+重试)
vs RocketMQ业务方自建扫表 vs MQ 内置回查
工程要点扫表任务集群化、消费端幂等、死信队列、消息归档

关联题目

  • ⚠️ 《如何基于本地消息表实现分布式事务?》— 2026-09-23 Round 2 Q7, ⭐⭐⭐, 未明确"同一本地事务"和状态机
  • ⚠️ 《如何基于 MQ 实现分布式事务》— 2026-10-10 Round 1 Q5, ⭐⭐, 把 MQ 与分布式事务当对立概念,四大模式完全没提

关联知识

本地事务 + 异步消息换分布式事务
寄信前先记笔记本,每天检查没寄的,朋友收到打勾
✦ 记 忆 口 诀 ✦
灵魂是同一本地事务 / 不丢失不重复不丢单 / 状态机每个转换幂等
关键可视化
本地消息表状态机
flowchart LR
  A[INIT 本地事务已提交 消息待发送] --> B[SENT 消息已发送 等待消费]
  B --> C[DELIVERED 消费端已确认]
  C --> D[ACKED 业务已完成]
  A -.定时扫表.-> B
  B -.消费端 ACK.-> C
  C -.业务处理完成.-> D
本地消息表完整流程
flowchart TB
  A[1 本地事务阶段] --> B[业务 SQL + 写消息表 同一本地事务]
  B --> C[2 消息投递阶段]
  C --> D[定时扫表 status 等于 INIT]
  D --> E[投递 MQ]
  E --> F[更新状态 SENT]
  F --> G[3 消费确认阶段]
  G --> H[消费端幂等处理]
  H --> I[更新状态 DELIVERED]
  I --> J[业务完成更新 ACKED]
可靠性三大目标
flowchart TB
  A[可靠性三大目标] --> B[不丢失]
  B --> B1[本地事务保证 业务成功等于消息必发]
  A --> C[不重复]
  C --> C1[消息表主键去重]
  C --> C2[消费端幂等]
  A --> D[不丢单]
  D --> D1[定时扫表重试]
  D --> D2[消费端重试]
本地消息表 vs RocketMQ 事务消息
flowchart TB
  A[对比维度] --> A1[扫表逻辑]
  A --> A2[业务侵入]
  A --> A3[可靠性]
  A --> A4[部署复杂度]
  A --> A5[适合场景]
  A1 --> A1V[本地消息表 业务方自建 vs RocketMQ MQ 内置回查]
  A2 --> A2V[本地消息表 中 vs RocketMQ 低]
  A3 --> A3V[本地消息表 依赖业务方 vs RocketMQ 依赖 MQ]
  A4 --> A4V[本地消息表 低 vs RocketMQ 中]
  A5 --> A5V[本地消息表 业务方运维能力强 vs RocketMQ 想用成熟 MQ]
知识关系

⬆️ 前置(Prerequisite)

distributed-tx-overview

🔄 延伸(Extends)

cross-border-payment-design

⚡ 对比(Contrast)

TCC— TCC 强一致高侵入;本地消息表最终一致低侵入
🎯 概念 📏 规则 ⚠️ 误区 🔍 追问 ✨ 口诀 共 0 张卡,点击翻面