分布式事务工程实践:从 2PC 到 Saga 的架构演进与生产落地

在微服务架构中,一个业务操作往往横跨多个服务与数据库。如何保证跨服务数据一致性,是分布式系统中最古老也最具挑战性的问题之一。本文从工程实践出发,深入分析 2PC、TCC、Saga、事务消息等主流方案的实现原理、故障模式与生产级优化策略。

一、问题的本质:CAP 与一致性的光谱

传统单机数据库通过 WAL(Write-Ahead Logging)和锁机制完美实现 ACID。但在分布式环境下,网络分区不可避免,根据 CAP 定理,我们在 Partition Tolerance 前提下,只能在 Consistency 与 Availability 之间权衡。

实际上,工程中很少需要真正的"强一致性"。大多数场景能接受的是:

  • 最终一致性(Eventual Consistency):短暂不一致,但保证最终对齐
  • 因果一致性(Causal Consistency):因果相关的操作保持顺序
  • 会话一致性(Session Consistency):同一会话内读写一致

因此,分布式事务的核心工程问题转化为:在特定一致性级别下,如何设计最小侵入、最高可用的补偿机制。

二、两阶段提交(2PC):经典但脆弱的协约

2.1 协议流程

2PC 将事务分为两个阶段:

  1. Prepare 阶段:协调者询问所有参与者是否可以提交,参与者执行事务但不提交,写入 undo/redo 日志,锁定资源,回复 Yes/No。
  2. Commit 阶段:如果全部 Yes,协调者发 Commit;否则发 Abort。参与者执行并释放锁。
# 简化的 2PC 协调者伪代码
class TwoPhaseCoordinator:
    def __init__(self, txn_id, participants):
        self.txn_id = txn_id
        self.participants = participants
        self.state = "INIT"

    async def execute(self, operations):
        # Phase 1: Prepare
        self.state = "PREPARING"
        votes = []
        for op in operations:
            try:
                vote = await op.participant.prepare(op)
                votes.append(vote)
            except Exception as e:
                votes.append("NO")
                logging.warning(f"Participant {op.participant.id} prepare failed: {e}")

        # 决策
        if all(v == "YES" for v in votes):
            self.state = "COMMITTING"
            # Phase 2: Commit
            for op in operations:
                await op.participant.commit(op)
            self.state = "COMMITTED"
        else:
            self.state = "ABORTING"
            for op in operations:
                await op.participant.rollback(op)
            self.state = "ABORTED"

2.2 致命缺陷:阻塞与单点

2PC 的工程困境在于:

  • 同步阻塞:参与者在 Prepare 后持有锁等待协调者决策,协调者故障会导致长时间阻塞
  • 脑裂场景:网络分区时少数派参与者可能单方面 Abort,而多数派已 Commit
  • 协调者单点:协调者故障且无双机热备时,整个事务系统不可用

Google Spanner 通过 Paxos 组复制解决了协调者单点问题,但其 TrueTime API 依赖硬件时钟同步,工程门槛极高。对于一般业务系统,2PC 更适合同机房、低延迟、参与方少的内部场景(如 MySQL XA)。

2.3 生产优化:超时与推测执行

实际部署 2PC 时,必须加入以下工程优化:

  • 参与者推测超时:若 Prepare 后 T 秒未收到决策,可查询协调者状态或 Abort 自身
  • 协调者日志持久化:决策写入 WAL 后再发送,确保重启后可恢复
  • 参与者幂等设计:Commit/Rollback 可重复执行

三、TCC:业务层面的资源预留

TCC(Try-Confirm-Cancel)将分布式事务下放到业务层,避免了数据库锁竞争,适用于跨异构存储的场景。

3.1 三阶段语义

  • Try:资源冻结/预留,检查业务规则(如余额>100),但不实际扣减
  • Confirm:真正执行操作,使用 Try 阶段的预留资源
  • Cancel:释放预留资源,执行补偿

以电商下单为例:

// 库存服务 TCC 实现
type InventoryTCC struct {
    db *sql.DB
}

