Go语言分布式事务实战:Saga模式与TCC补偿机制实现

微服务架构下,跨服务的数据一致性问题无法依赖数据库本地事务解决。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/

(0)
小编小编
上一篇 20小时前
下一篇 20小时前

相关推荐