Skip to content

分布式事务与实现

  分布式事务是微服务架构中最难啃的骨头——一次业务操作跨越多个服务、多个数据库,如何保证数据一致性? 单机时代一个 @Transactional 搞定,微服务时代同样的操作变成了分布式难题。

一、为什么需要分布式事务

1.1 单机事务 vs 分布式事务

单机事务(一个数据库):

  @Transactional
  public void createOrder() {
      orderMapper.insert(order);      // 写订单表
      inventoryMapper.deduct(1001);   // 扣库存
      accountMapper.deduct(100);      // 扣余额
  }
  // 要么全成功,要么全回滚,一个 @Transactional 搞定


分布式事务(微服务架构):

  订单服务 ──→ 订单数据库   ──→ INSERT 订单 ✅
  库存服务 ──→ 库存数据库   ──→ 扣减库存 ✅
  账户服务 ──→ 账户数据库   ──→ 扣减余额 ❌ 失败!

  此时订单已创建,库存已扣减,但余额扣款失败——
  怎么回滚?@Transactional 管不了跨服务、跨数据库!

1.2 典型场景

场景涉及服务事务要求
下单订单 + 库存 + 账户 + 优惠券扣库存、扣余额、减优惠券、创建订单,全部成功或全部回滚
转账账户 A + 账户 BA 扣钱 + B 加钱,要么都成功,要么都不做
秒杀库存 + 订单 + 消息队列扣库存成功后,订单创建和消息发送也必须成功
退款订单 + 账户 + 库存订单状态更新、余额退回、库存恢复

二、CAP 理论与 BASE 理论

2.1 CAP 定理

CAP:一个分布式系统最多只能同时满足三项中的两项

  C(Consistency,一致性):所有节点同一时间看到相同的数据
  A(Availability,可用性):每次请求都能获得非错误的响应
  P(Partition Tolerance,分区容错性):网络分区时系统仍能工作

  分布式系统 P 是必须的(网络分区无法避免),
  所以只能在 C 和 A 之间做取舍:

  CP 系统(强一致,牺牲可用性):
    → ZooKeeper、Etcd、Consul
    → 网络分区时,宁可不可用也不返回不一致数据

  AP 系统(高可用,牺牲强一致):
    → Eureka、Nacos(AP 模式)
    → 网络分区时,宁可数据不一致也要保证可用

2.2 BASE 理论

BASE = 对 CAP 中 AP 方案的补充

  BA(Basically Available,基本可用):
    系统出现故障时,允许损失部分可用性,但保证核心功能可用
    例如:下单页面正常,但优惠券查询暂时不可用

  S(Soft State,软状态):
    允许系统中的数据存在中间状态,该状态不影响系统整体可用性
    例如:订单状态为"支付中"(中间状态),最终会变为"已支付"或"已取消"

  E(Eventually Consistent,最终一致性):
    不要求数据实时一致,但保证最终会达到一致
    例如:支付成功后,账户余额可能 1 秒后才更新,但最终一定正确

  BASE 的核心思想:放弃强一致性,追求最终一致性
CAP → BASE 的演进:

  CAP 告诉我们:分布式系统做不到完美的强一致 + 高可用
  BASE 告诉我们:那就退一步,追求最终一致

  大部分互联网业务不需要强一致——
  用户支付成功后,订单状态晚 1 秒变成"已支付",完全可以接受

三、分布式事务四大方案

3.1 方案全景图

方案一致性性能复杂度适用场景
2PC(XA)强一致传统单体拆微服务,跨库事务
TCC强一致资金、交易等核心业务
Saga最终一致长流程、非资金类业务
AT(Seata)最终一致对业务侵入小的场景
可靠消息最终一致异步解耦、通知类场景
最大努力通知最终一致外部系统回调、通知

3.2 方案选择决策树

你的业务场景                              → 推荐方案
─────────────────────────────────────────────────────────────
强一致性、资金交易(转账、支付)            → TCC
强一致性、简单跨库事务(单体拆微服务)       → 2PC(XA)
长流程、非资金类(订单创建→审核→发货)      → Saga
无侵入、快速接入、非强一致                 → AT(Seata)
异步解耦、通知类(短信、邮件)              → 可靠消息
外部系统回调(支付回调、物流通知)           → 最大努力通知

