
Go Saga模式:分布式事务摘要: 本篇讲解Go语言Saga分布式事务模式实现编排式Saga用中央协调器驱动多服务流程设计补偿事务回滚已执行步骤用状态机控制Saga流转处理超时与重试分享补偿事务顺序错误导致数据不一致的踩坑经验对比编排式Saga、协作式Saga、两阶段提交三种方案。开篇故事去年做旅游预订系统一笔订单要扣减机票库存、扣款、生成酒店凭证三个服务。库存扣了支付扣款失败机票库存没回滚超卖了50单。客服手动处理了一整天赔了差价。当时用本地事务做不了三个服务三个数据库跨库事务性能太差。调研后选了Saga模式把长事务拆成多个本地事务每步配一个补偿事务失败时按相反顺序补偿。上线后超卖问题再没出现过。这篇把编排式Saga、补偿事务设计、状态机驱动都写清楚。一、编排式Saga:中央协调器Saga有两种风格。编排式(Orchestration)有个中央协调器按顺序调用各服务失败就触发补偿。协作式(Choreography)没有协调器各服务监听事件自行决定下一步。编排式流程清晰可控适合复杂流程。这里实现编排式。packagesagaimport(contextfmtlogtime)// SagaStep Saga的一个步骤// Execute执行正向操作Compensate执行补偿操作typeSagaStepinterface{// Name 步骤名称用于日志和状态记录Name()string// Execute 执行正向操作如扣库存Execute(ctx context.Context,data*SagaData)error// Compensate 执行补偿操作如恢复库存Compensate(ctx context.Context,data*SagaData)error}// SagaData Saga执行过程中的共享数据// 各步骤把结果写进来后续步骤和补偿用typeSagaDatastruct{OrderIDstring// 订单IDUserIDstring// 用户IDAmountfloat64// 金额TicketRefstring// 机票预订号步骤1产出PaymentRefstring// 支付流水号步骤2产出HotelRefstring// 酒店凭证号步骤3产出}// BookTicketStep 预订机票步骤typeBookTicketStepstruct{}// Name 步骤名称func(s*BookTicketStep)Name()string{returnBookTicket}// Execute 扣减机票库存func(s*BookTicketStep)Execute(ctx context.Context,data*SagaData)error{// 模拟调用机票服务扣库存data.TicketReffmt.Sprintf(TKT-%s,data.OrderID)log.Printf(机票库存已扣减, 订单%s, 凭证%s,data.OrderID,data.TicketRef)returnnil}// Compensate 恢复机票库存func(s*BookTicketStep)Compensate(ctx context.Context,data*SagaData)error{// 模拟调用机票服务恢复库存log.Printf(机票库存已恢复, 凭证%s,data.TicketRef)returnnil}// ChargePaymentStep 扣款步骤typeChargePaymentStepstruct{}func(s*ChargePaymentStep)Name()string{returnChargePayment}// Execute 扣款func(s*ChargePaymentStep)Execute(ctx context.Context,data*SagaData)error{// 模拟支付失败的情况用于演示补偿ifdata.Amount10000{returnfmt.Errorf(余额不足, 金额%.2f,data.Amount)}data.PaymentReffmt.Sprintf(PAY-%s,data.OrderID)log.Printf(扣款成功, 金额%.2f, 流水%s,data.Amount,data.PaymentRef)returnnil}// Compensate 退款func(s*ChargePaymentStep)Compensate(ctx context.Context,data*SagaData)error{log.Printf(已退款, 流水%s,data.PaymentRef)returnnil}// BookHotelStep 预订酒店步骤typeBookHotelStepstruct{}func(s*BookHotelStep)Name()string{returnBookHotel}func(s*BookHotelStep)Execute(ctx context.Context,data*SagaData)error{data.HotelReffmt.Sprintf(HTL-%s,data.OrderID)log.Printf(酒店凭证已生成, 凭证%s,data.HotelRef)returnnil}func(s*BookHotelStep)Compensate(ctx context.Context,data*SagaData)error{log.Printf(酒店凭证已取消, 凭证%s,data.HotelRef)returnnil}// Orchestrator Saga中央协调器// 按顺序执行步骤失败时按相反顺序补偿typeOrchestratorstruct{steps[]SagaStep// 正向步骤列表}// NewOrchestrator 创建协调器funcNewOrchestrator(steps...SagaStep)*Orchestrator{returnOrchestrator{steps:steps}}// Execute 执行整个Saga流程// 成功返回nil失败自动补偿已执行步骤func(o*Orchestrator)Execute(ctx context.Context,data*SagaData)error{// 记录已成功执行的步骤用于补偿completed:make([]int,0,len(o.steps))// 正向执行每个步骤fori,step:rangeo.steps{iferr:step.Execute(ctx,data);err!nil{// 某步失败开始补偿log.Printf(步骤 %s 执行失败: %v, 开始补偿,step.Name(),err)o.compensate(ctx,completed,data)returnfmt.Errorf(saga failed at step %s: %w,step.Name(),err)}// 记录成功步骤索引completedappend(completed,i)log.Printf(步骤 %s 执行成功,step.Name())}returnnil}// compensate 按相反顺序补偿已执行步骤func(o*Orchestrator)compensate(ctx context.Context,completed[]int,data*SagaData){// 逆序遍历已完成的步骤fori:len(completed)-1;i0;i--{step:o.steps[completed[i]]// 补偿失败不中断继续补偿其他步骤iferr:step.Compensate(ctx,data);err!nil{log.Printf(步骤 %s 补偿失败: %v, 需人工介入,step.Name(),err)}}}协调器是Saga的大脑。它知道步骤顺序知道失败后怎么补偿。每个步骤只管自己的正向和补偿逻辑不关心其他步骤。二、状态机驱动与超时重试实际生产中Saga要持久化状态断电重启能恢复。用状态机驱动每个步骤的状态转换配超时和重试机制。packagesagaimport(contextfmtsynctime)// StepStatus 步骤状态typeStepStatusstringconst(StatusPending StepStatuspending// 待执行StatusRunning StepStatusrunning// 执行中StatusCompleted StepStatuscompleted// 已完成StatusCompensat StepStatuscompensating// 补偿中StatusFailed StepStatusfailed// 失败)// StepState 步骤状态记录typeStepStatestruct{Indexint// 步骤索引Namestring// 步骤名Status StepStatus// 当前状态Attemptsint// 重试次数MaxRetryint// 最大重试次数Timeout time.Duration// 超时时间}// SagaState Saga整体状态typeSagaStatestruct{IDstring// Saga实例IDData*SagaData// 共享数据Steps[]StepState// 各步骤状态mu sync.Mutex// 并发保护}// StateMachine 状态机驱动的Saga协调器typeStateMachinestruct{orchestrator*Orchestrator states sync.Map// sagaID - *SagaState}// NewStateMachine 创建状态机funcNewStateMachine(o*Orchestrator)*StateMachine{returnStateMachine{orchestrator:o}}// ExecuteWithRetry 带超时和重试的执行func(sm*StateMachine)ExecuteWithRetry(ctx context.Context,sagaIDstring,data*SagaData,)error{// 初始化Saga状态state:SagaState{ID:sagaID,Data:data,Steps:make([]StepState,len(sm.orchestrator.steps)),}// 每个步骤默认重试3次超时5秒fori,step:rangesm.orchestrator.steps{state.Steps[i]StepState{Index:i,Name:step.Name(),Status:StatusPending,MaxRetry:3,Timeout:5*time.Second,}}sm.states.Store(sagaID,state)// 执行步骤带重试completed:make([]int,0,len(state.Steps))fori:rangestate.Steps{iferr:sm.executeStep(ctx,state,i);err!nil{// 步骤最终失败开始补偿sm.compensateFrom(ctx,state,completed)returnfmt.Errorf(saga %s failed at step %d: %w,sagaID,i,err)}completedappend(completed,i)}returnnil}// executeStep 执行单个步骤带超时和重试func(sm*StateMachine)executeStep(ctx context.Context,state*SagaState,indexint,)error{step:sm.orchestrator.steps[index]stepState:state.Steps[index]// 带超时的contextforattempt:0;attemptstepState.MaxRetry;attempt{stepState.Attemptsattempt1stepState.StatusStatusRunning// 执行上下文带超时timeoutCtx,cancel:context.WithTimeout(ctx,stepState.Timeout)err:step.Execute(timeoutCtx,state.Data)cancel()iferrnil{stepState.StatusStatusCompletedreturnnil}// 失败记录并准备重试fmt.Printf(步骤 %s 第%d次执行失败: %v\n,step.Name(),attempt1,err)// 指数退避等待ifattemptstepState.MaxRetry-1{time.Sleep(time.Duration(1attempt)*time.Second)}}stepState.StatusStatusFailedreturnfmt.Errorf(step %s failed after %d attempts,step.Name(),stepState.MaxRetry)}// compensateFrom 从指定步骤开始按相反顺序补偿func(sm*StateMachine)compensateFrom(ctx context.Context,state*SagaState,completed[]int,){fori:len(completed)-1;i0;i--{idx:completed[i]step:sm.orchestrator.steps[idx]state.Steps[idx].StatusStatusCompensat// 补偿也带超时重试timeoutCtx,cancel:context.WithTimeout(ctx,5*time.Second)err:step.Compensate(timeoutCtx,state.Data)cancel()iferr!nil{fmt.Printf(步骤 %s 补偿失败: %v, 需人工介入\n,step.Name(),err)}}}状态机的好处是可恢复。Saga执行到一半挂了重启后从状态机记录的位置继续。每个步骤的重试次数和状态都持久化了不会重复执行。三、独家踩坑:补偿事务顺序错误导致数据不一致这个坑上线一周后才暴露。Saga有三个步骤: 扣库存、扣款、发券。某次扣款失败补偿本该按退款、恢复库存的逆序执行。结果补偿代码写错了顺序先恢复库存再退款库存恢复了但退款失败库存多了钱没退财务对账差了一笔。// 错误补偿: 顺序写反了func(o*BadOrchestrator)compensate(completed[]int){// 正序遍历先补偿第一步(恢复库存)再补偿第二步(退款)// 退款失败时库存已经恢复状态不一致fori:0;ilen(completed);i{o.steps[completed[i]].Compensate(ctx,data)}}补偿必须严格按正向执行的逆序。正向是A、B、C补偿就是C、B、A。而且补偿本身要幂等因为补偿可能被重试。packagesagaimport(contextsync)// IdempotentCompensator 幂等补偿器// 用已执行记录保证补偿只生效一次typeIdempotentCompensatorstruct{mu sync.Mutex compensatedmap[string]bool// 已补偿的步骤记录}// NewIdempotentCompensator 创建幂等补偿器funcNewIdempotentCompensator()*IdempotentCompensator{returnIdempotentCompensator{compensated:make(map[string]bool),}}// Compensate 幂等补偿// 同一步骤多次调用只生效一次func(c*IdempotentCompensator)Compensate(ctx context.Context,step SagaStep,data*SagaData,)error{c.mu.Lock()deferc.mu.Unlock()// 已补偿过直接返回成功ifc.compensated[step.Name()]{returnnil}// 执行补偿iferr:step.Compensate(ctx,data);err!nil{returnerr}// 记录已补偿c.compensated[step.Name()]truereturnnil}经验是补偿顺序和幂等都要有测试覆盖。上线前用混沌测试故意让中间步骤失败验证补偿是否按逆序执行、是否幂等。补偿失败的步骤要告警人工介入不能静默吞掉。四、对比分析方案一致性性能复杂度可恢复适用场景两阶段提交强一致差中弱传统数据库编排式Saga最终一致好高强复杂跨服务流程协作式Saga最终一致好中中简单事件驱动流程两阶段提交保证强一致但资源锁定时间长性能差微服务里基本不用。编排式Saga有中央协调器流程清晰可控状态可持久化恢复适合订单、支付这类复杂流程。协作式Saga靠事件驱动没有协调器耦合低但流程隐式排查问题困难适合简单流程。总结与预告Saga把跨服务长事务拆成多个本地事务每步配补偿事务失败时逆序补偿。编排式用中央协调器控制流程状态机驱动并持久化状态支持超时重试和断点恢复。补偿顺序必须严格逆序补偿本身要幂等。下一篇聊DDD领域驱动设计看限界上下文怎么划分、聚合根怎么设计。