微服务架构下,跨服务的数据一致性问题无法依赖数据库本地事务解决。Saga模式通过将长事务拆分为多个本地事务并配合补偿操作,保证最终一致性。TCC(Try-Confirm-Cancel)模式则通过两阶段提交实现强一致性的业务级事务。本文用Go语言实现两种分布式事务模式,覆盖协调器设计、幂等性保障和故障恢复。
分布式事务问题与Saga模式原理
以电商订单场景为例:创建订单需要调用库存服务扣减库存、支付服务扣款、积分服务增加积分。任一步骤失败都需要回滚已执行的操作。传统ACID事务无法跨服务边界,Saga模式将全局事务拆分为有序的本地事务序列,每一步都有对应的补偿操作。
Saga执行流程:
正向执行:
T1: 创建订单 → C1: 取消订单
T2: 扣减库存 → C2: 恢复库存
T3: 扣款 → C3: 退款
T4: 加积分 → C4: 扣积分
正常流程: T1 → T2 → T3 → T4
异常流程: T1 → T2 → T3(失败) → C2 → C1
Saga有两种协调方式:编排式(Orchestration)由中央协调器控制流程;协同式(Choreography)各服务通过事件驱动自行协调。编排式更易维护和调试,适合复杂业务流程。
Saga协调器实现:编排式事务流程
定义事务步骤接口和协调器:
package saga
import (
"context"
"errors"
"log"
)
// Step 定义Saga事务步骤
type Step struct {
Name string
Execute func(ctx context.Context) error
Compensate func(ctx context.Context) error
}
// Coordinator Saga协调器
type Coordinator struct {
steps []Step
}
func NewCoordinator(steps []Step) *Coordinator {
return &Coordinator{steps: steps}
}
// Execute 执行Saga事务
func (c *Coordinator) Execute(ctx context.Context) error {
completed := make([]int, 0, len(c.steps))
for i, step := range c.steps {
if err := step.Execute(ctx); err != nil {
log.Printf("步骤 %s 执行失败: %v,开始补偿", step.Name, err)
// 逆序执行已完成步骤的补偿操作
for j := len(completed) - 1; j >= 0; j-- {
compStep := c.steps[completed[j]]
if compErr := compStep.Compensate(ctx); compErr != nil {
log.Printf("补偿步骤 %s 失败: %v", compStep.Name, compErr)
// 记录补偿失败,需要人工介入
return errors.New("补偿失败,需人工介入: " + compStep.Name)
}
log.Printf("补偿步骤 %s 完成", compStep.Name)
}
return err
}
completed = append(completed, i)
log.Printf("步骤 %s 执行成功", step.Name)
}
return nil
}
电商订单Saga实现:
package order
import (
"context"
"fmt"
"myapp/saga"
)
type OrderSaga struct {
orderID string
userID string
amount float64
productID string
quantity int
}
func (s *OrderSaga) Build() *saga.Coordinator {
steps := []saga.Step{
{
Name: "create_order",
Execute: func(ctx context.Context) error {
return s.createOrder(ctx)
},
Compensate: func(ctx context.Context) error {
return s.cancelOrder(ctx)
},
},
{
Name: "deduct_inventory",
Execute: func(ctx context.Context) error {
return s.deductInventory(ctx)
},
Compensate: func(ctx context.Context) error {
return s.restoreInventory(ctx)
},
},
{
Name: "charge_payment",
Execute: func(ctx context.Context) error {
return s.chargePayment(ctx)
},
Compensate: func(ctx context.Context) error {
return s.refundPayment(ctx)
},
},
{
Name: "add_points",
Execute: func(ctx context.Context) error {
return s.addPoints(ctx)
},
Compensate: func(ctx context.Context) error {
return s.deductPoints(ctx)
},
},
}
return saga.NewCoordinator(steps)
}
func (s *OrderSaga) createOrder(ctx context.Context) error {
// 写入订单表,状态为PENDING
return db.ExecContext(ctx,
"INSERT INTO orders (id, user_id, amount, status) VALUES (?, ?, ?, 'PENDING')",
s.orderID, s.userID, s.amount)
}
TCC补偿事务模式与Go实现
TCC将每个操作分为Try(资源预留)、Confirm(确认执行)、Cancel(释放预留)三个阶段。与Saga不同,TCC在Try阶段就锁定资源,Confirm/Cancel阶段执行最终操作或释放。
TCC接口定义:
package tcc
import "context"
// TCCService TCC事务参与者接口
type TCCService interface {
Try(ctx context.Context, bid string, data interface{}) error
Confirm(ctx context.Context, bid string) error
Cancel(ctx context.Context, bid string) error
}
// Transaction TCC全局事务
type Transaction struct {
XID string // 全局事务ID
Services []TCCParticipant // 参与者列表
}
type TCCParticipant struct {
Service TCCService
BranchID string // 分支事务ID
Data interface{}
Status string // TRYING, CONFIRMED, CANCELLED
}
TCC协调器实现:
package tcc
import (
"context"
"errors"
"log"
"sync"
)
type Coordinator struct {
tx *Transaction
mu sync.Mutex
}
func (c *Coordinator) Execute(ctx context.Context) error {
c.mu.Lock()
defer c.mu.Unlock()
// Phase 1: Try
tried := make([]int, 0)
for i, p := range c.tx.Services {
if err := p.Service.Try(ctx, p.BranchID, p.Data); err != nil {
log.Printf("Try阶段 %s 失败: %v", p.BranchID, err)
// 逆序Cancel已Try的参与者
for j := len(tried) - 1; j >= 0; j-- {
idx := tried[j]
if cerr := c.tx.Services[idx].Service.Cancel(ctx, c.tx.Services[idx].BranchID); cerr != nil {
log.Printf("Cancel %s 失败: %v", c.tx.Services[idx].BranchID, cerr)
}
}
return err
}
tried = append(tried, i)
c.tx.Services[i].Status = "TRYING"
}
// Phase 2: Confirm
for i, p := range c.tx.Services {
if err := p.Service.Confirm(ctx, p.BranchID); err != nil {
log.Printf("Confirm阶段 %s 失败: %v", p.BranchID, err)
// Confirm失败需要重试或人工介入
return errors.New("confirm失败: " + p.BranchID)
}
c.tx.Services[i].Status = "CONFIRMED"
}
return nil
}
库存服务TCC实现:
type InventoryTCC struct {
repo *InventoryRepo
}
func (s *InventoryTCC) Try(ctx context.Context, bid string, data interface{}) error {
req := data.(*DeductReq)
// 冻结库存:将可用库存移到冻结库存
return s.repo.FreezeStock(ctx, req.ProductID, req.Quantity, bid)
}
func (s *InventoryTCC) Confirm(ctx context.Context, bid string) error {
// 确认扣减:冻结库存直接扣减
return s.repo.ConfirmDeduct(ctx, bid)
}
func (s *InventoryTCC) Cancel(ctx context.Context, bid string) error {
// 取消:冻结库存退回可用库存
return s.repo.UnfreezeStock(ctx, bid)
}
幂等性设计与重试机制
分布式事务的每个操作都可能因网络问题重试执行,必须保证幂等性——同一操作执行多次结果一致。
幂等性实现方案:
// 基于唯一请求ID的幂等控制
type IdempotentExecutor struct {
redis *redis.Client
}
func (e *IdempotentExecutor) Execute(ctx context.Context, reqID string, fn func() error) error {
// 使用Redis SETNX获取执行锁
key := fmt.Sprintf("idempotent:%s", reqID)
ok, err := e.redis.SetNX(ctx, key, "processing", 30*time.Second).Result()
if err != nil {
return err
}
if !ok {
// 检查是否已有结果
result, err := e.redis.Get(ctx, key+"result").Result()
if err == nil && result == "success" {
return nil // 已成功执行,幂等返回
}
return errors.New("请求正在处理中")
}
// 执行业务逻辑
if err := fn(); err != nil {
e.redis.Set(ctx, key+"result", "failed", 30*time.Second)
return err
}
e.redis.Set(ctx, key+"result", "success", 24*time.Hour)
return nil
}
指数退避重试:
func RetryWithBackoff(ctx context.Context, maxRetries int, fn func() error) error {
var lastErr error
for i := 0; i < maxRetries; i++ {
if err := fn(); err != nil {
lastErr = err
backoff := time.Duration(1< 10*time.Second {
backoff = 10 * time.Second
}
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(backoff):
continue
}
}
return nil
}
return fmt.Errorf("重试%d次后仍失败: %w", maxRetries, lastErr)
}
分布式事务监控与故障恢复
分布式事务需要持久化事务状态,支持宕机后恢复。将事务日志写入数据库,协调器重启后从断点继续执行。
事务状态持久化:
type TransactionLog struct {
XID string `json:"xid"`
StepName string `json:"step_name"`
Status string `json:"status"` // STARTED, COMPLETED, COMPENSATED, FAILED
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
// 保存事务步骤状态
func SaveStepStatus(ctx context.Context, log *TransactionLog) error {
_, err := db.ExecContext(ctx,
`INSERT INTO saga_logs (xid, step_name, status, created_at, updated_at)
VALUES (?, ?, ?, ?, ?)
ON DUPLICATE KEY UPDATE status=VALUES(status), updated_at=VALUES(updated_at)`,
log.XID, log.StepName, log.Status, log.CreatedAt, log.UpdatedAt)
return err
}
// 恢复未完成的事务
func RecoverPendingTransactions(ctx context.Context) error {
rows, err := db.QueryContext(ctx,
"SELECT xid, step_name, status FROM saga_logs WHERE status IN ('STARTED', 'FAILED')")
if err != nil {
return err
}
defer rows.Close()
for rows.Next() {
var log TransactionLog
if err := rows.Scan(&log.XID, &log.StepName, &log.Status); err != nil {
continue
}
// 根据状态决定继续执行或补偿
if log.Status == "FAILED" {
// 触发补偿流程
go compensateTransaction(ctx, log.XID)
} else if log.Status == "STARTED" {
// 从断点继续执行
go resumeTransaction(ctx, log.XID)
}
}
return nil
}
定期清理已完成的事务日志避免表膨胀。配合Prometheus采集事务成功率和补偿触发率指标,异常率超过阈值时触发告警。对于补偿失败的案例,需要人工介入处理,确保业务数据最终一致。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-yu-yan-fen-bu-shi-shi-wu-shi-zhan-saga-mo-shi-yu-tcc-bu/