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

在单体架构中,一个业务操作只需在一个数据库中完成 ACID 事务即可保证一致性。但在微服务架构下,一个业务操作往往跨越多个服务和多个数据库。以最经典的"电商下单"为例,它需要依次操作订单服务、库存服务、积分服务和支付服务——这些服务各自拥有独立的数据库实例。

分布式事务解决的核心问题是:在多个独立的数据库节点之间保证数据的一致性,即使面临网络分区、节点宕机、请求超时等故障场景。

本文将从 CAP 定理和 BASE 理论出发,深入讲解四种主流的分布式事务方案(2PC、TCC、Saga、本地消息表),提供完整的 Go 语言实现代码,并分析各方案的适用场景与选型建议。

二、理论基础:CAP 与 BASE

2.1 CAP 定理

CAP 定理(又称 Brewer 定理)指出,在一个分布式系统中,以下三个特性最多只能同时满足两个:

  • 一致性 (Consistency):所有节点在同一时刻看到的数据完全相同。读操作总能读到最新写入的结果。
  • 可用性 (Availability):每个请求都能在有限时间内收到非错误的响应(但不保证读到最新数据)。
  • 分区容错性 (Partition tolerance):系统在网络分区(节点间通信中断)的情况下仍能继续运行。

由于网络分区在分布式系统中几乎必然发生(网线被拔、交换机故障、网络拥塞),实际上我们总是在 CP 和 AP 之间做选择:

  • CP 架构:保证一致性,但在分区期间牺牲可用性(如 ZooKeeper、etcd)
  • AP 架构:保证可用性,但允许数据短期内不一致(如 Cassandra、DynamoDB)

2.2 BASE 理论

BASE 是对 CAP 中 AP 方案的延伸,其核心理念是通过牺牲强一致性来获得高可用性,允许数据在一段时间内处于不一致状态,但最终达到一致:

  • Basically Available(基本可用):系统在故障时允许损失部分可用性(如响应时间变长、降级服务)
  • Soft State(软状态):允许系统中的数据存在中间状态(如"处理中"),且该中间状态不影响系统整体可用性
  • Eventually Consistent(最终一致性):经过一段时间后,所有副本的数据会自动达到一致状态

BASE 理论是大多数分布式事务方案(TCC、Saga、本地消息表)的哲学基础——我们不再追求每时每刻的强一致,而是保证最终一致。

三、两阶段提交(2PC)

3.1 工作原理

两阶段提交引入两种角色:一个协调者 (Coordinator)和多个参与者 (Participant)。整个事务分为两个阶段执行:

  1. 准备阶段 (Voting Phase):协调者向所有参与者发送 Prepare 请求。参与者执行事务操作但不提交(锁定资源、写 Undo/Redo 日志),然后返回 Yes 或 No。
  2. 提交阶段 (Commit Phase):如果所有参与者都返回 Yes,协调者发送 Commit 指令,参与者正式提交;如果任一参与者返回 No 或超时,协调者发送 Rollback 指令,参与者回滚。

3.2 Go 语言实现

package twopc

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

