This repository has been archived on 2026-05-24. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
all-in-kingsoft/hhs/MS/03-数据一致性/02-分布式事务.md
T
2026-05-18 00:17:59 +08:00

19 KiB
Raw Blame History

tags, create time
tags create time
microservice
distributed-transactions
saga
tcc
outbox
seata
idempotency
dlq
2026-05-17 15:00

分布式事务

概述

每个服务拥有独立数据库,跨服务的"一次操作"实际上涉及多个本地事务。如何保证这些本地事务要么全部成功、要么全部回滚,就是分布式事务要解决的问题。

graph LR
    A["用户下单"] --> B["订单服务<br/>写库"]
    B --> C["库存服务<br/>扣减"]
    C --> D["支付服务<br/>扣款"]
    
    style A fill:#e3f2fd
    style B fill:#fff3e0
    style C fill:#fff3e0
    style D fill:#fff3e0

分布式事务方案全景

方案 一致性级别 性能 复杂度 适用场景
本地事务 + MQ 事件 最终一致 ⭐⭐⭐⭐⭐ ⭐ 绝大多数场景
Saga 模式 最终一致 ⭐⭐⭐ ⭐⭐ 长流程业务
TCC 强最终一致 ⭐⭐⭐ ⭐⭐⭐⭐ 对一致性要求较高的场景
AT 模式 (Seata) 伪强一致 ⭐⭐ ⭐ 不想改业务代码时

[!tip] 选择策略 先默认用本地事务 + 异步事件(最简单、最高效),只有在业务明确需要 Saga 或 TCC 时才升级。80% 的场景,本地事务 + MQ 就足够了。

方案一:本地事务 + 消息队列

核心思想:将"数据变更 + 发消息"合并到一个本地事务中。

Outbox 模式(推荐)

sequenceDiagram
    participant App as 应用服务
    participant DB as 数据库
    participant Outbox as Outbox 表
    participant MQ as 消息队列
    participant Sub as 订阅方

    App->>DB: BEGIN 事务
    App->>DB: 写业务数据
    App->>Outbox: 写入待发送消息
    App->>DB: COMMIT
    
    loop 定时任务
        MQ->>Outbox: 扫描 status='pending'
        Outbox-->>MQ: 返回消息列表
        MQ->>Sub: 投递消息
        MQ->>Outbox: 更新为 'sent'
    end
    
    Note right of Sub: 至少一次投递 + 消费者幂等

Outbox 表建表示例

CREATE TABLE outbox (
    id          BIGSERIAL PRIMARY KEY,
    topic       VARCHAR(255) NOT NULL,     -- 消息主题
    payload     JSONB      NOT NULL,       -- 消息体
    status      VARCHAR(20)  NOT NULL DEFAULT 'pending', -- pending / sent / failed
    error_msg   TEXT,                       -- 失败原因
    created_at  TIMESTAMP    NOT NULL DEFAULT NOW(),
    sent_at     TIMESTAMP
);

-- 加速定时扫描查询
CREATE INDEX idx_outbox_pending ON outbox (status, created_at) 
WHERE status = 'pending';
// Go 示例:出事务内同时写业务数据和 outbox 记录
tx, _ := db.Begin()
tx.Exec("INSERT INTO orders (user_id, total) VALUES ($1, $2)", userID, total)
tx.Exec(`INSERT INTO outbox (topic, payload, status) 
         VALUES ('order.created', $1, 'pending')`, jsonPayload)
tx.Commit()

// 后台 goroutine 轮询并推送
func outboxWorker(ctx context.Context, ticker *time.Ticker) {
    for {
        select {
        case <-ctx.Done():
            return
        case <-ticker.C:
            sendPendingMessages(ctx)
        }
    }
}

消费幂等性

[!warning] 关键保障 消息可能重复投递(网络超时、MQ 重投),消费者必须做到幂等——处理一次和处理多次的结果完全相同。

三种常见策略:

策略 实现方式 适用场景
数据库唯一约束 msg_id 做 UNIQUE 最可靠,推荐首选
Redis 去重键 SET dedup:{msg_id} 1 NX EX 86400 高吞吐场景
乐观锁版本控制 UPDATE SET qty = qty - N WHERE version = V 金额调整类操作
// 推荐方案:利用 UNIQUE 约束做幂等保障
_, err := db.Exec(`
    INSERT INTO order_events (msg_id, order_id, action, amount)
    VALUES ($1, $2, $3, $4)
    ON CONFLICT (msg_id) DO NOTHING
`, msgID, orderID, action, amount)

if !isUniqueViolation(err) {
    log.Error("process message failed", err)
    return
}
// 执行业务逻辑——到这里说明消息是新到达的

事务消息 vs Outbox

[!question] 思考一下 Outbox 模式有一个天然缺点:定时轮询有延迟(通常秒级)。如果你需要毫秒级的延迟敏感型通知(如支付结果实时推送到用户端),该怎么办?

如果使用的是 RocketMQ,可以绕过 Outbox 模式直接用事务消息:

// 发送事务消息
txMsg := rocketmq.NewTransactionMessage("order-created", payload)
localTx := &MyLocalTxChecker{}

// half 消息发送 → 本地事务执行 → 提交/回查
res, _ := producer.SendMessageInTransaction(txMsg, localTx)

RocketMQ 事务消息的完整交互流程:

sequenceDiagram
    participant Producer as 生产者
    participant Broker as MQ Broker
    participant DB as 业务数据库
    participant Consumer as 消费者

    Producer->>Broker: ① 发送 Half Message (消费者不可见)
    Broker-->>Producer: ACK

    Producer->>DB: ② 执行本地事务
    alt 本地事务成功
        Producer->>Broker: ③a Commit (消息对消费者可见)
    else 本地事务失败
        Producer->>Broker: ③b Rollback (消息丢弃)
    end

    Note over Producer,Broker: ④ 如果步骤 2 超时未回复
    Broker->>Producer: ④ 回查请求 (CheckLocalTx)
    Producer->>DB: 查询本地事务状态
    DB-->>Producer: 返回 COMMIT / ROLLBACK
    Producer->>Broker: 回查结果

[!important] 回查机制的意义 当生产者执行本地事务过程中发生崩溃或网络中断,Broker 收不到 Commit/Rollback 指令,就会通过回查主动询问生产者。这确保了 half 消息最终一定会收敛到可见或被丢弃,不会出现悬空状态。

可靠性要点清单:

  • ✅ 业务数据库和 Outbox 必须在同一个本地事务中写入
  • ✅ 消费者必须有幂等保护(否则至少一次投递 = 无限次重复)
  • ✅ 定时扫描任务需要设置最大重试次数,超过后进入人工处理
  • ✅ 生产环境建议配置 DLQ(死信队列),将连续消费失败的消息归档

方案二:Saga 模式

Saga 适用于跨多个服务的长流程操作,将大事务拆成一系列本地小事务,每个步骤都有对应的补偿操作。

编排式 vs 编舞式

graph TB
    subgraph ORCHESTRATION["编排式 — Coordinator 中心化"]
        CO[Coordinator] --> S1[OrderSvc: Create]
        CO --> S2[InventorySvc: Reserve]
        CO --> S3[PaymentSvc: Charge]
        
        S1 -.->|失败→Cancel| CO
        S2 -.->|失败→Cancel| CO
        S3 -.->|失败→Cancel| CO
    end

    subgraph CHOREOGRAPHY["编舞式 — 事件驱动"]
        E1[OrderCreated] --> S11[OrderService]
        S11 --> E2[StockReserved]
        E2 --> S12[InventoryService]
        S12 --> E3[PaymentCharged]
        E3 --> S13[PaymentService]
        
        S13 -.-> E4[PaidFailed] -.-> S11
    end
维度 编排式 编舞式
控制流 中心化 Coordinator 各服务通过事件自发响应
可观测性 ✅ 集中管理全流程 ❌ 流程散布在各服务
耦合度 依赖 Coordinator 服务间仅感知事件
适合规模 5~15 步的 Saga 简单链路 (< 5 步)

Saga 补偿设计原则

每个正向操作必须有对应的反向补偿:

正向操作 补偿操作
创建订单 取消订单
预留库存 释放库存
扣款 退款
发送通知 撤销通知(一般不需要)