四、2PC(两阶段提交)

4.1 原理

2PC = 协调者 + 参与者,分两个阶段

  阶段一:投票(Prepare)
    协调者 → 所有参与者:能执行吗?
    参与者 → 协调者:能(YES)/ 不能(NO)

  阶段二:提交(Commit)
    如果所有参与者都回复 YES:
      协调者 → 所有参与者:提交!
    如果任意参与者回复 NO:
      协调者 → 所有参与者:回滚!
正常流程:

  协调者                          参与者 A          参与者 B
    │                               │                 │
    │──── Prepare ─────────────────→│                 │
    │──── Prepare ──────────────────────────────────→│
    │                               │                 │
    │←─── YES(锁定资源)──────────│                 │
    │←─── YES(锁定资源)──────────────────────────│
    │                               │                 │
    │──── Commit ──────────────────→│                 │
    │──── Commit ──────────────────────────────────→│
    │                               │                 │
    │←─── ACK ─────────────────────│                 │
    │←─── ACK ─────────────────────────────────────│
    │                               │                 │
    ✅ 事务提交成功

4.2 2PC 的致命问题

❌ 同步阻塞:Prepare 阶段,参与者锁定资源,必须等协调者通知
              如果协调者挂了,参与者一直阻塞,资源无法释放

❌ 单点故障:协调者挂了,整个系统瘫痪
              参与者不知道应该提交还是回滚

❌ 数据不一致:阶段二,协调者发了 Commit 后挂了
              有的参与者收到 Commit 提交了,有的没收到,数据不一致

❌ 性能差:所有参与者都要锁定资源等协调者
              一次事务 = 2 次网络往返 × N 个参与者

4.3 Seata XA 模式(2PC 实现)

yaml
# application.yml
seata:
  tx-service-group: my_tx_group
  service:
    vgroup-mapping:
      my_tx_group: default
  data-source-proxy-mode: XA  # 使用 XA 模式
java
@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private InventoryClient inventoryClient;
    @Autowired
    private AccountClient accountClient;

    @GlobalTransactional  // Seata 分布式事务注解
    public void createOrder(Order order) {
        // ① 创建订单(本地事务)
        orderMapper.insert(order);

        // ② 扣减库存(远程调用,Seata 自动代理)
        inventoryClient.deduct(order.getProductId(), order.getQuantity());

        // ③ 扣减余额(远程调用)
        accountClient.deduct(order.getUserId(), order.getAmount());

        // 任意一步失败,Seata 自动回滚所有操作
    }
}

XA 模式适用场景: 传统单体应用拆分为微服务,数据库支持 XA 协议(MySQL 5.7+、PostgreSQL、Oracle),对一致性要求高但并发量不大。

五、TCC(Try-Confirm-Cancel)

5.1 原理

TCC = 把一次业务操作拆成三个步骤

  Try(尝试):预留资源,检查 + 锁定
    → 冻结库存 100 件
    → 冻结余额 100 元

  Confirm(确认):确认提交,使用预留的资源
    → 扣减冻结的库存
    → 扣减冻结的余额

  Cancel(取消):取消操作,释放预留的资源
    → 恢复冻结的库存
    → 恢复冻结的余额
TCC 完整流程:

  业务方(TM)                   库存服务                    账户服务
      │                            │                          │
      │──── Try ──────────────────→│                          │
      │  freezeStock(1001, 10)     │ 冻结 10 件库存             │
      │←─── OK ───────────────────│                          │
      │                            │                          │
      │──── Try ────────────────────────────────────────────→│
      │  freezeBalance(1001, 100)  │                          │ 冻结 100 元
      │←─── OK ─────────────────────────────────────────────│
      │                            │                          │
      │  ⚠️ 业务校验失败,需要回滚   │                          │
      │                            │                          │
      │──── Cancel ───────────────→│                          │
      │  unfreezeStock(1001, 10)   │ 恢复 10 件库存             │
      │←─── OK ───────────────────│                          │
      │                            │                          │
      │──── Cancel ──────────────────────────────────────────→│
      │  unfreezeBalance(1001,100) │                          │ 恢复 100 元
      │←─── OK ─────────────────────────────────────────────│

