一、分布式事务导论

在微服务架构中,一个业务操作往往需要跨越多个服务和数据库。例如电商下单场景需要同时操作订单服务、库存服务、积分服务和支付服务,这就要求我们在多个独立的数据库之间保证数据一致性。分布式事务正是为了解决这一挑战而诞生的话题。

与单机数据库的ACID事务不同,分布式事务面临网络分区、节点故障、时钟漂移等额外挑战。本文将从CAP定理和BASE理论出发,深入讲解两阶段提交(2PC)、TCC、Saga和本地消息表五种分布式事务方案,并提供完整的Go语言实现。

二、理论基础:CAP与BASE

2.1 CAP定理

CAP定理指出,在一个分布式系统中,一致性(Consistency)、可用性(Availability)、分区容错性(Partition tolerance)三者不可兼得,最多只能同时满足其中两项:

  • 一致性(C):所有节点在同一时刻看到的数据完全相同
  • 可用性(A):每个请求都能在有限时间内收到非错误的响应
  • 分区容错性(P):系统在网络分区的情况下仍能继续运行

由于网络分区在分布式系统中几乎必然发生,实际的选择通常是在CP(强一致性)和AP(高可用)之间权衡。

2.2 BASE理论

BASE是对CAP中AP方案的延伸,其核心思想是通过牺牲强一致性来获得高可用性:

  • Basically Available(基本可用):系统在故障时允许损失部分可用性
  • Soft State(软状态):允许系统中的数据存在中间状态
  • Eventually Consistent(最终一致性):经过一段时间后,所有副本的数据达到一致

三、两阶段提交(2PC)

3.1 工作原理

2PC引入一个协调者(Coordinator)和多个参与者(Participant),分为两个阶段:

  1. 准备阶段(Voting Phase):协调者向所有参与者发送Prepare请求,参与者执行事务但不提交,锁定资源并记录Undo/Redo日志
  2. 提交阶段(Commit Phase):如果所有参与者都返回Yes,协调者发送Commit指令;任一参与者返回No则发送Rollback指令

3.2 Go语言实现

package twopc

import (
 "context"
 "errors"
 "fmt"
 "sync"
 "time"
)

// TransactionPhase 事务阶段
type TransactionPhase int

const (
 PhasePrepare TransactionPhase = iota
 PhaseCommit
 PhaseRollback
)

// Participant 参与者接口
type Participant interface {
 ID() string
 Prepare(ctx context.Context, txID string) error
 Commit(ctx context.Context, txID string) error
 Rollback(ctx context.Context, txID string) error
}

// Coordinator 两阶段提交协调者
type Coordinator struct {
 participants []Participant
 timeout      time.Duration
}

func NewCoordinator(timeout time.Duration) *Coordinator {
 return &Coordinator{
  participants: make([]Participant, 0),
  timeout:      timeout,
 }
}

func (c *Coordinator) Register(p Participant) {
 c.participants = append(c.participants, p)
}

// Execute 执行分布式事务
func (c *Coordinator) Execute(ctx context.Context, txID string) error {
 ctx, cancel := context.WithTimeout(ctx, c.timeout)
 defer cancel()

 // 第一阶段:准备阶段
 prepareErrors := make(chan error, len(c.participants))
 var wg sync.WaitGroup

 for _, p := range c.participants {
  wg.Add(1)
  go func(part Participant) {
   defer wg.Done()
   if err := part.Prepare(ctx, txID); err != nil {
    prepareErrors <- fmt.Errorf("participant %s prepare failed: %w", part.ID(), err)
   }
  }(p)
 }

 wg.Wait()
 close(prepareErrors)

 // 检查准备结果
 for err := range prepareErrors {
  c.rollbackAll(ctx, txID)
  return err
 }

 // 第二阶段:提交阶段
 for _, p := range c.participants {
  wg.Add(1)
  go func(part Participant) {
   defer wg.Done()
   if err := part.Commit(ctx, txID); err != nil {
    prepareErrors <- fmt.Errorf("participant %s commit failed: %w", part.ID(), err)
   }
  }(p)
 }

 wg.Wait()
 return nil
}

func (c *Coordinator) rollbackAll(ctx context.Context, txID string) {
 var wg sync.WaitGroup
 for _, p := range c.participants {
  wg.Add(1)
  go func(part Participant) {
   defer wg.Done()
   part.Rollback(ctx, txID)
  }(p)
 }
 wg.Wait()
}