[!warning] 补偿操作的幂等性 补偿操作也必须幂等——CancelOrder 可能被触发多次。用订单状态的流转来保证(如只有 PENDING 才能转到 CANCELLED)。

编排式实现示例 (Go)

编排式的核心是一个 Saga Coordinator,它维护着当前步骤的执行状态和已完成的正向/反向操作列表:

// Step 定义一个 Saga 步骤
type Step struct {
    Name        string
    Execute     func(ctx context.Context) error  // 正向操作
    Compensation func(ctx context.Context) error  // 补偿操作
}

// Orchestrator 编排器
type Orchestrator struct {
    steps       []Step
    executed    []string // 已完成步骤名,用于补偿时逆向遍历
}

func (o *Orchestrator) AddStep(name string, exec, compens func(context.Context) error) {
    o.steps = append(o.steps, Step{Name: name, Execute: exec, Compensation: compens})
}

// Execute 按顺序执行所有步骤,任一步骤失败则逆向补偿
func (o *Orchestrator) Execute(ctx context.Context) error {
    for i, step := range o.steps {
        if err := step.Execute(ctx); err != nil {
            log.Error("step failed", z.String("step", step.Name), z.Err(err))
            // 逆向补偿:从后往前执行已完步骤的补偿操作
            o.rollback(ctx, i-1)
            return fmt.Errorf("saga failed at %s: %w", step.Name, err)
        }
        o.executed = append(o.executed, step.Name)
    }
    return nil
}

func (o *Orchestrator) rollback(ctx context.Context, upTo int) {
    for i := len(o.executed) - 1; i >= 0 && i <= upTo; i-- {
        // 找到对应步骤的补偿函数并执行
        step := o.steps[i]
        if err := step.Compensation(ctx); err != nil {
            log.Warn("compensation also failed — need manual intervention", z.String("step", step.Name))
            // 补偿失败必须告警!人工介入处理
        }
    }
}

使用方式——构建一个下单 Saga:

saga := &Orchestrator{}

saga.AddStep("create_order",
    func(ctx context.Context) error { return orderSvc.Create(ctx, req) },       // 正向
    func(ctx context.Context) error { return orderSvc.Cancel(ctx, orderID) },   // 反向
)
saga.AddStep("reserve_stock",
    func(ctx context.Context) error { return stockSvc.Reserve(ctx, productIDs) },
    func(ctx context.Context) error { return stockSvc.Release(ctx, productIDs) },
)
saga.AddStep("charge_payment",
    func(ctx context.Context) error { return paySvc.Charge(ctx, amount) },
    func(ctx context.Context) error { return paySvc.Refund(ctx, orderID) },
)

err := saga.Execute(ctx)

[!note] 关键理解:补偿的局限性

  • 补偿不是撤销:取消订单不会把商品退回到原始状态,而是创建一条反向业务记录
  • 补偿可能失败:如果退款接口也挂了,你需要重试 + 告警 + 人工介入机制
  • 时间窗口问题:如果正向前向到了第 5 步、补偿到了第 2 步时第 2 步也失败了,第 3~4 步的正向操作已经发生且无法撤销

方案三:TCC (Try-Confirm-Cancel)

TCC 在每个事务步骤中实现三个接口:

Try:   预留资源(冻结余额 / 锁定库存)
Confirm: 确认使用资源(正式扣减 / 正式锁定)
Cancel:  释放资源(解冻余额 / 解锁库存)

TCC 时序图

sequenceDiagram
    participant Orch as Coordinator
    participant O as Order Service
    participant I as Inventory Service
    participant P as Payment Service
    
    Orch->>O: Try(CreateOrder)
    O-->>Orch: OK
    
    Orch->>I: Try(ReserveStock)
    I-->>Orch: OK
    
    Orch->>P: Try(ChargeBalance)
    P-->>Orch: OK
    
    Orch->>O: Confirm
    Orch->>I: Confirm
    Orch->>P: Confirm
    
    Note over Orch,P: 全部 Confirm → 事务完成

TCC vs Saga 对比