5.2 TCC 实现

java
/**
 * 库存服务——TCC 接口
 */
public interface InventoryTccService {

    /**
     * Try:冻结库存
     * @param businessId 业务流水号,用于幂等
     * @param productId  商品 ID
     * @param count      冻结数量
     */
    @Transactional
    boolean tryFreezeStock(String businessId, Long productId, Integer count);

    /**
     * Confirm:确认扣减(使用冻结的库存)
     */
    @Transactional
    boolean confirmFreezeStock(String businessId, Long productId, Integer count);

    /**
     * Cancel:取消冻结(恢复库存)
     */
    @Transactional
    boolean cancelFreezeStock(String businessId, Long productId, Integer count);
}
java
@Service
public class InventoryTccServiceImpl implements InventoryTccService {

    @Autowired
    private InventoryMapper inventoryMapper;
    @Autowired
    private FreezeLogMapper freezeLogMapper;

    @Override
    public boolean tryFreezeStock(String businessId, Long productId, Integer count) {
        // ① 幂等校验:防止重复 Try
        if (freezeLogMapper.exists(businessId)) {
            return true;  // 已处理过,直接返回成功
        }

        // ② 扣减可用库存,增加冻结库存
        // UPDATE inventory SET available = available - ?, frozen = frozen + ?
        // WHERE product_id = ? AND available >= ?
        int rows = inventoryMapper.freeze(productId, count);
        if (rows == 0) {
            return false;  // 库存不足
        }

        // ③ 记录冻结流水(用于 Confirm 和 Cancel 的幂等控制)
        FreezeLog log = new FreezeLog();
        log.setBusinessId(businessId);
        log.setProductId(productId);
        log.setCount(count);
        log.setStatus("FROZEN");
        freezeLogMapper.insert(log);

        return true;
    }

    @Override
    public boolean confirmFreezeStock(String businessId, Long productId, Integer count) {
        // ① 幂等校验
        FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
        if (log == null) return false;
        if ("CONFIRMED".equals(log.getStatus())) return true;  // 已确认

        // ② 扣减冻结库存
        // UPDATE inventory SET frozen = frozen - ?, total = total - ?
        // WHERE product_id = ?
        inventoryMapper.confirmDeduct(productId, count);

        // ③ 更新流水状态
        freezeLogMapper.updateStatus(businessId, "CONFIRMED");

        return true;
    }

    @Override
    public boolean cancelFreezeStock(String businessId, Long productId, Integer count) {
        // ① 幂等校验
        FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
        if (log == null) {
            // 没有 Try 记录,说明 Try 根本没执行,直接返回成功
            return true;
        }
        if ("CANCELLED".equals(log.getStatus())) return true;  // 已取消

        // ② 恢复库存
        // UPDATE inventory SET available = available + ?, frozen = frozen - ?
        // WHERE product_id = ?
        inventoryMapper.unfreeze(productId, count);

        // ③ 更新流水状态
        freezeLogMapper.updateStatus(businessId, "CANCELLED");

        return true;
    }
}
java
/**
 * 订单服务——TCC 事务编排(TM)
 */
@Service
public class OrderTccService {

    @Autowired
    private InventoryTccClient inventoryTccClient;
    @Autowired
    private AccountTccClient accountTccClient;
    @Autowired
    private OrderService orderService;

    public void createOrder(Order order) {
        String businessId = order.getOrderId();  // 使用订单号作为业务流水号

        // ① Try 阶段
        boolean stockFrozen = inventoryTccClient.tryFreezeStock(
            businessId, order.getProductId(), order.getQuantity());
        if (!stockFrozen) {
            throw new RuntimeException("库存不足");
        }

        boolean balanceFrozen = accountTccClient.tryFreezeBalance(
            businessId, order.getUserId(), order.getAmount());
        if (!balanceFrozen) {
            // 余额不足,回滚库存
            inventoryTccClient.cancelFreezeStock(
                businessId, order.getProductId(), order.getQuantity());
            throw new RuntimeException("余额不足");
        }

        try {
            // ② 创建订单
            orderService.createOrder(order);

            // ③ Confirm 阶段:全部 Try 成功,确认提交
            inventoryTccClient.confirmFreezeStock(
                businessId, order.getProductId(), order.getQuantity());
            accountTccClient.confirmFreezeBalance(
                businessId, order.getUserId(), order.getAmount());

        } catch (Exception e) {
            // ④ Cancel 阶段:任意步骤失败,回滚所有
            inventoryTccClient.cancelFreezeStock(
                businessId, order.getProductId(), order.getQuantity());
            accountTccClient.cancelFreezeBalance(
                businessId, order.getUserId(), order.getAmount());
            throw new RuntimeException("下单失败", e);
        }
    }
}