func (t *InventoryTCC) Try(ctx context.Context, req *TryRequest) error {
    return t.db.WithTx(ctx, func(tx *sql.Tx) error {
        // 检查库存
        var available int
        err := tx.QueryRow("SELECT available FROM inventory WHERE sku = ? FOR UPDATE", req.SKU).Scan(&available)
        if err != nil {
            return err
        }
        if available < req.Quantity {
            return ErrInsufficientStock
        }
        // 冻结库存
        _, err = tx.Exec("UPDATE inventory SET available = available - ?, frozen = frozen + ? WHERE sku = ?", 
            req.Quantity, req.Quantity, req.SKU)
        // 记录 Try 幂等键
        _, err = tx.Exec("INSERT INTO tcc_record(txn_id, status, sku, qty) VALUES(?, 'TRY', ?, ?)",
            req.TxnID, req.SKU, req.Quantity)
        return err
    })
}

func (t *InventoryTCC) Confirm(ctx context.Context, txnID string) error {
    return t.db.WithTx(ctx, func(tx *sql.Tx) error {
        // 将 frozen 转为实际扣减
        _, err := tx.Exec(`UPDATE inventory SET frozen = frozen - qty 
            FROM tcc_record WHERE tcc_record.txn_id = ? AND inventory.sku = tcc_record.sku`,
            txnID)
        _, err = tx.Exec("UPDATE tcc_record SET status = 'CONFIRMED' WHERE txn_id = ?", txnID)
        return err
    })
}

func (t *InventoryTCC) Cancel(ctx context.Context, txnID string) error {
    return t.db.WithTx(ctx, func(tx *sql.Tx) error {
        // 释放冻结库存
        _, err := tx.Exec(`UPDATE inventory SET available = available + qty, frozen = frozen - qty 
            FROM tcc_record WHERE tcc_record.txn_id = ? AND inventory.sku = tcc_record.sku`,
            txnID)
        _, err = tx.Exec("UPDATE tcc_record SET status = 'CANCELLED' WHERE txn_id = ?", txnID)
        return err
    })
}

3.2 工程难点:悬挂与空回滚

TCC 有三个典型陷阱:

  • 空回滚:Try 未执行就收到 Cancel,Cancel 必须识别并记录,避免后续 Confirm 执行
  • 悬挂:Cancel 先于 Try 执行(网络乱序),Try 执行后永远不被 Confirm,导致资源泄漏
  • 幂等:Confirm/Cancel 可能重复调用,需保证结果一致

解决方案是引入三张状态表:

CREATE TABLE tcc_transaction (
    txn_id      VARCHAR(64) PRIMARY KEY,
    status      ENUM('TRYING','CONFIRMING','CANCELING','CONFIRMED','CANCELLED'),
    expire_at   TIMESTAMP,
    created_at  TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE tcc_branch (
    id          BIGINT AUTO_INCREMENT PRIMARY KEY,
    txn_id      VARCHAR(64) NOT NULL,
    resource_id VARCHAR(64) NOT NULL,  -- 参与者标识
    status      ENUM('TRY','CONFIRM','CANCEL'),
    snapshot    JSON,                   -- 业务快照,用于幂等
    UNIQUE KEY uk_txn_resource (txn_id, resource_id)
);

-- 索引:用于超时扫描
CREATE INDEX idx_tcc_txn_expire ON tcc_transaction(status, expire_at);

空回滚防护:Cancel 时先插入 tcc_transaction 记录(status=CANCELING),Try 时若发现已有 CANCELING 记录,则拒绝执行。

悬挂防护:Try 时先检查是否已有 CANCELING/CANCELLED 记录,若有则直接 Cancel 自身。

四、Saga:长事务的最终一致性编排

Saga 模式将长事务拆分为一系列本地事务,每个本地事务触发下一个,失败时执行补偿操作。它适合业务流程长、涉及外部系统(如第三方支付)的场景。

4.1 两种编排方式

编排式 Saga(Choreography):各服务通过事件驱动,无中心协调者。

OrderService → OrderCreated → PaymentService
                                      ↓
            InventoryService ← PaymentSucceeded
                    ↓
            ShippingService ← StockReserved

协调式 Saga(Orchestration):中央编排器(Saga Orchestrator)协调各步骤。

// 协调式 Saga 编排器
type SagaOrchestrator struct {
    store   SagaStore
    executor StepExecutor
}

type SagaDefinition struct {
    Name   string
    Steps  []SagaStep
}

type SagaStep struct {
    Name        string
    Action      func(ctx context.Context, data *SagaData) error
    Compensation func(ctx context.Context, data *SagaData) error
}

func (s *SagaOrchestrator) Execute(ctx context.Context, def SagaDefinition, input *SagaData) error {
    sagaID := generateSagaID()
    s.store.CreateSaga(ctx, sagaID, def.Name, RUNNING)

    for i, step := range def.Steps {
        // 记录步骤开始
        s.store.StepStart(ctx, sagaID, step.Name)

        err := step.Action(ctx, input)

        if err != nil {
            s.store.StepFailed(ctx, sagaID, step.Name, err)
            // 逆序补偿已执行的步骤
            for j := i - 1; j >= 0; j-- {
                if compErr := def.Steps[j].Compensation(ctx, input); compErr != nil {
                    s.store.CompensationFailed(ctx, sagaID, def.Steps[j].Name, compErr)
                    // 补偿失败:记录状态,等待人工介入或自动重试
                    s.store.SetSagaStatus(ctx, sagaID, COMPENSATION_FAILED)
                    return fmt.Errorf("compensation failed at step %s: %w", def.Steps[j].Name, compErr)
                }
                s.store.StepCompensated(ctx, sagaID, def.Steps[j].Name)
            }
            s.store.SetSagaStatus(ctx, sagaID, COMPENSATED)
            return fmt.Errorf("saga failed at step %s: %w", step.Name, err)
        }

        s.store.StepCompleted(ctx, sagaID, step.Name)
    }

    s.store.SetSagaStatus(ctx, sagaID, COMPLETED)
    return nil
}

4.2 工程要点

  • 幂等性:每个步骤和补偿操作都必须是幂等的,编排器可能重试
  • 可交换性:步骤应尽量设计为"正向执行 = 补偿的逆",避免状态泄漏
  • 超时管理:长时间 Saga 需设计超时中断与人工介入机制
  • 可视化:生产环境必须实现 Saga 状态机监控,定位卡在哪个步骤

Seata 框架的 Saga 状态机设计器是一个不错的工程参考,它将 Saga 定义为 JSON 状态图,支持可视化编排与持久化执行。

五、事务消息:异步一致性的工程利器

事务消息是 RocketMQ、Kafka 等 MQ 提供的一种特殊能力,用于解决"本地事务执行"与"消息投递"的一致性问题。它适用于最终一致性可接受、对延迟不敏感的场景(如订单创建后发送通知、积分变更等)。

5.1 RocketMQ 事务消息实现

// RocketMQ 事务消息生产者
@Bean
public TransactionMQProducer transactionProducer() {
    TransactionMQProducer producer = new TransactionMQProducer("tx_producer_group");
    producer.setNamesrvAddr("localhost:9876");
    producer.setTransactionListener(new TransactionListener() {

        // 执行本地事务
        @Override
        public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
            try {
                // 执行数据库操作
                orderService.createOrder((Order) arg);
                return LocalTransactionState.COMMIT_MESSAGE;
            } catch (Exception e) {
                return LocalTransactionState.ROLLBACK_MESSAGE;
            }
        }

        // 回查本地事务状态(MQ 未收到确认时触发)
        @Override
        public LocalTransactionState checkLocalTransaction(MessageExt msg) {
            String orderId = msg.getKeys();
            Order order = orderService.findByOrderId(orderId);
            if (order != null) {
                return LocalTransactionState.COMMIT_MESSAGE;
            }
            return LocalTransactionState.ROLLBACK_MESSAGE;
        }
    });
    return producer;
}

