一、分布式事务导论
在微服务架构中,一个业务操作往往需要跨越多个服务和数据库。例如电商下单场景需要同时操作订单服务、库存服务、积分服务和支付服务,这就要求我们在多个独立的数据库之间保证数据一致性。分布式事务正是为了解决这一挑战而诞生的话题。
与单机数据库的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),分为两个阶段:
- 准备阶段(Voting Phase):协调者向所有参与者发送Prepare请求,参与者执行事务但不提交,锁定资源并记录Undo/Redo日志
- 提交阶段(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 工作原理
本地消息表的核心思想是将分布式事务拆分为本地事务,通过保证本地事务和消息投递的一致性来实现最终一致性:
- 业务执行和消息插入在同一个本地事务中完成
- 定时任务扫描未发送的消息,投递到消息队列
- 下游消费者处理消息,成功后回调确认
- 失败时通过重试机制保证最终送达
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 | 最终一致 | 高 | 中等 | 长流程业务,如电商下单 |
| 本地消息表 | 最终一致 | 高 | 低 | 跨服务数据同步,如积分发放 |
| 最大努力通知 | 弱一致 | 最高 | 最低 | 对一致性要求极低的场景 |
八、最佳实践建议
- 优先考虑最终一致性:大多数业务场景无需强一致,BASE理论提供的最终一致性已经足够,且能带来更好的性能和可用性
- TCC用于核心链路:对一致性要求较高的核心链路(如支付、扣款),推荐使用TCC模式
- Saga应对长流程:涉及3个以上服务的复杂业务流程,使用Saga模式配合补偿机制
- 本地消息表作为兜底:作为最终手段保证操作消息不丢失,可与上述任何方案配合使用
- 始终保证幂等性:网络调用不可靠,所有分布式事务操作必须设计为幂等的
- 记录并监控补偿失败:补偿操作也可能失败,必须记录日志并及时告警,必要时人工介入
- 超时要短,重试要谨慎:参与者操作的超时时间应远短于全局超时,重试需有上限和退避策略
- 考虑引入成熟框架:如Seata(支持AT/TCC/SAGA/XA模式)、DTM(Go语言分布式事务管理器)可大幅降低实现成本
九、总结
分布式事务没有银弹,每种方案都有其适用场景和代价。在实际工程中,我们应该:
- 首先判断是否真的需要分布式事务——能通过合理的服务设计规避最好
- 需要时优先选择性能更好的弱一致性方案
- 核心链路使用TCC或Saga + 本地消息表的双重保障
- 建立完善的监控和补偿机制,确保问题可追溯、可修复
CAP定理告诉我们完美不可得,而工程师的艺术在于做出最合理的取舍。希望本文提供的理论分析和Go代码示例,能为你的分布式系统设计提供有价值的参考。

发表评论 取消回复