5.3 TCC 的空回滚、幂等、悬挂

分布式事务的三大经典问题:

① 空回滚(Empty Rollback):
  Try 超时,TM 发起 Cancel,但 Try 的请求还没到达 RM
  → RM 收到 Cancel,但没有对应的 Try 记录
  → 解决:Cancel 阶段发现没有 Try 记录时,直接返回成功(允许空回滚)

② 幂等(Idempotent):
  Try/Confirm/Cancel 可能被多次调用(网络重试、TM 重试)
  → 解决:使用 businessId 做唯一键,已处理则直接返回成功

③ 悬挂(Suspension):
  Cancel 先到达且执行完毕,之后 Try 才到达
  → Try 执行了但 Cancel 已经做过了,资源被错误冻结
  → 解决:Cancel 执行后标记 businessId 为已取消,Try 时检查标记,拒绝执行
java
// 悬挂处理:Cancel 后标记,Try 时检查
public boolean tryFreezeStock(String businessId, Long productId, Integer count) {
    // 检查是否已被 Cancel(防止悬挂)
    if (freezeLogMapper.isCancelled(businessId)) {
        return false;  // 已被取消,拒绝 Try
    }
    // ... 正常 Try 逻辑
}

public boolean cancelFreezeStock(String businessId, Long productId, Integer count) {
    // 空回滚处理
    FreezeLog log = freezeLogMapper.getByBusinessId(businessId);
    if (log == null) {
        // 记录一条已取消的日志,防止后续 Try 悬挂
        freezeLogMapper.insertCancelled(businessId);
        return true;  // 空回滚,允许成功
    }
    // ... 正常 Cancel 逻辑
}

5.4 TCC 优缺点

✅ 优点:
  - 强一致性:Confirm 全部成功才算成功,任一失败就全部 Cancel
  - 性能优于 2PC:Try 阶段只锁定资源,不阻塞等待
  - 业务可控:开发者自定义 Try/Confirm/Cancel 逻辑

❌ 缺点:
  - 代码侵入大:每个业务都要写 Try/Confirm/Cancel 三个方法
  - 开发成本高:空回滚、幂等、悬挂都要处理
  - 不适合长流程:Confirm 阶段失败需要人工介入

六、Saga 模式

6.1 原理

Saga = 把一个大事务拆成多个有序的小事务,每个小事务都有对应的补偿操作

  正向:T1 → T2 → T3 → ... → Tn
  补偿:C1 ← C2 ← C3 ← ... ← Cn

  如果 T3 失败,依次执行 C2 → C1 回滚

6.2 两种实现方式

① 编排式(Choreography,事件驱动):
  每个服务做完自己的事,发布事件,下一个服务监听事件

  订单服务 → 创建订单 → 发布"订单已创建"事件
  库存服务 → 监听事件 → 扣库存 → 发布"库存已扣减"事件
  账户服务 → 监听事件 → 扣余额 → 结束

  优点:松耦合,适合简单流程
  缺点:流程分散在各服务中,难以追踪


② 编排式(Orchestration,命令驱动):
  一个 Saga 编排器统一调度所有步骤

  Saga 编排器:
    ① 调用订单服务创建订单
    ② 调用库存服务扣库存
    ③ 调用账户服务扣余额
    ④ 失败时按相反顺序调用补偿操作

  优点:流程集中管理,清晰可追踪
  缺点:编排器成为新的耦合点

6.3 Seata Saga 实现