// Participant 分布式事务参与者接口
type Participant interface {
 ID() string
 // Prepare 执行准备操作,返回 nil 表示同意,返回 error 表示拒绝
 Prepare(ctx context.Context, txID string) error
 // Commit 正式提交事务
 Commit(ctx context.Context, txID string) error
 // Rollback 回滚事务
 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 {
 // ========== 第一阶段:准备阶段 ==========
 fmt.Printf("[%s] === 开始第一阶段:准备阶段 ===\n", txID)
 
 prepareResult := 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()
   pCtx, cancel := context.WithTimeout(ctx, c.timeout)
   defer cancel()
   if err := part.Prepare(pCtx, txID); err != nil {
    prepareResult <- fmt.Errorf("参与者 %s 准备阶段失败: %w", part.ID(), err)
   }
  }(p)
 }

 wg.Wait()
 close(prepareResult)

 // 检查是否有参与者拒绝
 for err := range prepareResult {
  fmt.Printf("[%s] 准备阶段失败: %s,开始回滚\n", txID, err)
  c.rollbackAll(ctx, txID)
  return fmt.Errorf("事务中止: %w", err)
 }

 fmt.Printf("[%s] 所有参与者准备成功\n", txID)

 // ========== 第二阶段:提交阶段 ==========
 fmt.Printf("[%s] === 开始第二阶段:提交阶段 ===\n", txID)

 commitResult := make(chan error, len(c.participants))
 for _, p := range c.participants {
  wg.Add(1)
  go func(part Participant) {
   defer wg.Done()
   cCtx, cancel := context.WithTimeout(ctx, c.timeout)
   defer cancel()
   if err := part.Commit(cCtx, txID); err != nil {
    commitResult <- fmt.Errorf("参与者 %s 提交阶段失败: %w", part.ID(), err)
   }
  }(p)
 }

 wg.Wait()
 close(commitResult)

 // 收集提交中的失败(理论上不应发生,需人工介入或重试)
 var commitErrors []error
 for err := range commitResult {
  commitErrors = append(commitErrors, err)
 }

 if len(commitErrors) > 0 {
  fmt.Printf("[%s] 警告:部分参与者提交失败(需人工介入或自动重试):\n", txID)
  for _, err := range commitErrors {
   fmt.Printf("  - %s\n", err)
  }
  return fmt.Errorf("事务部分提交失败,需补偿处理")
 }

 fmt.Printf("[%s] 事务提交成功\n", txID)
 return nil
}

// rollbackAll 回滚所有参与者
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()
   rCtx, cancel := context.WithTimeout(ctx, c.timeout)
   defer cancel()
   if err := part.Rollback(rCtx, txID); err != nil {
    fmt.Printf("[%s] 参与者 %s 回滚失败: %v\n", txID, part.ID(), err)
   }
  }(p)
 }
 wg.Wait()
 fmt.Printf("[%s] 全量回滚完成\n", txID)
}

// PaymentService 支付服务参与者示例
type PaymentService struct {
 name string
}

func (s *PaymentService) ID() string { return s.name }

func (s *PaymentService) Prepare(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 冻结用户账户金额\n", txID, s.name)
 // 实际场景:UPDATE account SET frozen = frozen + ? WHERE user_id = ?
 return nil
}

func (s *PaymentService) Commit(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 确认扣减账户余额\n", txID, s.name)
 // 实际场景:UPDATE account SET balance = balance - ?, frozen = frozen - ? WHERE user_id = ?
 return nil
}

func (s *PaymentService) Rollback(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 解冻用户账户金额\n", txID, s.name)
 // 实际场景:UPDATE account SET frozen = frozen - ? WHERE user_id = ?
 return nil
}

// InventoryService 库存服务参与者示例
type InventoryService struct {
 name string
}

func (s *InventoryService) ID() string { return s.name }

func (s *InventoryService) Prepare(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 锁定库存\n", txID, s.name)
 // 实际场景:SELECT ... FOR UPDATE 锁定库存记录
 return nil
}

func (s *InventoryService) Commit(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 确认扣减库存\n", txID, s.name)
 // 实际场景:UPDATE inventory SET stock = stock - ? WHERE sku = ?
 return nil
}

func (s *InventoryService) Rollback(ctx context.Context, txID string) error {
 fmt.Printf("[%s] %s: 释放库存锁定\n", txID, s.name)
 return nil
}

3.3 2PC 的优缺点

优点:原理简单,有成熟的数据库原生支持(MySQL XA、PostgreSQL PREPARE TRANSACTION),保证强一致性。

缺陷:

  • 同步阻塞:准备阶段所有参与者需要锁定资源直到第二阶段完成,长时间持锁严重影响并发性能
  • 协调者单点故障:如果协调者在准备阶段完成后、提交阶段前宕机,参与者将永远持有锁且无法得知最终决定
  • 脑裂风险:网络分区导致部分参与者未收到 Commit 指令时,已提交和未提交节点之间出现数据不一致
  • 性能瓶颈:所有参与者必须串行等待最慢的那个,不适合高并发短时间的事务