5.2 Kafka 事务实现

Kafka 通过两阶段提交变体实现跨分区原子写入:

// Kafka 事务生产者
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-service-tx-001");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 幂等性必须开启

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();

try {
    producer.beginTransaction();
    // 消费-处理-生产模式(consume-transform-produce)
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
    for (ConsumerRecord<String, String> record : records) {
        OrderEvent event = processOrder(record);
        producer.send(new ProducerRecord<>("order-events", event.getOrderId(), event.toJson()));
    }
    // 提交 consumer offsets 作为事务的一部分
    producer.sendOffsetsToTransaction(currentOffsets, consumer.groupMetadata());
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
    throw e;
}

Kafka 事务保证"要么全写入、要么全不写入",消费者需要设置 isolation.level=read_committed 才能只看到已提交消息。

六、Percolator 与分布式事务的工程新范式

Google Percolator 是工业界最具影响力的分布式事务实现之一,Bigtable、TiDB、CockroachDB 的单机事务机制都源于此。

6.1 核心思想:行级锁 + MVCC + 时间戳

Percolator 将事务分解为对单行(Row)的读写操作,利用 Bigtable 的单行原子性:

  1. 写入阶段:事务写入"锁列"(Lock Column)和"写入列"(Write Column),写入列写入新数据的主键+时间戳指针,锁列保存事务的 Primary Key 位置
  2. 提交阶段:从 Primary Key 开始,将锁列转为"提交时间戳"列(Write Column),释放锁
  3. 清理阶段:如果事务_writer_发现锁过期,可"推翻"该事务并清理
// 简化的 Percolator 行事务模型
pub struct PercolatorTxn {
    timestamp: Timestamp,  // Oracle 分配的单调递增时间戳
}

impl PercolatorTxn {
    // 写入数据 + 加锁
    pub fn put(&mut self, row: RowKey, data: Vec<u8>) -> Result<()> {
        let mut batch = WriteBatch::new();
        // 写 Write Column: data -> {commit_ts, data}
        batch.put_cf(WRITE_CF, &row, &self.encode_write(self.timestamp, &data));
        // 写 Lock Column: 保存 primary key
        batch.put_cf(LOCK_CF, &row, &self.encode_lock(&self.primary_key));
        self.db.write(batch)?;
        Ok(())
    }

    // 提交事务
    pub fn commit(self) -> Result<()> {
        // 1. 提交 Primary Key
        let commit_ts = self.oracle.timestamp_one()?;
        let primary_lock = self.read_lock(&self.primary_key)?;
        self.commit_single_key(self.primary_key.clone(), primary_lock, commit_ts)?;

        // 2. 并行提交所有 Secondary Keys
        for row in &self.secondaries {
            let lock = self.read_lock(row)?;
            self.commit_single_key(row.clone(), lock, commit_ts)?;
        }
        Ok(())
    }

    fn commit_single_key(&self, row: RowKey, lock: Lock, commit_ts: Timestamp) -> Result<()> {
        // CAS 操作:将锁转为已提交状态
        let mut batch = WriteBatch::new();
        batch.put_cf(WRITE_CF, &row, &self.encode_write(commit_ts, &lock.data));
        batch.delete_cf(LOCK_CF, &row);
        self.db.write(batch)
    }
}

6.2 TiDB 的工程优化

TiDB 在 Percolator 基础上做了多项生产优化:

  • 异步提交(Async Commit):只等待 Primary Key 提交即返回,Secondary Key 异步提交,大幅降低延迟
  • 1PC(一阶段提交):若所有 KV 在同一 Region,直接提交无需走两阶段
  • 悲观事务模式:读取时即加锁,避免乐观事务的冲突回滚
  • Memory Lock:内存中维护锁检测,减少 Raft 提交的 IO 次数
-- TiDB 中开启异步提交
SET GLOBAL tidb_enable_async_commit = ON;
SET GLOBAL tidb_enable_1pc = ON;