java
// Saga 状态机定义(JSON 格式)
// 定义流程:创建订单 → 扣库存 → 扣余额
json
{
  "Name": "createOrderSaga",
  "StartState": "CreateOrder",
  "States": {
    "CreateOrder": {
      "Type": "ServiceTask",
      "ServiceName": "orderService.createOrder",
      "CompensateState": "CancelOrder",
      "Next": "DeductStock"
    },
    "CancelOrder": {
      "Type": "ServiceTask",
      "ServiceName": "orderService.cancelOrder",
      "End": true
    },
    "DeductStock": {
      "Type": "ServiceTask",
      "ServiceName": "inventoryService.deductStock",
      "CompensateState": "CompensateStock",
      "Next": "DeductBalance"
    },
    "CompensateStock": {
      "Type": "ServiceTask",
      "ServiceName": "inventoryService.increaseStock",
      "Next": "CancelOrder"
    },
    "DeductBalance": {
      "Type": "ServiceTask",
      "ServiceName": "accountService.deductBalance",
      "CompensateState": "CompensateBalance",
      "End": true
    },
    "CompensateBalance": {
      "Type": "ServiceTask",
      "ServiceName": "accountService.increaseBalance",
      "Next": "CompensateStock"
    }
  }
}
java
@Service
public class OrderSagaService {

    @Autowired
    private SagaStateMachineEngine sagaEngine;

    public void createOrder(Order order) {
        // 启动 Saga 状态机,Seata 自动编排
        sagaEngine.start("createOrderSaga", order);
    }
}

6.4 Saga 优缺点

✅ 优点:
  - 高性能:每个小事务独立提交,不长时间锁定资源
  - 长流程友好:适合多步骤、跨服务的长事务
  - 松耦合:服务间通过事件或编排器通信

❌ 缺点:
  - 最终一致性:不保证实时一致,中间状态可见
  - 补偿复杂:补偿操作需要业务逻辑配合(如退款可能有手续费)
  - 隔离性弱:事务 A 扣了库存但还没提交,事务 B 可能读到中间状态

七、Seata AT 模式

7.1 原理

AT 模式 = 自动挡分布式事务,对业务代码零侵入

  核心机制:两阶段提交 + 自动生成回滚 SQL

  阶段一:
    1. 执行原始 SQL(INSERT / UPDATE / DELETE)
    2. 解析 SQL,获取前镜像(before image)
    3. 执行业务 SQL
    4. 获取后镜像(after image)
    5. 将前后镜像 + 行锁插入 undo_log 表
    6. 提交本地事务(释放数据库连接)

  阶段二(提交):
    1. 删除 undo_log 记录(异步,批量)

  阶段二(回滚):
    1. 从 undo_log 读取前镜像
    2. 校验后镜像是否匹配(防止脏写)
    3. 用前镜像生成反向 SQL 并执行
    4. 删除 undo_log 记录
AT 模式的全局锁机制:

  事务 A:UPDATE product SET stock = stock - 10 WHERE id = 1001
    → Seata 自动加全局锁:lock:product:1001

  事务 B:UPDATE product SET stock = stock - 5 WHERE id = 1001
    → 发现全局锁被占用 → 等待或超时

  事务 A 提交 → 释放全局锁 → 事务 B 获取锁 → 执行

7.2 AT 模式快速上手

xml
<!-- Seata 1.8.0 + Spring Boot 3.x -->
<dependency>
    <groupId>io.seata</groupId>
    <artifactId>seata-spring-boot-starter</artifactId>
    <version>2.1.0</version>
</dependency>
yaml
# application.yml
seata:
  enabled: true
  tx-service-group: default_tx_group
  service:
    vgroup-mapping:
      default_tx_group: default
  registry:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
  config:
    type: nacos
    nacos:
      server-addr: 127.0.0.1:8848
      group: SEATA_GROUP
  data-source-proxy-mode: AT  # AT 模式
