一、为什么需要分布式事务
在单体架构中,一个业务操作只需在一个数据库中完成 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)。整个事务分为两个阶段执行:
- 准备阶段 (Voting Phase):协调者向所有参与者发送 Prepare 请求。参与者执行事务操作但不提交(锁定资源、写 Undo/Redo 日志),然后返回 Yes 或 No。
- 提交阶段 (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),核心思路是将分布式事务拆解为一系列本地事务 + 可靠消息投递:
- 业务执行和消息写入在同一个本地事务中完成——要么一起成功,要么一起失败
- 定时任务扫描"待发送"状态的消息,投递到消息队列(如 Kafka、RabbitMQ、RocketMQ)
- 下游消费者处理消息,成功后将消息标记为"已完成"
- 投递失败时,定时任务在下一次扫描时重试(需设置最大重试次数,超过后告警人工介入)
关键保障:本地事务的原子性保证了业务操作和消息记录的一致性。消息投递的最终一致性通过重试 + 下游幂等消费来保证。
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 生态 |
八、最佳实践建议
- 能不用就不用:首先审视是否真的需要分布式事务。通过合理的服务拆分(将需要强一致的逻辑合并到一个服务中)、事件驱动 + 异步对账等方式规避,比使用分布式事务更高效。
- 核心链路用 TCC:对资金安全零容忍的核心链路(支付、退款、结算),TCC 的预留-确认机制能提供最接近强一致的保证。
- 业务编排用 Saga:涉及 3 个以上服务的长流程,使用 Saga + 补偿操作。推荐编排式以便集中管控和监控。
- 消息通知用本地消息表:对最终一致性容忍度高但要求可靠投递的场景(发积分、推搜索引擎、发 App 推送),本地消息表 + MQ 是最简单可靠的方案。
- 始终保证幂等:任何分布式事务操作(Try/Confirm/Cancel/Action/Compensate)都可能被网络重试触发,必须设计为幂等的。
- 可观测性是生命线:每个事务步骤的状态变更都应记录并上报监控系统。当补偿失败或重试超时时,必须触发告警通知运维人员人工介入。
- 考虑成熟框架:Java 生态的 Seata(支持 AT/TCC/SAGA/XA),Go 生态的 DTM(dtm-lifez.cn)和 Go-Saga 可直接使用,大幅降低开发与维护成本。
九、总结
分布式事务没有"银弹"。CAP 定理告诉我们完美不可得——强一致、高可用、分区容错三者不可兼得。在实际工程中:
- 先试图规避分布式事务(合并服务、事件驱动对账)
- 需要时优先选择最终一致性方案以获得更好的性能
- 核心金融链路使用 TCC + 本地消息表双重保障
- 建立完善的事务状态监控和补偿告警机制
理解每种方案的底层原理(强一致 vs 最终一致、锁机制 vs 补偿机制),再结合具体业务场景做出合理选择,是分布式系统设计中最为核心的能力之一。

发表评论 取消回复