四、TCC(Try-Confirm-Cancel)

4.1 工作原理

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

  • Try(尝试执行):做业务层面的检查与资源预留。例如下单时不是直接扣库存,而是将库存状态设为"冻结";支付时不是直接扣款,而是将金额从余额移到"冻结金额"字段。
  • Confirm(确认执行):真正执行业务操作。将冻结转为实际扣减。Confirm 被设计为几乎不应失败—— Try 已经完成了所有校验,Confirm 只是执行确定性的数据变更。
  • Cancel(取消执行):当 Try 失败时,对已成功的 Try 操作执行补偿(反向操作),释放预留的资源。

TCC 的核心思路是数据库层面不做全局锁,而是通过业务层面的"预留+确认"来实现事务。每个阶段都是独立的本地事务,因此避免了 2PC 的长时间锁等待问题。

4.2 Go 语言实现

package tcc

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

// TccAction TCC 参与者接口
type TccAction interface {
 // Try 预留资源,执行业务检查
 Try(ctx context.Context, params map[string]interface{}) error
 // Confirm 确认执行,真正扣减资源(此阶段不应失败)
 Confirm(ctx context.Context, params map[string]interface{}) error
 // Cancel 取消执行,释放预留的资源
 Cancel(ctx context.Context, params map[string]interface{}) error
}

// TccStep TCC 事务中的一个步骤
type TccStep struct {
 Name   string
 Action TccAction
 Params map[string]interface{}
}

// TccOrchestrator TCC 事务协调器
type TccOrchestrator struct {
 timeout time.Duration
}

func NewTccOrchestrator(timeout time.Duration) *TccOrchestrator {
 return &TccOrchestrator{timeout: timeout}
}

// Execute 执行 TCC 事务
func (o *TccOrchestrator) Execute(ctx context.Context, steps []TccStep) error {
 ctx, cancel := context.WithTimeout(ctx, o.timeout)
 defer cancel()

 completedSteps := make([]int, 0, len(steps))

 // ========== 第一阶段:Try ==========
 for i, step := range steps {
  fmt.Printf("[TCC] Try 阶段 - 步骤 %d: %s\n", i, step.Name)
  tCtx, tCancel := context.WithTimeout(ctx, 10*time.Second)

  if err := step.Action.Try(tCtx, step.Params); err != nil {
   tCancel()
   fmt.Printf("[TCC] Try 失败 - %s: %v\n", step.Name, err)
   // 逆序 Cancel 已成功的步骤
   o.cancelCompleted(ctx, steps, completedSteps)
   return fmt.Errorf("步骤 '%s' Try 失败: %w", step.Name, err)
  }

  tCancel()
  completedSteps = append(completedSteps, i)
 }

 // ========== 第二阶段:Confirm ==========
 for i, step := range steps {
  fmt.Printf("[TCC] Confirm 阶段 - 步骤 %d: %s\n", i, step.Name)
  cCtx, cCancel := context.WithTimeout(ctx, 10*time.Second)

  if err := step.Action.Confirm(cCtx, step.Params); err != nil {
   cCancel()
   // Confirm 失败是严重问题——理论上不应发生
   fmt.Printf("[TCC] CRITICAL - Confirm 失败: %s: %v\n", step.Name, err)
   // 需要:记录日志、告警、启动后台重试任务
   return fmt.Errorf("步骤 '%s' Confirm 失败(需人工介入): %w", step.Name, err)
  }

  cCancel()
 }

 fmt.Println("[TCC] 事务提交成功")
 return nil
}

// cancelCompleted 逆序取消已完成的步骤
func (o *TccOrchestrator) cancelCompleted(ctx context.Context, steps []TccStep, completed []int) {
 for i := len(completed) - 1; i >= 0; i-- {
  step := steps[completed[i]]
  cCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
  if err := step.Action.Cancel(cCtx, step.Params); err != nil {
   fmt.Printf("[TCC] Cancel 失败(需重试): %s: %v\n", step.Name, err)
  }
  cancel()
 }
}

