分布式事务工程实践:从 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 将事务分为两个阶段:
- Prepare 阶段:协调者询问所有参与者是否可以提交,参与者执行事务但不提交,写入 undo/redo 日志,锁定资源,回复 Yes/No。
- 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 的单行原子性:
- 写入阶段:事务写入"锁列"(Lock Column)和"写入列"(Write Column),写入列写入新数据的主键+时间戳指针,锁列保存事务的 Primary Key 位置
- 提交阶段:从 Primary Key 开始,将锁列转为"提交时间戳"列(Write Column),释放锁
- 清理阶段:如果事务_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 仅当业务确实需要强一致性,且团队具备足够的分布式系统调试能力。最后,对账系统是分布式事务的"最后防线",是任何方案都必须配套建设的生产基础设施。

发表评论 取消回复