// AccountParticipant 账务参与者示例
type AccountParticipant struct {
 id      string
 balance float64
 locked  float64
}

func (a *AccountParticipant) ID() string { return a.id }

func (a *AccountParticipant) Prepare(ctx context.Context, txID string) error {
 fmt.Printf("[%s] 准备阶段:锁定账户 %s\n", txID, a.id)
 return nil
}

func (a *AccountParticipant) Commit(ctx context.Context, txID string) error {
 fmt.Printf("[%s] 提交阶段:确认账户 %s 变更\n", txID, a.id)
 return nil
}

func (a *AccountParticipant) Rollback(ctx context.Context, txID string) error {
 fmt.Printf("[%s] 回滚:释放账户 %s 锁定\n", txID, a.id)
 return nil
}

3.3 2PC的缺陷

  • 同步阻塞:准备阶段参与者需要锁定资源,长时间等待会导致性能下降
  • 单点故障:协调者宕机会导致参与者一直持有锁
  • 数据不一致:第二阶段部分参与者收不到Commit指令时产生不一致

四、TCC(Try-Confirm-Cancel)

4.1 工作原理

TCC是一种补偿型事务,将每个参与者的操作拆分为三个阶段:

  • Try:预留资源,做业务可执行的检查(如余额是否充足)
  • Confirm:确认执行,真正执行业务(如扣减余额),此阶段不应失败
  • Cancel:取消执行,释放预留的资源(如退还余额)

4.2 Go语言实现

package tcc

import (
 "context"
 "errors"
 "fmt"
 "time"
)

// TCCAction TCC参与者接口
type TCCAction interface {
 Try(ctx context.Context, bizParams map[string]interface{}) error
 Confirm(ctx context.Context, bizParams map[string]interface{}) error
 Cancel(ctx context.Context, bizParams map[string]interface{}) error
}

// TCCTransaction TCC事务管理器
type TCCTransaction struct {
 actions []TCCAction
 timeout time.Duration
}

func NewTCCTransaction(timeout time.Duration) *TCCTransaction {
 return &TCCTransaction{
  actions: make([]TCCAction, 0),
  timeout: timeout,
 }
}

func (t *TCCTransaction) AddAction(action TCCAction) {
 t.actions = append(t.actions, action)
}

func (t *TCCTransaction) Execute(ctx context.Context, params []map[string]interface{}) error {
 ctx, cancel := context.WithTimeout(ctx, t.timeout)
 defer cancel()

 completed := make([]int, 0)

 // 第一阶段:Try
 for i, action := range t.actions {
  bizParams := params[i]
  if err := action.Try(ctx, bizParams); err != nil {
   // Try失败,Cancel已成功的操作
   for j := len(completed) - 1; j >= 0; j-- {
    t.actions[completed[j]].Cancel(ctx, params[completed[j]])
   }
   return fmt.Errorf("action %d try failed: %w", i, err)
  }
  completed = append(completed, i)
 }

 // 第二阶段:Confirm(理论上不应失败)
 for i, action := range t.actions {
  if err := action.Confirm(ctx, params[i]); err != nil {
   // Confirm失败需重试,不应放弃
   return fmt.Errorf("action %[1]d confirm failed CRITICAL: %w", i, err)
  }
 }

 return nil
}

// AccountService 账务服务TCC实现示例
type AccountService struct {
 db     map[string]float64 // 模拟用户余额
 locked map[string]float64 // 冻结金额
}

func NewAccountService() *AccountService {
 return &AccountService{
  db:     map[string]float64{"alice": 1000, "bob": 500},
  locked: make(map[string]float64),
 }
}