4.3 TCC 实现的三大坑

1. 空回滚 (Empty Rollback)

场景:参与者收到 Cancel 请求,但它从未收到过 Try 请求(Try 阶段就超时或丢失了)。此时应允许 Cancel 直接返回成功——你不能对一个从未执行过的操作做补偿。

解决方案:操作记录表记录事务状态。Cancel 时先查询记录,若无 Try 记录则直接返回成功。

2. 幂等性 (Idempotency)

场景:网络抖动导致 Confirm/Cancel 被多次重试调用。如果 Confirm 不是幂等的,就会重复扣款或重复释放资源。

解决方案:使用事务状态机(initial → trying → confirmed / cancelled),每次操作前先检查当前状态,已处于目标状态则直接返回成功。

3. 悬挂 (Suspension)

场景:Cancel 先于 Try 到达参与者。由于网络延迟不可控,可能出现先收到 Cancel(因网络绕路),后收到 Try 的情况。如果不加处理,Try 会在 Cancel 之后执行,导致资源被预留但永远无法 Confirm。

解决方案:Try 执行前检查是否已有 Cancel 记录。如果已 Cancel,则拒绝执行 Try 并返回"事务已取消"。或者先创建事务记录(Status=Trying),Cancel 时检查记录是否存在。

五、Saga 模式

5.1 工作原理

Saga 是一种长事务管理方案,将一个大事务拆分为一系列本地事务,每个本地事务都有对应的补偿操作 (Compensation)。如果某个步骤失败,则按照逆序执行已完成步骤的补偿操作来回滚。

Saga 有两种协调方式:

  • 编排式 Saga (Choreography):各服务通过事件驱动自行协调。每个服务完成本地事务后发布领域事件,下一个服务监听该事件并执行自己的本地事务。没有中央协调者。
  • 协调式 Saga (Orchestration):有一个中央 Saga 协调器(Saga Orchestrator),按顺序调用各服务的本地事务。协调器负责记录执行进度,失败时触发对应的补偿操作。

5.2 协调式 Saga 的 Go 实现

package saga

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

// Step Saga 中的一个步骤
type Step struct {
 Name string
 // Action 正常业务操作
 Action func(ctx context.Context) error
 // Compensate 补偿操作(Action 成功但后续步骤失败时调用)
 Compensate func(ctx context.Context) error
}

// Saga 事务定义
type Saga struct {
 Name     string
 Steps    []Step
 // Compensating 标志当前是否正在执行补偿(用于持久化状态)
 compensating bool
}

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

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

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

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

  if err := step.Action(stepCtx); err != nil {
   cancel()
   fmt.Printf("[Saga:%s] ✗ 步骤 '%s' 失败: %v\n", s.Name, step.Name, err)
   // 执行补偿
   s.compensating = true
   s.compensateAll(ctx, i, completedIndices)
   return fmt.Errorf("Saga '%s' 在步骤 '%s' 中止: %w", s.Name, step.Name, err)
  }

  cancel()
  completedIndices = append(completedIndices, i)
  fmt.Printf("[Saga:%s] ✓ 步骤 '%s' 完成\n", s.Name, step.Name)
 }

 fmt.Printf("[Saga:%s] ✓✓✓ 整个 Saga 执行成功\n", s.Name)
 return nil
}

// compensateAll 逆序执行补偿
func (s *Saga) compensateAll(ctx context.Context, failedIndex int, completed []int) {
 fmt.Printf("[Saga:%s] ░░░ 开始补偿(逆序回滚 %d 个已完成步骤)\n", s.Name, failedIndex)
 for i := failedIndex - 1; i >= 0; i-- {
  step := s.Steps[i]
  compCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
  fmt.Printf("[Saga:%s] ↩ 补偿步骤: %s\n", s.Name, step.Name)
  if err := step.Compensate(compCtx); err != nil {
   // 补偿失败是严重问题——需记录并告警
   fmt.Printf("[Saga:%s] ✗✗✗ 补偿失败 '%s': %v(需人工介入)\n", s.Name, step.Name, err)
  }
  cancel()
 }
 fmt.Printf("[Saga:%s] 补偿完成\n", s.Name)
}