sql
-- 每个业务的数据库都需要创建 undo_log 表
CREATE TABLE `undo_log` (
    `id`         BIGINT(20) NOT NULL AUTO_INCREMENT,
    `branch_id`  BIGINT(20) NOT NULL,
    `xid`        VARCHAR(100) NOT NULL,
    `context`    VARCHAR(128) NOT NULL,
    `rollback_info` LONGBLOB NOT NULL,
    `log_status` INT(11) NOT NULL,
    `log_created` DATETIME NOT NULL,
    `log_modified` DATETIME NOT NULL,
    PRIMARY KEY (`id`),
    KEY `idx_unionkey` (`xid`, `branch_id`)
) ENGINE = InnoDB DEFAULT CHARSET = utf8mb4;
java
@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private InventoryClient inventoryClient;
    @Autowired
    private AccountClient accountClient;

    @GlobalTransactional(name = "createOrder", timeoutMills = 300000)
    public void createOrder(Order order) {
        // 业务代码完全不需要修改,和单机事务写法一模一样
        // Seata 自动代理数据源,自动记录 undo_log,自动回滚

        orderMapper.insert(order);                              // 本地事务
        inventoryClient.deduct(order.getProductId(), order.getQuantity());  // 远程调用
        accountClient.deduct(order.getUserId(), order.getAmount());         // 远程调用

        // 任意步骤失败,Seata 自动回滚所有操作
    }
}

7.3 AT 模式优缺点

✅ 优点:
  - 零侵入:业务代码不需要任何修改,和单机事务一样写
  - 自动回滚:Seata 自动生成反向 SQL,不需要手写补偿
  - 性能较好:阶段一就提交本地事务,不长期占用连接

❌ 缺点:
  - 依赖数据库:undo_log 表是必须的,不支持非关系型数据库
  - 全局锁:依赖 Seata Server 的全局锁,有单点风险
  - 隔离性:默认读未提交,可能读到脏数据
  - 仅支持 ACID 数据库:MySQL、PostgreSQL、Oracle

八、可靠消息最终一致性

8.1 原理

可靠消息 = 本地事务 + 消息表 + 消息队列 + 重试机制

  ① 本地事务:执行业务操作 + 往"消息表"插入一条消息(同一个事务)
  ② 定时任务:扫描消息表,将"待发送"的消息发送到 MQ
  ③ 消费者:消费消息,执行业务逻辑
  ④ 确认删除:消费者执行成功后,回调确认,消息表标记为"已发送"

8.2 实现

sql
-- 消息表(本地事务表)
CREATE TABLE `tx_message` (
    `id`          BIGINT AUTO_INCREMENT PRIMARY KEY,
    `message_id`  VARCHAR(64) NOT NULL,          -- 消息 ID(幂等键)
    `topic`       VARCHAR(64) NOT NULL,           -- MQ Topic
    `tag`         VARCHAR(64) DEFAULT NULL,       -- MQ Tag
    `body`        TEXT NOT NULL,                  -- 消息体 JSON
    `status`      VARCHAR(16) NOT NULL DEFAULT 'PENDING',  -- PENDING/SENT/CONFIRMED
    `retry_count` INT DEFAULT 0,                 -- 重试次数
    `max_retry`   INT DEFAULT 10,                -- 最大重试次数
    `next_retry`  DATETIME NOT NULL,             -- 下次重试时间
    `created_at`  DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
    UNIQUE KEY `uk_message_id` (`message_id`)
);
java
@Service
public class OrderService {

    @Autowired
    private OrderMapper orderMapper;
    @Autowired
    private TxMessageMapper txMessageMapper;
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    /**
     * 下单:本地事务 + 消息表
     */
    @Transactional
    public void createOrder(Order order) {
        // ① 业务操作:创建订单
        orderMapper.insert(order);

        // ② 同一事务中,插入消息表
        TxMessage message = new TxMessage();
        message.setMessageId(order.getOrderId());
        message.setTopic("order-topic");
        message.setTag("order-created");
        message.setBody(JSON.toJSONString(order));
        message.setStatus("PENDING");
        message.setNextRetry(LocalDateTime.now());

        txMessageMapper.insert(message);

        // ③ 事务提交后,发送消息(这里简化,实际由定时任务发送)
        // 如果发送失败,定时任务会重试
    }
}
java
/**
 * 定时任务:扫描消息表,发送消息
 */
@Component
public class TxMessageSender {

