Go微服务中分布式事务Saga模式实现:从编排到补偿的完整方案

微服务架构下分布式事务的核心矛盾

单体应用中,数据库本地事务(ACID)保证了一致性。微服务拆分后,一个业务操作可能跨越订单、库存、支付、通知等多个服务,每个服务独享自己的数据库,本地事务无法覆盖跨服务的数据一致性。直接使用2PC(两阶段提交)会引入同步阻塞和单点故障,在微服务场景下不实用。Saga模式通过将长事务拆分为一系列本地事务,配合补偿操作,在最终一致性和可用性之间取得平衡。

Saga编排模式 vs 协调模式

Saga有两种实现风格:

编排模式(Choreography):各服务通过事件驱动自行决定下一步动作。优点是松耦合,缺点是流程逻辑分散在各服务中,难以全局追踪。

协调模式(Orchestration):由一个Saga协调器统一编排步骤,各服务只暴露本地事务接口和补偿接口。优点是流程可见、便于监控,缺点是协调器是单点。

在业务流程较复杂(超过5个步骤)或需要人工审批环节的场景中,协调模式的优势更明显。以下实现基于协调模式。

Go语言Saga协调器实现

package saga

import (
	"context"
	"fmt"
	"log"
)

type Step struct {
	Name       string
	Execute    func(ctx context.Context) error
	Compensate func(ctx context.Context) error
}

type Saga struct {
	Name  string
	Steps []Step
}

type Result struct {
	Success     bool
	FailedStep  string
	Compensated []string
	Err         error
}

func (s *Saga) Execute(ctx context.Context) Result {
	result := Result{Success: false}
	completedSteps := make([]Step, 0, len(s.Steps))

	for i, step := range s.Steps {
		if err := step.Execute(ctx); err != nil {
			log.Printf("[Saga] step %q failed: %v", step.Name, err)
			result.FailedStep = step.Name
			result.Err = err
			result.Compensated = s.compensate(ctx, completedSteps)
			return result
		}
		completedSteps = append(completedSteps, step)
		log.Printf("[Saga] step %q done (%d/%d)", step.Name, i+1, len(s.Steps))
	}

	result.Success = true
	return result
}

func (s *Saga) compensate(ctx context.Context, completedSteps []Step) []string {
	compensated := make([]string, 0)
	for i := len(completedSteps) - 1; i >= 0; i-- {
		step := completedSteps[i]
		if step.Compensate == nil {
			continue
		}
		if err := step.Compensate(ctx); err != nil {
			log.Printf("[Saga] step %q compensate failed: %v", step.Name, err)
			continue
		}
		compensated = append(compensated, step.Name)
	}
	return compensated
}

实战:电商下单流程的Saga实现

package order

import (
	"context"
	"fmt"
	"saga"
)

type OrderService struct {
	orderRepo       OrderRepository
	inventorySvc    InventoryService
	paymentSvc      PaymentService
	notificationSvc NotificationService
}

func (svc *OrderService) CreateOrder(ctx context.Context, req CreateOrderReq) error {
	s := &saga.Saga{
		Name: "create-order",
		Steps: []saga.Step{
			{
				Name: "create-order-record",
				Execute: func(ctx context.Context) error {
					return svc.orderRepo.Create(ctx, req)
				},
				Compensate: func(ctx context.Context) error {
					return svc.orderRepo.MarkCancelled(ctx, req.OrderID)
				},
			},
			{
				Name: "deduct-inventory",
				Execute: func(ctx context.Context) error {
					return svc.inventorySvc.Deduct(ctx, req.SKUID, req.Quantity)
				},
				Compensate: func(ctx context.Context) error {
					return svc.inventorySvc.Restore(ctx, req.SKUID, req.Quantity)
				},
			},
			{
				Name: "charge-payment",
				Execute: func(ctx context.Context) error {
					return svc.paymentSvc.Charge(ctx, req.UserID, req.Amount)
				},
				Compensate: func(ctx context.Context) error {
					return svc.paymentSvc.Refund(ctx, req.UserID, req.Amount)
				},
			},
			{
				Name: "send-notification",
				Execute: func(ctx context.Context) error {
					_ = svc.notificationSvc.SendOrderConfirm(ctx, req)
					return nil
				},
				Compensate: nil,
			},
		},
	}

	result := s.Execute(ctx)
	if !result.Success {
		return fmt.Errorf("order failed: step[%s], compensated: %v, err: %w",
			result.FailedStep, result.Compensated, result.Err)
	}
	return nil
}

Saga持久化:应对协调器崩溃

协调器本身是单点,进程崩溃后需要能恢复到崩溃前的状态。持久化方案是:将Saga的每一步执行状态写入数据库,崩溃恢复后从断点继续:

CREATE TABLE saga_instance (
    id           BIGINT PRIMARY KEY AUTO_INCREMENT,
    saga_name    VARCHAR(128) NOT NULL,
    status       ENUM('running', 'compensating', 'completed', 'failed') NOT NULL,
    current_step INT NOT NULL DEFAULT 0,
    payload      JSON,
    created_at   DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
    updated_at   DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    INDEX idx_status (status),
    INDEX idx_created (created_at)
);

每完成一个步骤,更新current_step和status。补偿阶段同样记录每步的补偿状态。重启时查询status为running或compensating的记录,从current_step继续执行。补偿失败的步骤标记为failed,由人工或定时任务重试。

Saga模式的使用边界

Saga模式不适用于所有分布式事务场景。以下情况需要重新评估:

1. 强一致性要求:如果业务要求多个服务的数据在同一时刻必须一致(如银行转账),Saga的最终一致性模型无法满足,需要考虑TCC模式。

2. 补偿操作不可逆:如果某个步骤的副作用无法补偿(如已发出的短信、已触发的IoT设备动作),Saga的补偿机制失效,需要设计幂等的替代补偿方案。

3. 步骤过多导致补偿链过长:超过8个步骤的Saga,补偿链的复杂度和失败概率会急剧上升,此时应考虑拆分为多个子Saga。

Go语言实现Saga的关键在于简洁——不需要引入重型框架,一个不到100行的协调器就能覆盖核心逻辑。真正的难点不在代码,在于补偿操作的设计:每个正向操作必须有对应的补偿操作,补偿操作必须幂等,补偿失败必须有兜底方案。这三个约束满足后,Saga模式在微服务架构下可以稳定运行。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/go-wei-fu-wu-zhong-fen-bu-shi-shi-wu-saga-mo-shi-shi-xian/

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

相关推荐