维度 TCC Saga
一致性强度 较强(资源被占用期间不允许其他事务使用) 较弱(中间态数据可见)
开发成本 高(每个业务方法实现 Try/Confirm/Cancel) 低(只需正向 + 反向操作)
性能 中(需要两阶段提交) 中高(单阶段本地事务)
适用场景 资金、库存等高敏感业务 订单流程、审批流等业务链

方案四:AT 模式 (Seata)

[!question] 什么时候该用 Seata?

如果你的团队不愿或没有能力为每个业务方法编写 TCC 接口,但又有跨库事务需求——Seata 的 AT 模式就是为此设计的。它是"零侵入"的最强候选方案。

Seata AT 模式的核心原理是 全局锁 + 二阶段提交,业务代码不需要任何改动:

第一阶段(Global Lock):
  1. 拦截 SQL → 解析前后镜像
  2. 申请全局锁(锁住被修改的行)
  3. 本地事务执行(含 undo_log 写入)
  4. 提交本地事务(全局锁暂不释放)

第二阶段(Global Commit/Rollback):
  Commit: 异步删除 undo_log,释放全局锁
  Rollback: 利用 undo_log 生成反向 SQL,回滚并释放全局锁

AT 模式时序图

sequenceDiagram
    participant TC as Transaction Coordinator<br/>(Seata Server)
    participant TM as Transaction Manager<br/>(应用 A)
    participant RS1 as Resource 1<br/>(订单 DB)
    participant RS2 as Resource 2<br/>(库存 DB)

    TM->>TC: 开启全局事务 (xid)
    TM->>RS1: BEGIN; UPDATE orders SET ...
    Note over RS1: 自动捕获 BEFORE/AFTER 镜像<br/>写入 undo_log
    
    TM->>RS2: BEGIN; UPDATE stock SET qty = qty - 1
    Note over RS2: 同样写 undo_log
    
    TM->>RS1: COMMIT (本地事务完成)
    TM->>RS2: COMMIT (本地事务完成)
    
    TM->>TC: 报告二阶段提交完成
    TC-->>RS1: 异步清理 undo_log
    TC-->>RS2: 异步清理 undo_log

AT 与 XA 的区别

很多人混淆 AT 和 XA——它们都是两阶段提交,但实现方式完全不同:

维度 XA 模式 AT 模式 (Seata)
锁粒度 数据库级(长事务锁) 行级(短期全局锁)
隔离级别 需要 SERIALIZABLE 基于脏读 + 全局锁
性能 较差(锁持有时间长) 较好(本地事务快速提交后异步清理)
侵入性 需配置 DataSource Proxy 零侵入,自动代理 JDBC

使用注意事项

// Seata Go SDK — 几乎零侵入
import "github.com/seata/seata-sdk-go-client/tx"

func CreateOrder(ctx context.Context, req *Request) error {
    // 只需添加这一个注解
    tx.GlobalTransaction()
    
    // 下面的代码完全不变——Seata 自动代理数据源
    orderRepo.Create(ctx, &order)
    stockRepo.Decrease(ctx, productID, qty)
    
    return nil
}

[!warning] 生产环境避坑清单

  1. undo_log 表必须建:seata_undo_log 是 Seata 自动创建的,用于存储前后镜像。务必保证这个表可用(通常和业务库在同一实例)
  2. 不要混用本地事务和全局事务:同一个方法里同时出现 @Transactional 和 @GlobalTransactional 会导致行为不可预测
  3. 异常处理:全局事务中抛出的异常会被 Seata 捕获并触发回滚,确保异常能正确向上传递
  4. 性能考量:全局锁虽然比 XA 短,但在高并发场景下仍然是瓶颈。优先用 MQ 最终一致,实在做不到再上 AT
  5. 只支持部分数据库:MySQL、PostgreSQL、Oracle 支持良好;TiDB 等 NewSQL 数据库兼容性有限

[!tip] 选型建议

AT 模式最适用的场景:遗留系统微服务化改造。当原有单体拆分成多个服务时,原有的 @Transactional 跨库调用突然失效,引入 AT 模式可以快速过渡——等业务稳定后再逐步迁移到 MQ 事件驱动架构。

补充保障:重试与死信队列

分布式系统中,没有哪个组件永远可靠。消息消费失败、网络闪断、下游超时——这些都不可避免。你需要一套完整的错误处理链路:

flowchart LR
    MQ["消息队列"] --> C["消费者"]
    C --> FAIL{消费成功}
    FAIL -- 是 --> OK["正常结束"]
    FAIL -- 否 --> RETRY{重试次数未满}
    RETRY -- 是 --> DELAY["延迟重试<br/>指数退避"]
    DELAY --> C
    RETRY -- 否 --> DLQ["进入死信队列<br/>人工介入"]
    
    DLQ --> ALERT["告警通知"]
    DLQ --> REPLAY["手动重放或修复后补发"]

三种重试策略对比

策略 适用场景 示例
立即重试 瞬时故障(连接池未就绪) 等 100ms 再试 3 次
延迟重试 需要时间窗口等待恢复 1s → 2s → 4s → 8s
定时任务补偿 批量漏掉的消息 每 5 分钟扫描 pending 表

死信队列 (DLQ) 设计

// 消费者主循环:带重试和 DLQ 保护
func (c *consumer) Process(ctx context.Context, msg Message) error {
    for attempt := 0; attempt < c.maxRetries; attempt++ {
        err := c.handle(ctx, msg)
        if err == nil {
            return nil // 成功
        }
        
        if isRetryable(err) {
            time.Sleep(time.Duration(attempt+1) * time.Second)
            continue // 重试
        }
        
        // 非可重试错误,直接进 DLQ
        return c.sendToDLQ(msg, err)
    }
    
    // 超过最大重试次数,进 DLQ
    return c.sendToDLQ(msg, fmt.Errorf("exhausted %d retries", c.maxRetries))
}

[!quote] 运维心态

"一个没有死信队列的消息消费者,就像一辆没有备胎的车——任何一次不可预期的故障都会让你抛锚在路上。"

如何选择分布式事务方案?

[!question]- 决策树

面对跨服务的数据一致性问题,你不需要每次都重新选型。按照下面的思路逐层判断即可:

flowchart TD
    START["跨服务写操作"] --> Q1{核心资金或资产}
    
    Q1 -- 是且要求强一致 --> TCC["TCC"]
    Q1 -- 否 --> Q2{流程步数超过5}
    
    Q2 -- 是长流程 --> SAGA["Saga"]
    Q2 -- 否简单链路 --> Q3{容忍秒级延迟}
    
    Q3 -- 可以 --> MQ["本地事务+MQ(Outbox)"]
    Q3 -- 不可以毫秒级 --> TXMSG["RocketMQ事务消息"]
    
    Q3 -- 无法接受最终一致 --> SEATA["Seata AT模式"]
    
    TCC --> END["完成选择"]
    SAGA --> END
    MQ --> END
    TXMSG --> END
    SEATA --> WARNING["过渡方案,后续应迁移至事件驱动架构"]
    WARNING --> END
    
    style TCC fill:#e8f5e9,stroke:#4caf50,stroke-width:2px
    style SAGA fill:#fff3e0,stroke:#ff9800,stroke-width:2px
    style MQ fill:#e3f2fd,stroke:#1976d2,stroke-width:2px
    style TXMSG fill:#e3f2fd,stroke:#1976d2,stroke-width:2px
    style SEATA fill:#fce4ec,stroke:#e91e63,stroke-width:2px

决策速查矩阵

你的业务特征 推荐方案 理由
日订单量百万级,通知类场景 Outbox + MQ 性能最好,开发成本最低
支付/清算/账务核心链路 TCC 资源占用期间不允许其他事务修改
审批流、物流状态流转(多步) Saga 编排式 流程清晰,补偿机制完备
已有 RocketMQ,对延迟敏感 事务消息 绕过 Outbox 轮询延迟
遗留系统快速拆分 Seata AT 零侵入,快速过渡
只需要"最终一致" 本地事务 + MQ 80% 的场景这就是答案

[!tip] 黄金法则

"能用异步事件解决的,不要用两阶段提交;能本地事务 + MQ 解决的,不要上框架。"

复杂度越高,出问题的概率越大。每一层抽象都引入了新的故障点——协调器挂了怎么办?补偿逻辑写错了怎么发现?回滚又失败了谁来做?保持简单是最好的工程纪律。

关联笔记