    @Autowired
    private TxMessageMapper txMessageMapper;
    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    @Scheduled(fixedDelay = 5000)  // 每 5 秒扫描一次
    public void sendPendingMessages() {
        // 查询待发送的消息(状态=PENDING,且到达重试时间)
        List<TxMessage> messages = txMessageMapper.selectPendingMessages(
            "PENDING", LocalDateTime.now(), 100);

        for (TxMessage msg : messages) {
            try {
                // 发送到 RocketMQ
                rocketMQTemplate.syncSend(
                    msg.getTopic() + ":" + msg.getTag(),
                    msg.getBody());

                // 更新状态为已发送
                txMessageMapper.updateStatus(msg.getMessageId(), "SENT");

            } catch (Exception e) {
                // 发送失败,更新重试信息
                int retries = msg.getRetryCount() + 1;
                if (retries > msg.getMaxRetry()) {
                    // 超过最大重试次数,标记为失败,人工介入
                    txMessageMapper.updateStatus(msg.getMessageId(), "FAILED");
                } else {
                    // 指数退避:下次重试 = 当前时间 + 2^retries 秒
                    txMessageMapper.updateRetry(
                        msg.getMessageId(), retries,
                        LocalDateTime.now().plusSeconds((long) Math.pow(2, retries)));
                }
            }
        }
    }
}

8.3 RocketMQ 事务消息(更优方案)

RocketMQ 事务消息 = 普通消息 + 事务状态回查

  ① 发送半消息(Half Message)→ MQ 保存但不可见
  ② 执行本地事务
  ③ 本地事务成功 → 提交半消息 → MQ 变为可见
     本地事务失败 → 回滚半消息 → MQ 删除
  ④ 如果 MQ 没收到确认(超时),MQ 回查本地事务状态
java
@Service
public class OrderService {

    @Autowired
    private RocketMQTemplate rocketMQTemplate;

    public void createOrder(Order order) {
        // RocketMQ 事务消息,一步到位
        rocketMQTemplate.sendMessageInTransaction(
            "order-topic:order-created",           // topic:tag
            MessageBuilder.withPayload(order).build(),  // 消息体
            order.getOrderId()                     // 业务参数
        );
    }

    /**
     * 本地事务执行回调
     */
    @RocketMQTransactionListener
    public class OrderTransactionListener implements RocketMQLocalTransactionListener {

        @Override
        public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            String orderId = (String) arg;
            try {
                // 执行本地事务:创建订单
                orderMapper.insert(parseOrder(msg));
                return RocketMQLocalTransactionState.COMMIT;  // 提交消息
            } catch (Exception e) {
                return RocketMQLocalTransactionState.ROLLBACK;  // 回滚消息
            }
        }

        @Override
        public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
            String orderId = (String) msg.getHeaders().get("orderId");
            // 回查:查询本地事务是否执行成功
            Order order = orderMapper.selectById(orderId);
            return order != null
                ? RocketMQLocalTransactionState.COMMIT
                : RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

九、最大努力通知

9.1 原理

最大努力通知 = 通知方尽最大努力通知,不保证一定成功

  ① 业务操作完成后,通知外部系统
  ② 如果通知失败,按一定策略重试(递增间隔)
  ③ 重试 N 次后仍失败,记录异常,人工介入

  典型场景:支付回调通知
java
@Component
public class PaymentNotifyService {

    @Autowired
    private RestTemplate restTemplate;

    private static final int MAX_RETRIES = 5;
    private static final long[] RETRY_INTERVALS = {1, 5, 15, 30, 60};  // 分钟

    public void notifyMerchant(String orderId, String notifyUrl, String notifyBody) {
        for (int i = 0; i < MAX_RETRIES; i++) {
            try {
                ResponseEntity<String> response = restTemplate.postForEntity(
                    notifyUrl, notifyBody, String.class);

                if (response.getStatusCode().is2xxSuccessful()
                    && "SUCCESS".equals(response.getBody())) {
                    // 通知成功
                    return;
                }
            } catch (Exception e) {
                log.warn("通知商户失败,第 {} 次重试,orderId={}", i + 1, orderId);
            }

            // 等待后重试
            try {
                TimeUnit.MINUTES.sleep(RETRY_INTERVALS[i]);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
                break;
            }
        }

        // 超过最大重试次数,记录异常,人工处理
        log.error("通知商户失败,已超过最大重试次数,orderId={}", orderId);
        // 发送告警
        alertService.send("支付通知失败", orderId);
    }
}

十、六大方案全面对比

维度2PC(XA)TCCSagaAT(Seata)可靠消息最大努力通知
一致性强一致强一致最终一致最终一致最终一致最终一致
性能⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
复杂度⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐⭐
代码侵入
回滚方式自动手动(Cancel)补偿自动(undo_log)补偿重试
锁资源时间
适用场景跨库事务资金交易长流程快速接入异步解耦外部回调
是否依赖数据库是(MySQL)

十一、生产级最佳实践

11.1 幂等设计

分布式事务的每个参与者都必须支持幂等:

  1. 使用全局唯一 ID(xid + branchId)作为幂等键
  2. 数据库唯一索引 + 状态机
  3. 所有 Confirm / Cancel / 补偿操作都支持重复调用

11.2 监控与告警

分布式事务需要完善的监控体系:

  1. 事务成功率:每分钟提交/回滚的比例
  2. 事务耗时:P50 / P99 / P999
  3. undo_log 堆积:未清理的 undo_log 数量
  4. 重试次数:消息重试次数分布
  5. 异常告警:事务失败、超时、回滚失败

11.3 避免长事务

分布式事务的黄金法则:事务越短越好

  1. 不要在事务中调用外部 API(RPC 除外)
  2. 不要在事务中发邮件、发短信
  3. 不要在事务中操作文件、生成 PDF
  4. 超过 1 秒的事务考虑异步化
  5. 超过 10 个参与者的考虑用 Saga 替代 TCC

11.4 定时任务补偿

java
/**
 * 分布式事务补偿任务:处理未完成的异常事务
 */
@Component
public class TransactionCompensationJob {

    @Scheduled(fixedDelay = 60000)  // 每分钟执行
    public void compensate() {
        // ① 处理长时间未确认的 TCC Try 记录(超过 30 分钟)
        List<FreezeLog> timeoutLogs = freezeLogMapper.selectTimeoutLogs(
            "FROZEN", LocalDateTime.now().minusMinutes(30));
        for (FreezeLog log : timeoutLogs) {
            // 自动 Cancel 释放资源
            inventoryService.cancelFreezeStock(
                log.getBusinessId(), log.getProductId(), log.getCount());
        }

        // ② 处理长时间未清理的 undo_log(超过 7 天)
        undoLogMapper.deleteExpired(LocalDateTime.now().minusDays(7));

        // ③ 处理失败的消息(超过 24 小时)
        List<TxMessage> failedMessages = txMessageMapper.selectFailedMessages(
            "FAILED", LocalDateTime.now().minusHours(24));
        for (TxMessage msg : failedMessages) {
            // 发送告警,人工处理
            alertService.send("消息发送失败", msg.getMessageId());
        }
    }
}

十二、面试要点

Q1:分布式事务有哪些实现方案?

  六大方案:2PC(XA)强一致、TCC(Try-Confirm-Cancel)强一致、Saga 补偿模式、Seata AT 模式(自动回滚)、可靠消息最终一致性、最大努力通知。按一致性要求从强到弱选择。

Q2:TCC 和 AT 模式有什么区别?

  TCC 需要手写 Try/Confirm/Cancel 三个方法,代码侵入大但灵活性高,适合资金类;AT 模式零侵入,Seata 自动生成反向 SQL,但依赖数据库,适合快速接入。

Q3:TCC 的空回滚、幂等、悬挂是什么?

  空回滚:Cancel 时 Try 还没执行;幂等:同一操作重复执行;悬挂:Cancel 先到,Try 后到。都需要通过 businessId + 状态机处理。

Q4:Seata 的 AT 模式隔离级别是什么?

  默认读未提交。如果业务要求读已提交,Seata 提供了 @GlobalLock 注解 + SELECT FOR UPDATE 来保证。

Q5:RocketMQ 事务消息和本地消息表怎么选?

  RocketMQ 事务消息更优雅,不需要额外维护消息表,但绑定 RocketMQ;本地消息表方案通用,任何 MQ 都适用,但需要定时任务扫描。

Q6:什么场景不适合用分布式事务?

  1) 非核心业务,允许少量不一致(如统计、日志);2) 超高并发场景(秒杀),分布式事务会成为瓶颈;3) 跨公司/跨系统,无法控制对方数据库。这些场景用异步 + 对账 + 人工兜底更实际。