---
title: 消息驱动的分布式事务：本地消息表 + 四大模式
type: concept
domain: distributed-transaction
tags: [distributed-transaction, local-message-table, mq, rocketmq, kafka-transaction, transactional-message, reliable-message, eventually-consistent, outbox]
status: learning
created: 2026-09-23
last_reviewed: 2026-10-10
method: feynman
related_questions:
  - "如何基于本地消息表实现分布式事务？"
  - "如何基于 MQ 实现分布式事务"
related_knowledge:
  - ./distributed-tx-overview.md
  - ./tcc-pitfalls.md
  - ./cross-border-payment-design.md
  - ./2pc-3pc.md
anki_cards: 12
interview_rounds:
  - "2026-09-23-round-2-Q7"
  - "2026-10-10-round-1-Q5"
---

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

> 本地消息表是"用本地事务+异步消息换分布式事务"的经典落地。但**基于 MQ 实现分布式事务不止本地消息表一种**——业界有四大主流模式（本地消息表 / 事务消息 / 可靠消息+幂等 / Kafka 事务），核心都是"业务与消息的强绑定"。共同前提：**消费端必须幂等 + 必须有对账兜底**。

## 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**：状态机每个状态转换必须幂等（重复执行结果一致）
- **结论 D**：vs 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 示例

```sql
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)
);
```

#### 扫表任务示例

```java
@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()));
        }
    }
}
```

#### 消费端幂等示例

```java
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 · 简化

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

## 消息驱动的分布式事务四大模式全景

> 本地消息表只是**四大主流模式**之一。业界完整对比如下——所有模式共同点是"用消息桥接跨服务的最终一致"。

### 模式对比总览

| 模式 | 强绑定方式 | 侵入性 | 依赖 | 适合场景 |
|------|-----------|-------|------|---------|
| **本地消息表（Outbox）** | 本地事务写入 | 中 | 业务方自建扫表 | 业务方运维能力强 |
| **事务消息（RocketMQ）** | half message + 回查 | 低 | MQ 支持事务特性 | 想代码简洁 |
| **可靠消息 + 幂等消费** | 无（不绑定） | 最低 | 简单 MQ | 允许业务/消息短暂不一致 |
| **Kafka 事务（2.5+）** | KafkaTransaction + 幂等 Producer | 中 | Kafka 2.5+ | Kafka 生态、跨 topic 原子写 |

### 模式 1：本地消息表（Outbox Pattern）—— 详解见上方"完整讲解"

- 业务库内新增 outbox 表
- 本地事务：`BEGIN → INSERT 业务数据 + INSERT outbox_msg → COMMIT`
- 后台定时任务扫 outbox 表投递 MQ，投递成功后标记 outbox
- **优势**：完全复用本地事务，不依赖 MQ 支持事务特性（Kafka/RabbitMQ 都能用）
- **代价**：需要额外的定时任务和 Outbox 表维护

### 模式 2：事务消息（RocketMQ 最典型）

**核心机制**：Producer 发的是**half message**（对 Consumer 不可见），Producer 执行本地事务后再决定提交/回滚。

**流程**：
```
1. Producer 发送 half message 到 Broker（对 Consumer 不可见）
2. Producer 执行本地事务
3. Producer 通知 Broker：
   - 本地事务成功 → commit half message（Consumer 可见）
   - 本地事务失败 → rollback half message（删除）
4. 若 Producer 在 3 之前崩溃：Broker 定时回查 Producer 的本地事务状态
```

**关键组件**：
- `TransactionProducer`：发送 half message
- `TransactionListener`：Producer 侧回查接口 `checkLocalTransaction(msg)` 返回 COMMIT / ROLLBACK / UNKNOWN
- **Broker 侧回查机制**：防止 Producer 崩溃导致 half message 悬挂

**Java 示例**：
```java
public class InventoryMessageListener implements TransactionListener {
    public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        // 执行本地事务：扣库存
        if (inventoryService.deduct(msg.getKeys())) {
            return LocalTransactionState.COMMIT_MESSAGE;  // 提交消息
        } else {
            return LocalTransactionState.ROLLBACK_MESSAGE;  // 回滚消息
        }
    }
    public LocalTransactionState checkLocalTransaction(MessageExt msg) {
        // 回查：根据业务状态返回 COMMIT/ROLLBACK/UNKNOWN
        return inventoryService.checkDeducted(msg.getKeys())
            ? LocalTransactionState.COMMIT_MESSAGE
            : LocalTransactionState.ROLLBACK_MESSAGE;
    }
}
```

### 模式 3：可靠消息 + 幂等消费

**核心机制**：Producer 保证消息可靠投递（`acks=all`、重试、持久化），Consumer 用**业务唯一键**做幂等去重。

**流程**：
```
1. Producer 发普通消息（非事务）
2. Producer 保证消息可靠（acks=all + 重试 + 持久化）
3. Consumer 用业务唯一键（订单号、支付流水号）幂等去重
4. 消费失败 → 重试（可能重复投递）
```

**特点**：
- **优点**：最简单，任何 MQ 都能用，无需 MQ 支持事务
- **缺点**：**不保证业务与消息的强绑定**（可能业务提交成功了但消息在投递中丢失）
- **适合**：能容忍"业务成功但消息偶发丢失"的场景（业务方有对账兜底）

### 模式 4：Kafka 事务（2.5+ 原生支持）

**核心 API**（4 个）：
```java
KafkaProducer producer = new KafkaProducer(props);
producer.initTransactions();              // 初始化事务
producer.beginTransaction();             // 开始事务
producer.send(new ProducerRecord<>("orders", order));
producer.sendOffsetsToTransaction(offsets, consumerGroupMetadata);  // 提交消费 offset
producer.commitTransaction();            // 提交（原子）
// 异常时：producer.abortTransaction();
```

**底层机制**：
- `transactional.id` → 唯一 PID（Producer ID）
- `__transaction_state` 内部 Topic 记录事务状态
- Consumer 端 `isolation.level=read_committed` 只读已提交事务

**局限**：
- **只能保证 Kafka 内部 exactly-once**，跨系统需业务配合
- 依赖 Kafka 2.5+，非主流 MQ 无法用

### 共同约束（无论哪种模式）

1. **消费端必须幂等**：MQ 是"至少投递一次"，可能重复消费。基于业务唯一键建唯一索引
2. **必须有对账/补偿兜底**：定时扫描未完成状态做补偿
3. **实现的是最终一致**：不是强一致。消费端可能存在"上游已提交、下游还没消费"的时间窗口（秒级到分钟级）
4. **强一致场景仍需 2PC/XA/Seata AT**：MQ 方案不能替代强一致协议

### 电商下单场景的实际落地（混合方案）

```
订单服务：
  BEGIN
    扣库存（本地事务）
    INSERT order + INSERT outbox_msg（本地消息表）
  COMMIT