// Try 冻结指定金额
func (s *AccountService) Try(ctx context.Context, params map[string]interface{}) error {
 userID := params["user_id"].(string)
 amount := params["amount"].(float64)

 if s.db[userID] < amount xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed xss=removed>

4.3 TCC注意事项

  • 空回滚防护:Try未收到就收到Cancel时,需要识别并记录,Cancel时直接返回成功
  • 幂等性:Confirm和Cancel可能因网络重试被多次调用,需保证幂等
  • 悬挂问题:Cancel先于Try执行,Try需检查是否存在已Cancel的记录

五、Saga模式

5.1 工作原理

Saga将一个长事务拆分为一系列本地事务,每个本地事务都有对应的补偿操作。如果某个步骤失败,则按逆序执行已完成步骤的补偿操作。

Saga有两种协调方式:

  • 编排式(Choreography):各服务通过事件驱动,每个服务完成自己的本地事务后发布事件,触发下一个服务
  • 协调式(Orchestration):由中央协调者按顺序调用各服务,失败时触发补偿

5.2 协调式Saga Go实现

package saga

import (
 "context"
 "errors"
 "fmt"
 "time"
)

// SagaStep Saga步骤
type SagaStep struct {
 Name       string
 Action     func(ctx context.Context) error
 Compensate func(ctx context.Context) error
}

// Saga Saga事务定义
type Saga struct {
 Name  string
 Steps []SagaStep
}

func NewSaga(name string) *Saga {
 return &Saga{
  Name:  name,
  Steps: make([]SagaStep, 0),
 }
}

func (s *Saga) AddStep(name string, action func(ctx context.Context) error, compensate func(ctx context.Context) error) {
 s.Steps = append(s.Steps, SagaStep{
  Name:       name,
  Action:     action,
  Compensate: compensate,
 })
}

// Execute 执行Saga事务
func (s *Saga) Execute(ctx context.Context) error {
 completedSteps := make([]int, 0)

 for i, step := range s.Steps {
  stepCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
  fmt.Printf("[Saga] 执行步骤 %d: %s\n", i, step.Name)

  if err := step.Action(stepCtx); err != nil {
   cancel()
   fmt.Printf("[Saga] 步骤 %d 失败: %v,开始补偿\n", i, err)

   // 逆序补偿
   for j := len(completedSteps) - 1; j >= 0; j-- {
    compStep := s.Steps[completedSteps[j]]
    compCtx, compCancel := context.WithTimeout(ctx, 30*time.Second)
    if compErr := compStep.Compensate(compCtx); compErr != nil {
     fmt.Printf("[Saga] 补偿步骤 '%s' 失败: %v(需人工介入)\n", compStep.Name, compErr)
    }
    compCancel()
   }
   return fmt.Errorf("saga '%s' failed at step '%s': %w", s.Name, step.Name, err)
  }

  cancel()
  completedSteps = append(completedSteps, i)
  fmt.Printf("[Saga] 步骤 %d 完成\n", i)
 }

 fmt.Printf("[Saga] 事务 '%s' 全部完成\n", s.Name)
 return nil
}

// 电商下单Saga示例
func OrderSagaExample() {
 ctx := context.Background()
 saga := NewSaga("创建订单")

 var orderID string
 var frozenStock []int
 var deductedPoints int

 // 步骤1:创建订单
 saga.AddStep(
  "创建订单",
  func(ctx context.Context) error {
   orderID = fmt.Sprintf("ORDER_%d", time.Now().Unix())
   fmt.Printf("创建订单: %s\n", orderID)
   return nil
  },
  func(ctx context.Context) error {
   fmt.Printf("取消订单: %s\n", orderID)
   orderID = ""
   return nil
  },
 )

 // 步骤2:扣减库存
 saga.AddStep(
  "扣减库存",
  func(ctx context.Context) error {
   frozenStock = []int{101, 102}
   fmt.Printf("扣减商品库存: %v\n", frozenStock)
   return nil
  },
  func(ctx context.Context) error {
   fmt.Printf("恢复库存: %v\n", frozenStock)
   frozenStock = nil
   return nil
  },
 )

 // 步骤3:扣减积分
 saga.AddStep(
  "扣减积分",
  func(ctx context.Context) error {
   deductedPoints = 100
   fmt.Printf("扣减积分: %d\n", deductedPoints)
   return nil
  },
  func(ctx context.Context) error {
   fmt.Printf("退还积分: %d\n", deductedPoints)
   deductedPoints = 0
   return nil
  },
 )

 // 步骤4:支付
 saga.AddStep(
  "支付",
  func(ctx context.Context) error {
   fmt.Printf("执行支付,订单: %s\n", orderID)
   return nil
  },
  func(ctx context.Context) error {
   fmt.Printf("退款完成\n")
   return nil
  },
 )

 if err := saga.Execute(ctx); err != nil {
  fmt.Printf("下单失败: %v\n", err)
 } else {
  fmt.Println("下单成功!")
 }
}

// 模拟支付失败的Saga
func OrderSagaWithFailureExample() {
 ctx := context.Background()
 saga := NewSaga("创建订单(模拟支付失败)")

 saga.AddStep("创建订单",
  func(ctx context.Context) error { fmt.Println("订单创建成功"); return nil },
  func(ctx context.Context) error { fmt.Println("订单已取消"); return nil },
 )
 saga.AddStep("扣减库存",
  func(ctx context.Context) error { fmt.Println("库存已扣减"); return nil },
  func(ctx context.Context) error { fmt.Println("库存已恢复"); return nil },
 )
 saga.AddStep("积分抵扣",
  func(ctx context.Context) error { fmt.Println("积分已抵扣"); return nil },
  func(ctx context.Context) error { fmt.Println("积分已退还"); return nil },
 )
 saga.AddStep("支付",
  func(ctx context.Context) error { return errors.New("余额不足,支付失败") },
  func(ctx context.Context) error { return nil },
 )

 if err := saga.Execute(ctx); err != nil {
  fmt.Printf("流程终止: %v\n", err)
 }
}

5.3 编排式Saga vs 协调式Saga

对比维度编排式(Choreography)协调式(Orchestration)
复杂度简单场景适用,步骤增多后复杂度高通过集中控制逻辑清晰
耦合度服务间通过事件相互感知服务只与协调者交互
单点风险无中心节点协调者是单点,需要考虑高可用
可观测性较难追踪整体流程协调者集中记录执行状态
适用场景步骤少、逻辑简单的短Saga步骤多、逻辑复杂的长Saga

六、本地消息表

6.1 工作原理

本地消息表的核心思想是将分布式事务拆分为本地事务,通过保证本地事务和消息投递的一致性来实现最终一致性:

  1. 业务执行和消息插入在同一个本地事务中完成
  2. 定时任务扫描未发送的消息,投递到消息队列
  3. 下游消费者处理消息,成功后回调确认
  4. 失败时通过重试机制保证最终送达

6.2 Go实现核心逻辑

package localmsg

import (
 "context"
 "database/sql"
 "fmt"
 "time"
)

// Message 本地消息
type Message struct {
 ID         string
 Topic      string
 Payload    string
 Status     int       // 0:待发送 1:已发送 2:已完成
 RetryCount int
 NextRetry  time.Time
 CreatedAt  time.Time
}

// MessageStore 消息存储接口
type MessageStore interface {
 Save(ctx context.Context, msg *Message, execTx func(tx *sql.Tx) error) error
 ScanPending(ctx context.Context, limit int) ([]*Message, error)
 MarkSent(ctx context.Context, id string) error
 MarkDone(ctx context.Context, id string) error
 IncrementRetry(ctx context.Context, id string) error
}

// DBMessageStore 基于数据库的实现
type DBMessageStore struct {
 db *sql.DB
}

func NewDBMessageStore(db *sql.DB) *DBMessageStore {
 return &DBMessageStore{db: db}
}

// Save 在同一个事务中保存消息和业务数据
func (s *DBMessageStore) Save(ctx context.Context, msg *Message, execTx func(tx *sql.Tx) error) error {
 tx, err := s.db.BeginTx(ctx, nil)
 if err != nil {
  return err
 }
 defer tx.Rollback()

 // 执行业务逻辑
 if err := execTx(tx); err != nil {
  return err
 }

 // 插入消息记录
 _, err = tx.ExecContext(ctx,
  "INSERT INTO local_message (id, topic, payload, status, created_at) VALUES (?, ?, ?, 0, ?)",
  msg.ID, msg.Topic, msg.Payload, time.Now(),
 )
 if err != nil {
  return fmt.Errorf("insert message failed: %w", err)
 }

 return tx.Commit()
}

// ScanPending 扫描待发送消息
func (s *DBMessageStore) ScanPending(ctx context.Context, limit int) ([]*Message, error) {
 rows, err := s.db.QueryContext(ctx,
  "SELECT id, topic, payload, status, retry_count, next_retry FROM local_message WHERE status = 0 AND next_retry <= ? LIMIT ?",
  time.Now(), limit,
 )
 if err != nil {
  return nil, err
 }
 defer rows.Close()

 msgs := make([]*Message, 0)
 for rows.Next() {
  m := &Message{}
  err := rows.Scan(&m.ID, &m.Topic, &m.Payload, &m.Status, &m.RetryCount, &m.NextRetry)
  if err != nil {
   return nil, err
  }
  msgs = append(msgs, m)
 }
 return msgs, nil
}

func (s *DBMessageStore) MarkSent(ctx context.Context, id string) error {
 _, err := s.db.ExecContext(ctx, "UPDATE local_message SET status = 1 WHERE id = ?", id)
 return err
}

func (s *DBMessageStore) MarkDone(ctx context.Context, id string) error {
 _, err := s.db.ExecContext(ctx, "UPDATE local_message SET status = 2 WHERE id = ?", id)
 return err
}

func (s *DBMessageStore) IncrementRetry(ctx context.Context, id string) error {
 _, err := s.db.ExecContext(ctx,
  "UPDATE local_message SET retry_count = retry_count + 1, next_retry = ? WHERE id = ?",
  time.Now().Add(60*time.Second), id,
 )
 return err
}

// MessageWorker 消息投递工作器
type MessageWorker struct {
 store  MessageStore
 mq     MQProducer
 stopCh chan struct{}
}

// MQProducer 消息队列生产者接口
type MQProducer interface {
 Publish(ctx context.Context, topic string, payload []byte) error
}

func NewMessageWorker(store MessageStore, mq MQProducer) *MessageWorker {
 return &MessageWorker{
  store:  store,
  mq:     mq,
  stopCh: make(chan struct{}),
 }
}

func (w *MessageWorker) Start() {
 ticker := time.NewTicker(5 * time.Second)
 defer ticker.Stop()

 for {
  select {
  case <-w.stopCh:
   return
  case <-ticker.C:
   w.process()
  }
 }
}

func (w *MessageWorker) Stop() {
 close(w.stopCh)
}

func (w *MessageWorker) process() {
 ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
 defer cancel()

 msgs, err := w.store.ScanPending(ctx, 100)
 if err != nil {
  fmt.Printf("扫描消息失败: %v\n", err)
  return
 }

 for _, msg := range msgs {
  if err := w.mq.Publish(ctx, msg.Topic, []byte(msg.Payload)); err != nil {
   fmt.Printf("投递消息 %s 失败: %v\n", msg.ID, err)
   w.store.IncrementRetry(ctx, msg.ID)
   continue
  }
  w.store.MarkSent(ctx, msg.ID)
  fmt.Printf("消息 %s 已投递\n", msg.ID)
 }
}

// 建表SQL
const CreateTableSQL = `CREATE TABLE IF NOT EXISTS local_message (
 id VARCHAR(64) PRIMARY KEY,
 topic VARCHAR(128) NOT NULL,
 payload TEXT NOT NULL,
 status TINYINT NOT NULL DEFAULT 0,
 retry_count INT NOT NULL DEFAULT 0,
 next_retry DATETIME NOT NULL,
 created_at DATETIME NOT NULL,
 INDEX idx_status_next_retry (status, next_retry)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;`

七、方案对比与选型建议

方案一致性性能实现复杂度适用场景
2PC强一致低(同步阻塞)中等数据库层分布式事务,如XA
TCC最终一致高(需三阶段实现)高一致性要求,如金融交易
Saga最终一致中等长流程业务,如电商下单
本地消息表最终一致跨服务数据同步,如积分发放
最大努力通知弱一致最高最低对一致性要求极低的场景

八、最佳实践建议

  1. 优先考虑最终一致性:大多数业务场景无需强一致,BASE理论提供的最终一致性已经足够,且能带来更好的性能和可用性
  2. TCC用于核心链路:对一致性要求较高的核心链路(如支付、扣款),推荐使用TCC模式
  3. Saga应对长流程:涉及3个以上服务的复杂业务流程,使用Saga模式配合补偿机制
  4. 本地消息表作为兜底:作为最终手段保证操作消息不丢失,可与上述任何方案配合使用
  5. 始终保证幂等性:网络调用不可靠,所有分布式事务操作必须设计为幂等的
  6. 记录并监控补偿失败:补偿操作也可能失败,必须记录日志并及时告警,必要时人工介入
  7. 超时要短,重试要谨慎:参与者操作的超时时间应远短于全局超时,重试需有上限和退避策略
  8. 考虑引入成熟框架:如Seata(支持AT/TCC/SAGA/XA模式)、DTM(Go语言分布式事务管理器)可大幅降低实现成本

九、总结

分布式事务没有银弹,每种方案都有其适用场景和代价。在实际工程中,我们应该:

  • 首先判断是否真的需要分布式事务——能通过合理的服务设计规避最好
  • 需要时优先选择性能更好的弱一致性方案
  • 核心链路使用TCC或Saga + 本地消息表的双重保障
  • 建立完善的监控和补偿机制,确保问题可追溯、可修复

CAP定理告诉我们完美不可得,而工程师的艺术在于做出最合理的取舍。希望本文提供的理论分析和Go代码示例,能为你的分布式系统设计提供有价值的参考。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部