// ===========================
// 电商下单 Saga 示例
// ===========================

// 模拟数据库状态
var (
 dbOrder    bool   // 订单是否创建
 dbStock    int    // 库存状态(扣减/恢复)
 dbPoints   int    // 积分状态
 dbPaid     bool   // 支付状态
 orderID    string
)

func PlaceOrderSaga() {
 ctx := context.Background()
 saga := NewSaga("创建订单-支付流程")

 // 步骤1: 创建订单
 saga.AddStep(
  "创建订单",
  func(ctx context.Context) error {
   dbOrder = true
   orderID = "ORD-20260919-001"
   fmt.Printf("  订单已创建: %s\n", orderID)
   return nil
  },
  func(ctx context.Context) error {
   dbOrder = false
   fmt.Printf("  订单已取消: %s\n", orderID)
   return nil
  },
 )

 // 步骤2: 扣减库存
 saga.AddStep(
  "扣减库存",
  func(ctx context.Context) error {
   dbStock += 2
   fmt.Printf("  库存已扣减 2 件(当前已扣: %d)\n", dbStock)
   return nil
  },
  func(ctx context.Context) error {
   dbStock -= 2
   fmt.Printf("  库存已恢复 2 件\n")
   return nil
  },
 )

 // 步骤3: 扣减积分
 saga.AddStep(
  "扣减积分",
  func(ctx context.Context) error {
   dbPoints += 50
   fmt.Printf("  积分已扣减 50(当前已扣: %d)\n", dbPoints)
   return nil
  },
  func(ctx context.Context) error {
   dbPoints -= 50
   fmt.Printf("  积分已退还 50\n")
   return nil
  },
 )

 // 步骤4: 执行支付(模拟场景:成功 / 余额不足失败)
 paymentShouldFail := false // 切换此值演示不同结果
 saga.AddStep(
  "支付扣款",
  func(ctx context.Context) error {
   if paymentShouldFail {
    return errors.New("用户余额不足")
   }
   dbPaid = true
   fmt.Printf("  支付成功: ¥299.00\n")
   return nil
  },
  func(ctx context.Context) error {
   dbPaid = false
   fmt.Printf("  退款完成: ¥299.00\n")
   return nil
  },
 )

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

// 模拟支付失败触发的完整回滚
func PlaceOrderWithPaymentFailure() {
 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 { fmt.Println("  退款处理"); return nil },
 )

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

5.3 编排式 vs 协调式

维度编排式 (Choreography)协调式 (Orchestration)
中心节点无,去中心化有 Saga Orchestrator
服务耦合服务需感知上下游事件服务只与协调器交互
流程可见性较难实时追踪(需拼凑事件链)协调器集中记录,易于监控
单点风险协调者故障需重新拉起
测试难度高(需模拟事件链路)中(可单独测试协调器逻辑)
适合场景服务少 (≤4个)、流程简单服务多、流程复杂、需强监控

六、本地消息表 (Local Message Table)

6.1 工作原理

本地消息表是最早被大规模应用的分布式事务方案之一(源自 eBay),核心思路是将分布式事务拆解为一系列本地事务 + 可靠消息投递

  1. 业务执行和消息写入在同一个本地事务中完成——要么一起成功,要么一起失败
  2. 定时任务扫描"待发送"状态的消息,投递到消息队列(如 Kafka、RabbitMQ、RocketMQ)
  3. 下游消费者处理消息,成功后将消息标记为"已完成"
  4. 投递失败时,定时任务在下一次扫描时重试(需设置最大重试次数,超过后告警人工介入)

关键保障:本地事务的原子性保证了业务操作和消息记录的一致性。消息投递的最终一致性通过重试 + 下游幂等消费来保证。

6.2 Go 实现

package localmsg

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

// Message 本地消息
type Message struct {
 ID         string    `json:"id"`
 Topic      string    `json:"topic"`       // 目标消息队列 Topic
 Payload    string    `json:"payload"`     // 消息内容(通常为 JSON)
 Status     int       `json:"status"`      // 0:待发送 1:已发送 2:已完成
 RetryCount int       `json:"retry_count"` // 已重试次数
 MaxRetry   int       `json:"max_retry"`   // 最大重试次数
 NextRetry  time.Time `json:"next_retry"`  // 下次重试时间
 CreatedAt  time.Time `json:"created_at"`
}

// MessageStore 消息持久化存储
type MessageStore interface {
 // SaveInTx 在同一事务中保存消息和业务数据
 SaveInTx(ctx context.Context, msg Message, bizExec func(tx *sql.Tx) error) error
 // ScanPending 扫描待发送的消息(就绪时间 ≤ now 且重试次数未超限)
 ScanPending(ctx context.Context, batchSize int) ([]Message, error)
 // MarkSent 标记为已发送
 MarkSent(ctx context.Context, id string) error
 // MarkDone 标记为已完成(消费方确认后回调)
 MarkDone(ctx context.Context, id string) error
 // IncrementRetry 增加重试计数并计算下次重试时间(指数退避)
 IncrementRetry(ctx context.Context, id string) error
}

// =========================================
// 表结构 SQL(应在应用启动时初始化)
// =========================================
const CreateMessageTableSQL = `
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,
    max_retry   INT           NOT NULL DEFAULT 5,
    next_retry  DATETIME      NOT NULL,
    created_at  DATETIME      NOT NULL,
    INDEX idx_status_next_retry (status, next_retry)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
`

// =========================================
// 消息投递 Worker
// =========================================

// MQProducer 消息队列生产者(可替换为 Kafka/RabbitMQ 实现)
type MQProducer interface {
 Publish(ctx context.Context, topic string, payload string) error
}

// MessageDeliveryWorker 消息投递工作器
type MessageDeliveryWorker struct {
 store     MessageStore
 producer  MQProducer
 batchSize int
 interval  time.Duration
}

func NewDeliveryWorker(store MessageStore, producer MQProducer) *MessageDeliveryWorker {
 return &MessageDeliveryWorker{
  store:     store,
  producer:  producer,
  batchSize: 100,
  interval:  5 * time.Second,
 }
}

// Run 启动投递循环
func (w *MessageDeliveryWorker) Run(ctx context.Context) {
 ticker := time.NewTicker(w.interval)
 defer ticker.Stop()

 for {
  select {
  case <-ctx.Done():
   return
  case <-ticker.C:
   w.processBatch(ctx)
  }
 }
}

func (w *MessageDeliveryWorker) processBatch(ctx context.Context) {
 msgs, err := w.store.ScanPending(ctx, w.batchSize)
 if err != nil {
  fmt.Printf("[Worker] 扫描消息失败: %v\n", err)
  return
 }

 for _, msg := range msgs {
  mCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
  if err := w.producer.Publish(mCtx, msg.Topic, msg.Payload); err != nil {
   fmt.Printf("[Worker] 投递消息 %s 失败: %v\n", msg.ID, err)
   // 重试次数未到上限则继续重试,超过则告警
   if msg.RetryCount >= msg.MaxRetry {
    fmt.Printf("[Worker] ⚠ 消息 %s 超过最大重试次数,需人工处理\n", msg.ID)
    continue
   }
   w.store.IncrementRetry(mCtx, msg.ID)
   cancel()
   continue
  }

  w.store.MarkSent(mCtx, msg.ID)
  cancel()
  fmt.Printf("[Worker] 消息 %s 投递成功\n", msg.ID)
 }
}

// =========================================
// 业务层:调用示例(用户注册送积分)
// =========================================

// User 用户信息
type User struct {
 ID    string `json:"id"`
 Name  string `json:"name"`
 Email string `json:"email"`
}

// RegisterUserWithBonus 注册送积分:保证用户记录和积分发放通知一致
func RegisterUserWithBonus(ctx context.Context, store MessageStore, user User) error {
 pointsMsg := map[string]interface{}{
  "user_id": user.ID,
  "points":  1000,
  "reason":  "新用户注册奖励",
 }
 payload, _ := json.Marshal(pointsMsg)

 return store.SaveInTx(ctx, Message{
  ID:        fmt.Sprintf("msg-%s-%d", user.ID, time.Now().UnixNano()),
  Topic:     "user.points.bonus",
  Payload:   string(payload),
  MaxRetry:  5,
  NextRetry: time.Now(),
 }, func(tx *sql.Tx) error {
  // 同一个事务中执行业务操作:插入用户记录
  _, err := tx.ExecContext(ctx,
   "INSERT INTO users (id, name, email) VALUES (?, ?, ?)",
   user.ID, user.Name, user.Email,
  )
  return err // 如果插入失败,消息也不会记录
 })
}

6.3 RocketMQ 事务消息

RocketMQ 内置了事务消息机制,实现原理类似本地消息表:发送"半消息"(Half Message,消费者不可见)→ 执行本地事务 → 提交(Commit)使消息对消费者可见,或回滚(Rollback)丢弃消息。这相当于 RocketMQ 替你实现了本地消息表的框架部分,业务代码只需实现事务执行器和事务回查接口。

七、方案对比与选型

方案一致性级别性能实现复杂度典型适用场景
2PC (XA)强一致性低(全局锁、同步阻塞)同构数据库间分布式事务(MySQL XA),短事务
TCC最终一致性(强保证)高(无全局锁)高(需实现三阶段 + 幂等 + 空回滚防护)核心金融交易(转账、扣款),高一致性高并发
Saga最终一致性中(需设计补偿操作)长流程业务(电商下单、旅行预订),多服务编排
本地消息表最终一致性低(只需一张消息表 + 投递 Worker)跨服务异步通知(发积分、推搜索索引、发推送)
MQ 事务消息最终一致性低(框架已封装)类似本地消息表,适合 RocketMQ 生态

八、最佳实践建议

  1. 能不用就不用:首先审视是否真的需要分布式事务。通过合理的服务拆分(将需要强一致的逻辑合并到一个服务中)、事件驱动 + 异步对账等方式规避,比使用分布式事务更高效。
  2. 核心链路用 TCC:对资金安全零容忍的核心链路(支付、退款、结算),TCC 的预留-确认机制能提供最接近强一致的保证。
  3. 业务编排用 Saga:涉及 3 个以上服务的长流程,使用 Saga + 补偿操作。推荐编排式以便集中管控和监控。
  4. 消息通知用本地消息表:对最终一致性容忍度高但要求可靠投递的场景(发积分、推搜索引擎、发 App 推送),本地消息表 + MQ 是最简单可靠的方案。
  5. 始终保证幂等:任何分布式事务操作(Try/Confirm/Cancel/Action/Compensate)都可能被网络重试触发,必须设计为幂等的。
  6. 可观测性是生命线:每个事务步骤的状态变更都应记录并上报监控系统。当补偿失败或重试超时时,必须触发告警通知运维人员人工介入。
  7. 考虑成熟框架:Java 生态的 Seata(支持 AT/TCC/SAGA/XA),Go 生态的 DTM(dtm-lifez.cn)和 Go-Saga 可直接使用,大幅降低开发与维护成本。

九、总结

分布式事务没有"银弹"。CAP 定理告诉我们完美不可得——强一致、高可用、分区容错三者不可兼得。在实际工程中:

  • 先试图规避分布式事务(合并服务、事件驱动对账)
  • 需要时优先选择最终一致性方案以获得更好的性能
  • 核心金融链路使用 TCC + 本地消息表双重保障
  • 建立完善的事务状态监控和补偿告警机制

理解每种方案的底层原理(强一致 vs 最终一致、锁机制 vs 补偿机制),再结合具体业务场景做出合理选择,是分布式系统设计中最为核心的能力之一。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部