-- 批量写入优化
BEGIN;
INSERT INTO orders ... ;
INSERT INTO order_items ... ;
COMMIT;
-- 若满足 1PC 条件,提交延迟从 ~50ms 降至 ~3ms

七、工程落地的四个关键决策

7.1 场景匹配矩阵

一致性要求 延迟敏感 涉及外部系统 推荐方案
高 是 否 本地事务 + 幂等重试
高 是 否 TCC(资源预留)
高 否 否 Seata AT / Percolator
中 否 否 Saga 编排
低 否 是 事务消息 + 对账
低 否 是 最大努力通知 + 人工兜底

7.2 幂等与对账:工程的生命线

无论选择哪种方案,都必须实现:

  • 业务幂等:每个操作可重复执行且结果相同(如 INSERT ... ON DUPLICATE KEY UPDATE 或唯一键约束)
  • 异步对账:定时比对上下游数据差异,如订单系统的"支付状态"与"订单状态"校准
  • 死信队列:无法自动补偿的消息进入告警,人工介入
// 幂等性实现示例
func (s *OrderService) CreateOrder(ctx context.Context, req *CreateOrderRequest) (*Order, error) {
    // 利用唯一索引保证幂等
    order := &Order{
        ID:        req.OrderID, // 客户端生成的 UUID
        UserID:    req.UserID,
        Amount:    req.Amount,
        Status:    "CREATED",
        CreatedAt: time.Now(),
    }

    _, err := s.db.NamedExecContext(ctx, `
        INSERT INTO orders (id, user_id, amount, status, created_at)
        VALUES (:id, :user_id, :amount, :status, :created_at)
        ON DUPLICATE KEY UPDATE updated_at = NOW()
    `, order)

    if err != nil {
        // 唯一键冲突 = 已存在,查询并返回
        if isDuplicateKeyError(err) {
            return s.FindOrder(ctx, req.OrderID)
        }
        return nil, err
    }
    return order, nil
}

7.3 性能:2PC 之外的工程选择

TCC 的 Try 阶段可能导致业务资源长时间占用(如库存冻结后用户取消),在生产中可优化为:

  • 准实时补偿:Try 成功后立即执行 Confirm,仅在 Cancel 时触发补偿(类似 AT 模式)
  • 库存冻结替代扣减:Try 阶段不做任何数据库变更,Confirm 时真正扣减,冲突时通过重试+超时放弃
  • Saga 的幂等性设计:使用 Saga 数据携带(Saga Data)模式,每个步骤读取数据而非依赖上游服务状态

7.4 监控与告警:不可忽视的工程闭环

分布式事务系统必须有完善的监控:

  • 事务成功率/回滚率:Tx Success Rate < 99.9% 需立即告警
  • 平均提交延迟:P99 > 500ms 需排查锁竞争/网络问题
  • 补偿失败率:补偿失败的告警必须配 oncall,避免数据不一致
  • 长事务检测:运行超过 10s 的事务应被 kill 或告警
  • 对账差异数:日均对账差异 ≠ 0 说明系统有 bug

八、总结:架构演进与工程权衡

方案 侵入性 一致性 性能 适用场景
2PC/XA 低 强 低 同库事务、MySQL XA
TCC 高 强 高 库存扣减、积分操作
Saga 中 最终 高 长流程、跨系统
事务消息 低 最终 高 异步通知、数据同步
Percolator 低 强 中 TiDB/CRDB 内置

核心原则:不要为了"强一致性"而牺牲可用性。在大多数互联网业务场景下,通过"乐观锁 + 幂等设计 + 异步对账 + 人工兜底"的工程组合,可以实现"足够好"的一致性,同时保持系统的高可用与高性能。

落地时建议优先选择对业务侵入最小的事务消息方案,引入 TCC/Saga 仅当业务确实需要强一致性,且团队具备足够的分布式系统调试能力。最后,对账系统是分布式事务的"最后防线",是任何方案都必须配套建设的生产基础设施。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部