后台任务：扫 outbox → 投递 MQ

支付服务：
  消费"扣库存成功"消息 → 创建订单（消费端幂等：订单号唯一索引）
  消费"下单成功"消息 → 发起支付
```

即使消息延迟，也不会重复扣库存，也不会重复下单。

### STEP 5 · 简化（四大模式一句话对比）

> **本地消息表**（业务方自建扫表，通用）｜ **事务消息**（MQ 内置回查，简洁）｜ **可靠消息+幂等**（最简单，不绑定）｜ **Kafka 事务**（Kafka 生态专属）。选型的核心判断是：**业务与消息需要多强的绑定关系**——强绑定选事务消息或本地消息表，弱绑定选可靠消息。

## 常见误区（从面试记录提炼）

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

## 延伸追问（面试追问预演）

1. **"如果定时扫表任务挂了，消息怎么办？"**
   - 要点：消息会一直停留在 INIT 状态，不会被投递。恢复方案：扫表任务集群化（多实例）+ 监控告警 + 人工触发兜底

2. **"扫表任务多实例并发运行时，怎么避免重复投递？"**
   - 要点：用 `FOR UPDATE SKIP LOCKED` 锁住消息行，其他实例跳过；或加分布式锁

3. **"本地消息表怎么和 RocketMQ 事务消息选型？"**
   - 要点：业务方运维能力强、想自建 → 本地消息表；想用成熟 MQ、代码简洁 → RocketMQ 事务消息

4. **"消息表会无限增长吗？怎么清理？"**
   - 要点：ACKED 状态的消息可以定期归档/删除；保留策略根据业务需求（如审计要求保留 6 个月）

5. **"消费端幂等怎么实现？"**
   - 要点：基于 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 与分布式事务当对立概念，四大模式完全没提

## 关联知识

- [分布式事务概览与模式分类](./distributed-tx-overview.md)
- [TCC 与三大坑](./tcc-pitfalls.md)
- [跨境支付方案设计](./cross-border-payment-design.md)

## Anki 候选卡片

**卡片 1**
- Q：本地消息表的核心思想是什么？
- A：用"本地事务+异步消息"换"分布式事务"，业务数据和消息在同一本地事务里写入，保证"业务成功=消息存在"

**卡片 2**
- Q：本地消息表的完整流程 3 步是什么？
- A：本地事务阶段（业务 SQL+写消息表）→ 消息投递阶段（定时扫表+MQ 投递）→ 消费确认阶段（消费端幂等处理+更新状态）

**卡片 3**
- Q：本地消息表的状态机是什么？
- A：INIT → SENT → DELIVERED → ACKED，每个状态转换必须幂等

**卡片 4**
- Q：本地消息表的可靠性三大目标怎么保证？
- A：不丢失（本地事务）/ 不重复（主键+消费端幂等）/ 不丢单（定时扫表+消费端重试）

**卡片 5**
- Q：本地消息表和 RocketMQ 事务消息的区别是什么？
- A：本地消息表业务方自建扫表，依赖业务方运维能力；RocketMQ 内置回查机制，业务方代码更简洁

**卡片 6**
- Q：扫表任务多实例并发运行时，怎么避免重复投递？
- A：用 `FOR UPDATE SKIP LOCKED` 锁住消息行，其他实例跳过；或加分布式锁

**卡片 7**
- Q：消费端幂等怎么实现？
- A：基于 msg_id 建唯一键，重复消息直接返回成功；或基于业务单号建幂等表

**卡片 8**
- Q：如果定时扫表任务挂了，消息怎么办？
- A：消息会一直停留在 INIT 状态。恢复方案：扫表任务集群化（多实例）+ 监控告警 + 人工触发兜底
