首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >分布式事务 SAGA 深度落地:补偿语义、异常防御与流程编排

分布式事务 SAGA 深度落地:补偿语义、异常防御与流程编排

原创
作者头像
Alan_751
发布2026-07-20 10:15:19
发布2026-07-20 10:15:19
410
举报

最近在做跨系统采购链路的对接,整条链路涉及订单、库存、支付三套独立的外部系统,全部通过HTTP API交互。最开始图省事,按顺序串行调用接口,用本地事务包一层,结果上线没多久就出了两次问题:中间某一步超时失败,前面已经执行成功的操作没法回滚,要么库存多扣要么资金单边账,资损风险极高。

调研了一圈分布式事务方案,2PC/XA对跨系统场景太重,强一致的阻塞开销我们也接受不了,最终选了SAGA柔性事务。落地的过程踩了不少坑——网上大部分资料只讲正向流程+反向补偿的Happy Path,真放到生产环境,网络乱序、超时重试、服务宕机这些问题全出来了。这里把核心的三块实践经验整理出来,附可直接复用的Python实现。

一、补偿的核心不是“反向调用”,是业务语义撤销

很多人刚接触SAGA的时候,会把补偿简单理解成“调反向接口”:扣了库存就加回去,付了款就退回去。这是最常见的误区,也是线上事故的高发点。

SAGA的本质是把一个长事务拆成N个独立的本地事务,每一步执行完就立即提交、释放连接资源。当第K步失败时,按逆序对前面K-1步执行补偿,最终让系统回到事务开始前的等价状态。注意这里说的是等价状态,不是“数据完全复原”——本地事务一旦提交就没法rollback,补偿是一个全新的业务操作,用反向语义抵消正向操作的业务影响。

举个最常见的下单链路例子:

  • 创建订单的补偿不是删除订单,是把订单状态置为已取消,保留单据留痕
  • 扣减库存的补偿不是直接加库存,是归还本次事务占用的对应数量
  • 支付扣款的补偿不是撤销扣款记录,是发起一笔原路退款

我们落地时给补偿操作定了三条硬标准:必须幂等、语义可逆、执行后不可再正向重入。这三条没守住,补偿就是新的事故源。

生产必踩的三个异常坑与防御方案

跨系统API调用天然存在网络延迟、丢包、乱序,只写“正向+补偿”两层逻辑,线上必出问题。我们前前后后踩了三遍坑,才把这三道防御补全。

1. 幂等性缺失:超时重试导致重复执行

踩坑现场:第一次压测模拟网络抖动,库存接口超时3秒,调用端按默认策略重试了2次,最终库存被扣了3次。

根因:HTTP超时场景下,调用方永远无法确定下游到底有没有执行业务逻辑。框架层的重试机制一触发,重复扣减、重复退款是必然结果。

解决方案:全链路透传saga_id + branch_id,下游基于唯一键做幂等控制。用数据库唯一键冲突判重,比先查后插性能更好,也能避免并发场景下的判断失效。

代码语言:python
复制
def deduct_inventory(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存扣减接口(幂等实现)"""
    unique_key = f"{saga_id}:{branch_id}:deduct"
    
    # 先插幂等记录,靠唯一键挡重复请求
    try:
        db.execute(
            "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
            "VALUES (%s, %s, 'deduct', 'PROCESSING')",
            (saga_id, branch_id)
        )
    except IntegrityError:
        # 重复请求直接返回已有状态,不执行业务
        status = db.query_status(saga_id, branch_id, 'deduct')
        return status == 'SUCCESS'
    
    # 执行业务逻辑
    try:
        affected = db.execute(
            "UPDATE inventory SET stock = stock - %s "
            "WHERE sku_id = %s AND stock >= %s",
            (quantity, sku_id, quantity)
        )
        if affected == 0:
            raise ValueError("库存不足")
        
        db.execute(
            "UPDATE saga_branch_log SET status = 'SUCCESS' "
            "WHERE saga_id = %s AND branch_id = %s AND operation = 'deduct'",
            (saga_id, branch_id)
        )
        return True
    except Exception as e:
        db.execute(
            "UPDATE saga_branch_log SET status = 'FAILED' "
            "WHERE saga_id = %s AND branch_id = %s AND operation = 'deduct'",
            (saga_id, branch_id)
        )
        raise e
2. 空补偿:正向请求没到,补偿先到了

踩坑现场:有一次下游库存系统网关拥堵,正向请求堵在队列里没进业务逻辑,我们这边超时触发回滚,补偿请求走了另一条低延迟链路先到了下游。业务代码查不到扣减记录直接抛错,整个回滚流程直接卡住。

根因:网络乱序+超时回滚机制叠加,补偿请求在时序上先于正向请求到达下游。

解决方案:补偿操作执行前先校验正向记录是否存在。如果正向记录不存在,直接记录“空补偿”标记后返回成功,不执行业务逻辑——反正后续正向请求到了也会被防悬挂逻辑拦住,不会有数据不一致。

代码语言:python
复制
def compensate_inventory(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存补偿接口(空补偿防御)"""
    # 先查正向操作有没有执行过
    forward_status = db.query_status(saga_id, branch_id, 'deduct')
    
    if forward_status is None:
        # 正向没执行 → 空补偿,记个日志直接返回成功
        db.execute(
            "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
            "VALUES (%s, %s, 'compensate', 'EMPTY_COMPENSATED')",
            (saga_id, branch_id)
        )
        return True
    
    if forward_status == 'FAILED':
        # 正向本身就失败了,不用补偿
        return True
    
    # 补偿本身也要做幂等
    comp_status = db.query_status(saga_id, branch_id, 'compensate')
    if comp_status in ('SUCCESS', 'EMPTY_COMPENSATED'):
        return True
    
    # 执行实际补偿逻辑
    db.execute(
        "UPDATE inventory SET stock = stock + %s WHERE sku_id = %s",
        (quantity, sku_id)
    )
    db.execute(
        "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
        "VALUES (%s, %s, 'compensate', 'SUCCESS')",
        (saga_id, branch_id)
    )
    return True
3. 悬挂:补偿都执行完了,正向请求才姗姗来迟

踩坑现场:空补偿的问题刚修完一周,又出了反向问题:补偿已经标记完空补偿返回了,堵了半分钟的正向请求终于到了下游,直接执行了库存扣减。此时整个SAGA事务已经结束,没人再触发补偿,平白无故少了库存。

根因:网络延迟导致正向请求严重滞后,空补偿标记写入后,正向请求才真正进入业务逻辑。

解决方案:正向操作执行前加一道前置校验,只要该分支已经有补偿记录(不管是正常补偿还是空补偿),直接拒绝执行正向操作。

代码语言:python
复制
def deduct_inventory_with_suspend_protect(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存扣减接口(防悬挂增强版)"""
    # 防悬挂前置校验:补偿已经执行过,正向直接拒绝
    comp_status = db.query_status(saga_id, branch_id, 'compensate')
    if comp_status in ('SUCCESS', 'EMPTY_COMPENSATED'):
        return False
    
    # 后续走正常的幂等+业务逻辑
    return _do_deduct(saga_id, branch_id, sku_id, quantity)

这三道防线合起来就是业内说的“子事务屏障”,看着都是很简单的判断,但是每一条都是线上踩坑踩出来的。少了任何一个,只要网络出点波动,迟早要出数据问题。

二、用状态机做流程编排,别写满屏if-else

最开始我们的SAGA实现非常朴素:一个for循环挨个调正向接口,失败了就倒着for循环调补偿。当时只有三步流程,看着也清晰。后来业务加了两个步骤,还要支持服务重启后断点续跑,代码里的if-else堆得没人敢改,出一次问题定位半天。

后来我们把整个流程重构成了状态机驱动,本质就是把SAGA的执行过程抽象成有限状态流转——每个全局事务有明确的状态,每个分支步骤也有独立状态,所有流程推进都靠状态驱动,逻辑瞬间就清晰了。

状态机最大的好处有两个:一是流程可视化,每个步骤的前后依赖、异常分支一目了然;二是可重入性,状态持久化到数据库之后,服务不管什么时候重启,扫表读出状态就能接着往下跑,不用自己写一堆断点恢复逻辑。

状态模型定义

我们把全局SAGA事务分成7种状态,每个分支步骤单独维护状态,流转规则全部收敛在状态机里,不允许业务代码随意改状态。

代码语言:python
复制
from dataclasses import dataclass, field
from enum import Enum
from typing import List, Dict, Callable, Optional
import uuid
import time

class SagaStatus(str, Enum):
    INIT = "INIT"
    EXECUTING = "EXECUTING"
    SUCCEEDED = "SUCCEEDED"
    FAILED = "FAILED"
    COMPENSATING = "COMPENSATING"
    COMPENSATED = "COMPENSATED"
    COMPENSATE_FAILED = "COMPENSATE_FAILED"

class BranchStatus(str, Enum):
    PENDING = "PENDING"
    SUCCESS = "SUCCESS"
    FAILED = "FAILED"
    COMPENSATED = "COMPENSATED"

@dataclass
class SagaBranch:
    name: str
    action: Callable          # 正向操作
    compensate: Callable      # 补偿操作
    action_params: Dict = field(default_factory=dict)
    compensate_params: Dict = field(default_factory=dict)
    status: BranchStatus = BranchStatus.PENDING
    retry_count: int = 0
    error_msg: str = ""

@dataclass
class SagaInstance:
    saga_id: str
    status: SagaStatus = SagaStatus.INIT
    branches: List[SagaBranch] = field(default_factory=list)
    current_index: int = 0    # 当前执行到的步骤下标
    created_at: float = field(default_factory=time.time)
    updated_at: float = field(default_factory=time.time)

状态机核心实现

核心就是一个advance()方法,不管什么时候调用,都会根据当前状态自动推进:正向流就往下走,失败了就切补偿流往回走,到终态就停止。

正向重试我们设的3次,补偿重试设的5次——补偿失败的影响比正向失败大得多,多给几次重试机会,实在不行再进入人工兜底状态。

代码语言:python
复制
class SagaStateMachine:
    def __init__(self, max_retry: int = 3, compensate_retry: int = 5):
        self.max_retry = max_retry
        self.compensate_retry = compensate_retry
        self._instances: Dict[str, SagaInstance] = {}
    
    def create_saga(self, branches: List[SagaBranch]) -> str:
        saga_id = str(uuid.uuid4())
        instance = SagaInstance(
            saga_id=saga_id,
            status=SagaStatus.EXECUTING,
            branches=branches,
            current_index=0
        )
        self._instances[saga_id] = instance
        return saga_id
    
    def advance(self, saga_id: str) -> SagaStatus:
        """推进状态机,可重复调用,断点续跑"""
        inst = self._instances.get(saga_id)
        if not inst:
            raise ValueError(f"Saga {saga_id} not found")
        
        if inst.status in (SagaStatus.SUCCEEDED, SagaStatus.COMPENSATED, SagaStatus.COMPENSATE_FAILED):
            return inst.status
        
        if inst.status == SagaStatus.EXECUTING:
            self._execute_forward(inst)
        elif inst.status in (SagaStatus.FAILED, SagaStatus.COMPENSATING):
            self._execute_backward(inst)
        
        inst.updated_at = time.time()
        return inst.status
    
    def _execute_forward(self, inst: SagaInstance):
        while inst.current_index < len(inst.branches):
            branch = inst.branches[inst.current_index]
            
            if branch.status == BranchStatus.PENDING:
                try:
                    # 透传事务标识,下游做幂等屏障
                    branch.action_params["saga_id"] = inst.saga_id
                    branch.action_params["branch_id"] = f"branch_{inst.current_index}"
                    
                    result = branch.action(**branch.action_params)
                    
                    if result:
                        branch.status = BranchStatus.SUCCESS
                        inst.current_index += 1
                    else:
                        branch.status = BranchStatus.FAILED
                        branch.error_msg = "业务返回失败"
                        inst.status = SagaStatus.FAILED
                        break
                        
                except Exception as e:
                    branch.retry_count += 1
                    branch.error_msg = str(e)
                    
                    if branch.retry_count >= self.max_retry:
                        branch.status = BranchStatus.FAILED
                        inst.status = SagaStatus.FAILED
                        break
                    # 没到重试上限就退出,下次advance再试
                    return
            else:
                inst.current_index += 1
        
        if inst.current_index >= len(inst.branches) and inst.status == SagaStatus.EXECUTING:
            inst.status = SagaStatus.SUCCEEDED

补偿流的逻辑就是从失败位置的前一步开始,逆序执行补偿操作。这里要注意:只对执行成功的分支做补偿,本身就失败或者没执行的分支跳过。

代码语言:python
复制
    def _execute_backward(self, inst: SagaInstance):
        if inst.status == SagaStatus.FAILED:
            inst.status = SagaStatus.COMPENSATING
            # 从失败步的前一步开始回滚
            inst.current_index -= 1
        
        while inst.current_index >= 0:
            branch = inst.branches[inst.current_index]
            
            if branch.status == BranchStatus.SUCCESS:
                try:
                    branch.compensate_params["saga_id"] = inst.saga_id
                    branch.compensate_params["branch_id"] = f"branch_{inst.current_index}"
                    
                    branch.compensate(**branch.compensate_params)
                    branch.status = BranchStatus.COMPENSATED
                    inst.current_index -= 1
                    
                except Exception as e:
                    branch.retry_count += 1
                    branch.error_msg = f"补偿失败: {str(e)}"
                    
                    if branch.retry_count >= self.compensate_retry + self.max_retry:
                        inst.status = SagaStatus.COMPENSATE_FAILED
                        # 进入终态,触发告警人工介入
                        return
                    return
            else:
                inst.current_index -= 1
        
        if inst.current_index < 0:
            inst.status = SagaStatus.COMPENSATED

这套状态机我们用了快半年,最大的感受就是省心。不管是服务重启、定时任务恢复、还是手动触发重试,调同一个advance()方法就行,不用再写各种分支判断逻辑。

三、编排式SAGA在跨系统API场景的工程落地

SAGA有两种主流实现:编排式和编舞式。很多文章会花大篇幅对比两者优劣,但放到我们跨系统API调用的场景里,选型根本不用纠结。

编舞式靠事件总线驱动,每个服务监听事件自己执行逻辑,去中心化。但跨系统对接的现实是:外部合作方不可能为了你接一套事件总线,人家只提供HTTP API,所有流程只能由调用方主动触发。所以编排式是唯一可行的方案——做一个中央协调器,统一指挥每个步骤的调用、失败回滚、状态持久化,责任边界也清晰。

编排器完整实现

我们在状态机的基础上封装了一层编排器,专门处理跨系统HTTP调用场景,统一做超时控制、异常封装、持久化和告警。

代码语言:python
复制
import requests
from typing import List, Dict, Any
import logging

logger = logging.getLogger(__name__)

class SagaOrchestrator:
    """
    SAGA编排器(编排式)
    负责跨系统API调用的事务协调
    """
    
    def __init__(self, storage=None, alert_callback=None):
        self.state_machine = SagaStateMachine()
        self.storage = storage          # 持久化接口,生产环境必须实现
        self.alert_callback = alert_callback  # 告警回调,补偿失败触发
    
    def execute(self, branches_def: List[Dict[str, Any]]) -> Dict[str, Any]:
        branches = []
        for idx, bdef in enumerate(branches_def):
            timeout = bdef.get("timeout", 5)
            action_url = bdef["action_url"]
            compensate_url = bdef["compensate_url"]
            
            # 封装HTTP调用为分支操作
            branch = SagaBranch(
                name=bdef["name"],
                action=lambda **kw: self._call_api(action_url, kw, timeout),
                compensate=lambda **kw: self._call_api(compensate_url, kw, timeout),
                action_params=bdef.get("params", {}),
                compensate_params=bdef.get("compensate_params", {})
            )
            branches.append(branch)
        
        saga_id = self.state_machine.create_saga(branches)
        logger.info(f"启动SAGA事务: {saga_id}")
        
        # 驱动状态机直到终态
        while True:
            status = self.state_machine.advance(saga_id)
            
            # 每次推进后持久化状态,宕机可恢复
            if self.storage:
                self.storage.save(self.state_machine._instances[saga_id])
            
            if status in (SagaStatus.SUCCEEDED, SagaStatus.COMPENSATED):
                break
            
            if status == SagaStatus.COMPENSATE_FAILED:
                logger.error(f"SAGA补偿失败,需人工介入: {saga_id}")
                if self.alert_callback:
                    self.alert_callback(saga_id, "COMPENSATE_FAILED")
                break
        
        inst = self.state_machine._instances[saga_id]
        return {
            "saga_id": saga_id,
            "status": status.value,
            "executed_steps": sum(1 for b in inst.branches if b.status != BranchStatus.PENDING),
            "compensated_steps": sum(1 for b in inst.branches if b.status == BranchStatus.COMPENSATED)
        }
    
    def _call_api(self, url: str, params: Dict, timeout: int) -> bool:
        """统一封装HTTP调用,异常标准化"""
        try:
            resp = requests.post(
                url,
                json=params,
                timeout=timeout,
                headers={"Content-Type": "application/json"}
            )
            resp.raise_for_status()
            data = resp.json()
            return data.get("success", False)
        except requests.Timeout:
            raise RuntimeError(f"接口超时: {url}")
        except requests.RequestException as e:
            raise RuntimeError(f"接口调用失败: {url}, 错误: {str(e)}")

使用示例

以典型的订单-库存-支付三系统调用为例,使用方式非常简洁:

代码语言:python
复制
if __name__ == "__main__":
    orchestrator = SagaOrchestrator()
    
    saga_flow = [
        {
            "name": "create_order",
            "action_url": "http://order-service/api/orders/create",
            "compensate_url": "http://order-service/api/orders/cancel",
            "params": {"order_no": "ORD20260720001", "amount": 999.00},
            "compensate_params": {"order_no": "ORD20260720001"},
            "timeout": 3
        },
        {
            "name": "deduct_inventory",
            "action_url": "http://inventory-service/api/stock/deduct",
            "compensate_url": "http://inventory-service/api/stock/release",
            "params": {"sku_id": "SKU001", "quantity": 2},
            "compensate_params": {"sku_id": "SKU001", "quantity": 2},
            "timeout": 3
        },
        {
            "name": "process_payment",
            "action_url": "http://payment-service/api/pay/charge",
            "compensate_url": "http://payment-service/api/pay/refund",
            "params": {"amount": 999.00, "channel": "alipay"},
            "compensate_params": {"amount": 999.00},
            "timeout": 8  # 支付接口超时阈值设高一点
        }
    ]
    
    result = orchestrator.execute(saga_flow)
    print(f"事务执行结果: {result}")

生产落地的几个补充建议

上面的代码是核心骨架,真要上线还有几个必须补齐的基础设施:

  1. 持久化+定时恢复:别把状态放内存里,必须落库。配合定时任务扫描超时、中断的事务,重新调用advance续跑。我们用MySQL存事务状态,APScheduler做分钟级扫描。
  2. 异步化改造:超过3步的长链路别同步阻塞等待。调用方发起事务后立即返回saga_id,后续通过回调或者主动查询获取结果,避免线程被长时间占用。
  3. 人工兜底入口:别迷信补偿一定成功。第三方通道维护、账户冻结、退款限额这些场景,补偿必然失败。必须做一个后台页面,展示补偿失败的事务,支持人工处理。
  4. 隔离性应对:SAGA没有隔离性,中间状态对其他事务可见。我们的做法是把最容易失败的步骤放最前面,减少补偿概率;同时资源用PENDING状态标记,避免中间状态被其他事务误读。
  5. 监控埋点:SAGA的成功率、补偿率、平均执行耗时都要进监控。补偿率突然升高,基本就是下游系统出问题了,能提前发现很多故障。

最后

SAGA这个方案,原理看着简单,落地的难度全在细节里。很多团队用不好SAGA,不是模式本身有问题,是只实现了正常流程,异常场景全漏了。幂等、空补偿、悬挂这三道防线,加上状态机做流程编排,基本就能覆盖99%的生产场景。

跨系统API调用这种场景,强一致是不现实的,最终一致是性价比最高的选择。SAGA不是银弹,但只要把异常处理做足,稳定性完全能满足业务需求。# 跨系统API调用落地SAGA模式:补偿机制、异常防御与流程编排实战总结

最近在做跨系统采购链路的对接,整条链路涉及订单、库存、支付三套独立的外部系统,全部通过HTTP API交互。最开始图省事,按顺序串行调用接口,用本地事务包一层,结果上线没多久就出了两次问题:中间某一步超时失败,前面已经执行成功的操作没法回滚,要么库存多扣要么资金单边账,资损风险极高。

调研了一圈分布式事务方案,2PC/XA对跨系统场景太重,强一致的阻塞开销我们也接受不了,最终选了SAGA柔性事务。落地的过程踩了不少坑——网上大部分资料只讲正向流程+反向补偿的Happy Path,真放到生产环境,网络乱序、超时重试、服务宕机这些问题全出来了。这里把核心的三块实践经验整理出来,附可直接复用的Python实现。

一、补偿的核心不是“反向调用”,是业务语义撤销

很多人刚接触SAGA的时候,会把补偿简单理解成“调反向接口”:扣了库存就加回去,付了款就退回去。这是最常见的误区,也是线上事故的高发点。

SAGA的本质是把一个长事务拆成N个独立的本地事务,每一步执行完就立即提交、释放连接资源。当第K步失败时,按逆序对前面K-1步执行补偿,最终让系统回到事务开始前的等价状态。注意这里说的是等价状态,不是“数据完全复原”——本地事务一旦提交就没法rollback,补偿是一个全新的业务操作,用反向语义抵消正向操作的业务影响。

举个最常见的下单链路例子:

  • 创建订单的补偿不是删除订单,是把订单状态置为已取消,保留单据留痕
  • 扣减库存的补偿不是直接加库存,是归还本次事务占用的对应数量
  • 支付扣款的补偿不是撤销扣款记录,是发起一笔原路退款

我们落地时给补偿操作定了三条硬标准:必须幂等、语义可逆、执行后不可再正向重入。这三条没守住,补偿就是新的事故源。

生产必踩的三个异常坑与防御方案

跨系统API调用天然存在网络延迟、丢包、乱序,只写“正向+补偿”两层逻辑,线上必出问题。我们前前后后踩了三遍坑,才把这三道防御补全。

1. 幂等性缺失:超时重试导致重复执行

踩坑现场:第一次压测模拟网络抖动,库存接口超时3秒,调用端按默认策略重试了2次,最终库存被扣了3次。

根因:HTTP超时场景下,调用方永远无法确定下游到底有没有执行业务逻辑。框架层的重试机制一触发,重复扣减、重复退款是必然结果。

解决方案:全链路透传saga_id + branch_id,下游基于唯一键做幂等控制。用数据库唯一键冲突判重,比先查后插性能更好,也能避免并发场景下的判断失效。

代码语言:python
复制
def deduct_inventory(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存扣减接口(幂等实现)"""
    unique_key = f"{saga_id}:{branch_id}:deduct"
    
    # 先插幂等记录,靠唯一键挡重复请求
    try:
        db.execute(
            "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
            "VALUES (%s, %s, 'deduct', 'PROCESSING')",
            (saga_id, branch_id)
        )
    except IntegrityError:
        # 重复请求直接返回已有状态,不执行业务
        status = db.query_status(saga_id, branch_id, 'deduct')
        return status == 'SUCCESS'
    
    # 执行业务逻辑
    try:
        affected = db.execute(
            "UPDATE inventory SET stock = stock - %s "
            "WHERE sku_id = %s AND stock >= %s",
            (quantity, sku_id, quantity)
        )
        if affected == 0:
            raise ValueError("库存不足")
        
        db.execute(
            "UPDATE saga_branch_log SET status = 'SUCCESS' "
            "WHERE saga_id = %s AND branch_id = %s AND operation = 'deduct'",
            (saga_id, branch_id)
        )
        return True
    except Exception as e:
        db.execute(
            "UPDATE saga_branch_log SET status = 'FAILED' "
            "WHERE saga_id = %s AND branch_id = %s AND operation = 'deduct'",
            (saga_id, branch_id)
        )
        raise e
2. 空补偿:正向请求没到,补偿先到了

踩坑现场:有一次下游库存系统网关拥堵,正向请求堵在队列里没进业务逻辑,我们这边超时触发回滚,补偿请求走了另一条低延迟链路先到了下游。业务代码查不到扣减记录直接抛错,整个回滚流程直接卡住。

根因:网络乱序+超时回滚机制叠加,补偿请求在时序上先于正向请求到达下游。

解决方案:补偿操作执行前先校验正向记录是否存在。如果正向记录不存在,直接记录“空补偿”标记后返回成功,不执行业务逻辑——反正后续正向请求到了也会被防悬挂逻辑拦住,不会有数据不一致。

代码语言:python
复制
def compensate_inventory(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存补偿接口(空补偿防御)"""
    # 先查正向操作有没有执行过
    forward_status = db.query_status(saga_id, branch_id, 'deduct')
    
    if forward_status is None:
        # 正向没执行 → 空补偿,记个日志直接返回成功
        db.execute(
            "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
            "VALUES (%s, %s, 'compensate', 'EMPTY_COMPENSATED')",
            (saga_id, branch_id)
        )
        return True
    
    if forward_status == 'FAILED':
        # 正向本身就失败了,不用补偿
        return True
    
    # 补偿本身也要做幂等
    comp_status = db.query_status(saga_id, branch_id, 'compensate')
    if comp_status in ('SUCCESS', 'EMPTY_COMPENSATED'):
        return True
    
    # 执行实际补偿逻辑
    db.execute(
        "UPDATE inventory SET stock = stock + %s WHERE sku_id = %s",
        (quantity, sku_id)
    )
    db.execute(
        "INSERT INTO saga_branch_log (saga_id, branch_id, operation, status) "
        "VALUES (%s, %s, 'compensate', 'SUCCESS')",
        (saga_id, branch_id)
    )
    return True
3. 悬挂:补偿都执行完了,正向请求才姗姗来迟

踩坑现场:空补偿的问题刚修完一周,又出了反向问题:补偿已经标记完空补偿返回了,堵了半分钟的正向请求终于到了下游,直接执行了库存扣减。此时整个SAGA事务已经结束,没人再触发补偿,平白无故少了库存。

根因:网络延迟导致正向请求严重滞后,空补偿标记写入后,正向请求才真正进入业务逻辑。

解决方案:正向操作执行前加一道前置校验,只要该分支已经有补偿记录(不管是正常补偿还是空补偿),直接拒绝执行正向操作。

代码语言:python
复制
def deduct_inventory_with_suspend_protect(saga_id: str, branch_id: str, sku_id: str, quantity: int) -> bool:
    """库存扣减接口(防悬挂增强版)"""
    # 防悬挂前置校验:补偿已经执行过,正向直接拒绝
    comp_status = db.query_status(saga_id, branch_id, 'compensate')
    if comp_status in ('SUCCESS', 'EMPTY_COMPENSATED'):
        return False
    
    # 后续走正常的幂等+业务逻辑
    return _do_deduct(saga_id, branch_id, sku_id, quantity)

这三道防线合起来就是业内说的“子事务屏障”,看着都是很简单的判断,但是每一条都是线上踩坑踩出来的。少了任何一个,只要网络出点波动,迟早要出数据问题。

二、用状态机做流程编排,别写满屏if-else

最开始我们的SAGA实现非常朴素:一个for循环挨个调正向接口,失败了就倒着for循环调补偿。当时只有三步流程,看着也清晰。后来业务加了两个步骤,还要支持服务重启后断点续跑,代码里的if-else堆得没人敢改,出一次问题定位半天。

后来我们把整个流程重构成了状态机驱动,本质就是把SAGA的执行过程抽象成有限状态流转——每个全局事务有明确的状态,每个分支步骤也有独立状态,所有流程推进都靠状态驱动,逻辑瞬间就清晰了。

状态机最大的好处有两个:一是流程可视化,每个步骤的前后依赖、异常分支一目了然;二是可重入性,状态持久化到数据库之后,服务不管什么时候重启,扫表读出状态就能接着往下跑,不用自己写一堆断点恢复逻辑。

状态模型定义

我们把全局SAGA事务分成7种状态,每个分支步骤单独维护状态,流转规则全部收敛在状态机里,不允许业务代码随意改状态。

代码语言:python
复制
from dataclasses import dataclass, field
from enum import Enum
from typing import List, Dict, Callable, Optional
import uuid
import time

class SagaStatus(str, Enum):
    INIT = "INIT"
    EXECUTING = "EXECUTING"
    SUCCEEDED = "SUCCEEDED"
    FAILED = "FAILED"
    COMPENSATING = "COMPENSATING"
    COMPENSATED = "COMPENSATED"
    COMPENSATE_FAILED = "COMPENSATE_FAILED"

class BranchStatus(str, Enum):
    PENDING = "PENDING"
    SUCCESS = "SUCCESS"
    FAILED = "FAILED"
    COMPENSATED = "COMPENSATED"

@dataclass
class SagaBranch:
    name: str
    action: Callable          # 正向操作
    compensate: Callable      # 补偿操作
    action_params: Dict = field(default_factory=dict)
    compensate_params: Dict = field(default_factory=dict)
    status: BranchStatus = BranchStatus.PENDING
    retry_count: int = 0
    error_msg: str = ""

@dataclass
class SagaInstance:
    saga_id: str
    status: SagaStatus = SagaStatus.INIT
    branches: List[SagaBranch] = field(default_factory=list)
    current_index: int = 0    # 当前执行到的步骤下标
    created_at: float = field(default_factory=time.time)
    updated_at: float = field(default_factory=time.time)

状态机核心实现

核心就是一个advance()方法,不管什么时候调用,都会根据当前状态自动推进:正向流就往下走,失败了就切补偿流往回走,到终态就停止。

正向重试我们设的3次,补偿重试设的5次——补偿失败的影响比正向失败大得多,多给几次重试机会,实在不行再进入人工兜底状态。

代码语言:python
复制
class SagaStateMachine:
    def __init__(self, max_retry: int = 3, compensate_retry: int = 5):
        self.max_retry = max_retry
        self.compensate_retry = compensate_retry
        self._instances: Dict[str, SagaInstance] = {}
    
    def create_saga(self, branches: List[SagaBranch]) -> str:
        saga_id = str(uuid.uuid4())
        instance = SagaInstance(
            saga_id=saga_id,
            status=SagaStatus.EXECUTING,
            branches=branches,
            current_index=0
        )
        self._instances[saga_id] = instance
        return saga_id
    
    def advance(self, saga_id: str) -> SagaStatus:
        """推进状态机,可重复调用,断点续跑"""
        inst = self._instances.get(saga_id)
        if not inst:
            raise ValueError(f"Saga {saga_id} not found")
        
        if inst.status in (SagaStatus.SUCCEEDED, SagaStatus.COMPENSATED, SagaStatus.COMPENSATE_FAILED):
            return inst.status
        
        if inst.status == SagaStatus.EXECUTING:
            self._execute_forward(inst)
        elif inst.status in (SagaStatus.FAILED, SagaStatus.COMPENSATING):
            self._execute_backward(inst)
        
        inst.updated_at = time.time()
        return inst.status
    
    def _execute_forward(self, inst: SagaInstance):
        while inst.current_index < len(inst.branches):
            branch = inst.branches[inst.current_index]
            
            if branch.status == BranchStatus.PENDING:
                try:
                    # 透传事务标识,下游做幂等屏障
                    branch.action_params["saga_id"] = inst.saga_id
                    branch.action_params["branch_id"] = f"branch_{inst.current_index}"
                    
                    result = branch.action(**branch.action_params)
                    
                    if result:
                        branch.status = BranchStatus.SUCCESS
                        inst.current_index += 1
                    else:
                        branch.status = BranchStatus.FAILED
                        branch.error_msg = "业务返回失败"
                        inst.status = SagaStatus.FAILED
                        break
                        
                except Exception as e:
                    branch.retry_count += 1
                    branch.error_msg = str(e)
                    
                    if branch.retry_count >= self.max_retry:
                        branch.status = BranchStatus.FAILED
                        inst.status = SagaStatus.FAILED
                        break
                    # 没到重试上限就退出,下次advance再试
                    return
            else:
                inst.current_index += 1
        
        if inst.current_index >= len(inst.branches) and inst.status == SagaStatus.EXECUTING:
            inst.status = SagaStatus.SUCCEEDED

补偿流的逻辑就是从失败位置的前一步开始,逆序执行补偿操作。这里要注意:只对执行成功的分支做补偿,本身就失败或者没执行的分支跳过。

代码语言:python
复制
    def _execute_backward(self, inst: SagaInstance):
        if inst.status == SagaStatus.FAILED:
            inst.status = SagaStatus.COMPENSATING
            # 从失败步的前一步开始回滚
            inst.current_index -= 1
        
        while inst.current_index >= 0:
            branch = inst.branches[inst.current_index]
            
            if branch.status == BranchStatus.SUCCESS:
                try:
                    branch.compensate_params["saga_id"] = inst.saga_id
                    branch.compensate_params["branch_id"] = f"branch_{inst.current_index}"
                    
                    branch.compensate(**branch.compensate_params)
                    branch.status = BranchStatus.COMPENSATED
                    inst.current_index -= 1
                    
                except Exception as e:
                    branch.retry_count += 1
                    branch.error_msg = f"补偿失败: {str(e)}"
                    
                    if branch.retry_count >= self.compensate_retry + self.max_retry:
                        inst.status = SagaStatus.COMPENSATE_FAILED
                        # 进入终态,触发告警人工介入
                        return
                    return
            else:
                inst.current_index -= 1
        
        if inst.current_index < 0:
            inst.status = SagaStatus.COMPENSATED

这套状态机我们用了快半年,最大的感受就是省心。不管是服务重启、定时任务恢复、还是手动触发重试,调同一个advance()方法就行,不用再写各种分支判断逻辑。

三、编排式SAGA在跨系统API场景的工程落地

SAGA有两种主流实现:编排式和编舞式。很多文章会花大篇幅对比两者优劣,但放到我们跨系统API调用的场景里,选型根本不用纠结。

编舞式靠事件总线驱动,每个服务监听事件自己执行逻辑,去中心化。但跨系统对接的现实是:外部合作方不可能为了你接一套事件总线,人家只提供HTTP API,所有流程只能由调用方主动触发。所以编排式是唯一可行的方案——做一个中央协调器,统一指挥每个步骤的调用、失败回滚、状态持久化,责任边界也清晰。

编排器完整实现

我们在状态机的基础上封装了一层编排器,专门处理跨系统HTTP调用场景,统一做超时控制、异常封装、持久化和告警。

代码语言:python
复制
import requests
from typing import List, Dict, Any
import logging

logger = logging.getLogger(__name__)

class SagaOrchestrator:
    """
    SAGA编排器(编排式)
    负责跨系统API调用的事务协调
    """
    
    def __init__(self, storage=None, alert_callback=None):
        self.state_machine = SagaStateMachine()
        self.storage = storage          # 持久化接口,生产环境必须实现
        self.alert_callback = alert_callback  # 告警回调,补偿失败触发
    
    def execute(self, branches_def: List[Dict[str, Any]]) -> Dict[str, Any]:
        branches = []
        for idx, bdef in enumerate(branches_def):
            timeout = bdef.get("timeout", 5)
            action_url = bdef["action_url"]
            compensate_url = bdef["compensate_url"]
            
            # 封装HTTP调用为分支操作
            branch = SagaBranch(
                name=bdef["name"],
                action=lambda **kw: self._call_api(action_url, kw, timeout),
                compensate=lambda **kw: self._call_api(compensate_url, kw, timeout),
                action_params=bdef.get("params", {}),
                compensate_params=bdef.get("compensate_params", {})
            )
            branches.append(branch)
        
        saga_id = self.state_machine.create_saga(branches)
        logger.info(f"启动SAGA事务: {saga_id}")
        
        # 驱动状态机直到终态
        while True:
            status = self.state_machine.advance(saga_id)
            
            # 每次推进后持久化状态,宕机可恢复
            if self.storage:
                self.storage.save(self.state_machine._instances[saga_id])
            
            if status in (SagaStatus.SUCCEEDED, SagaStatus.COMPENSATED):
                break
            
            if status == SagaStatus.COMPENSATE_FAILED:
                logger.error(f"SAGA补偿失败,需人工介入: {saga_id}")
                if self.alert_callback:
                    self.alert_callback(saga_id, "COMPENSATE_FAILED")
                break
        
        inst = self.state_machine._instances[saga_id]
        return {
            "saga_id": saga_id,
            "status": status.value,
            "executed_steps": sum(1 for b in inst.branches if b.status != BranchStatus.PENDING),
            "compensated_steps": sum(1 for b in inst.branches if b.status == BranchStatus.COMPENSATED)
        }
    
    def _call_api(self, url: str, params: Dict, timeout: int) -> bool:
        """统一封装HTTP调用,异常标准化"""
        try:
            resp = requests.post(
                url,
                json=params,
                timeout=timeout,
                headers={"Content-Type": "application/json"}
            )
            resp.raise_for_status()
            data = resp.json()
            return data.get("success", False)
        except requests.Timeout:
            raise RuntimeError(f"接口超时: {url}")
        except requests.RequestException as e:
            raise RuntimeError(f"接口调用失败: {url}, 错误: {str(e)}")

使用示例

以典型的订单-库存-支付三系统调用为例,使用方式非常简洁:

代码语言:python
复制
if __name__ == "__main__":
    orchestrator = SagaOrchestrator()
    
    saga_flow = [
        {
            "name": "create_order",
            "action_url": "http://order-service/api/orders/create",
            "compensate_url": "http://order-service/api/orders/cancel",
            "params": {"order_no": "ORD20260720001", "amount": 999.00},
            "compensate_params": {"order_no": "ORD20260720001"},
            "timeout": 3
        },
        {
            "name": "deduct_inventory",
            "action_url": "http://inventory-service/api/stock/deduct",
            "compensate_url": "http://inventory-service/api/stock/release",
            "params": {"sku_id": "SKU001", "quantity": 2},
            "compensate_params": {"sku_id": "SKU001", "quantity": 2},
            "timeout": 3
        },
        {
            "name": "process_payment",
            "action_url": "http://payment-service/api/pay/charge",
            "compensate_url": "http://payment-service/api/pay/refund",
            "params": {"amount": 999.00, "channel": "alipay"},
            "compensate_params": {"amount": 999.00},
            "timeout": 8  # 支付接口超时阈值设高一点
        }
    ]
    
    result = orchestrator.execute(saga_flow)
    print(f"事务执行结果: {result}")

生产落地的几个补充建议

上面的代码是核心骨架,真要上线还有几个必须补齐的基础设施:

  1. 持久化+定时恢复:别把状态放内存里,必须落库。配合定时任务扫描超时、中断的事务,重新调用advance续跑。我们用MySQL存事务状态,APScheduler做分钟级扫描。
  2. 异步化改造:超过3步的长链路别同步阻塞等待。调用方发起事务后立即返回saga_id,后续通过回调或者主动查询获取结果,避免线程被长时间占用。
  3. 人工兜底入口:别迷信补偿一定成功。第三方通道维护、账户冻结、退款限额这些场景,补偿必然失败。必须做一个后台页面,展示补偿失败的事务,支持人工处理。
  4. 隔离性应对:SAGA没有隔离性,中间状态对其他事务可见。我们的做法是把最容易失败的步骤放最前面,减少补偿概率;同时资源用PENDING状态标记,避免中间状态被其他事务误读。
  5. 监控埋点:SAGA的成功率、补偿率、平均执行耗时都要进监控。补偿率突然升高,基本就是下游系统出问题了,能提前发现很多故障。

最后

SAGA这个方案,原理看着简单,落地的难度全在细节里。很多团队用不好SAGA,不是模式本身有问题,是只实现了正常流程,异常场景全漏了。幂等、空补偿、悬挂这三道防线,加上状态机做流程编排,基本就能覆盖99%的生产场景。

跨系统API调用这种场景,强一致是不现实的,最终一致是性价比最高的选择。SAGA不是银弹,但只要把异常处理做足,稳定性完全能满足业务需求。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

评论
登录后参与评论
0 条评论
热度
最新
推荐阅读
目录
  • 一、补偿的核心不是“反向调用”,是业务语义撤销
    • 生产必踩的三个异常坑与防御方案
      • 1. 幂等性缺失:超时重试导致重复执行
      • 2. 空补偿:正向请求没到,补偿先到了
      • 3. 悬挂:补偿都执行完了,正向请求才姗姗来迟
  • 二、用状态机做流程编排,别写满屏if-else
    • 状态模型定义
    • 状态机核心实现
  • 三、编排式SAGA在跨系统API场景的工程落地
    • 编排器完整实现
    • 使用示例
    • 生产落地的几个补充建议
  • 最后
  • 一、补偿的核心不是“反向调用”,是业务语义撤销
    • 生产必踩的三个异常坑与防御方案
      • 1. 幂等性缺失:超时重试导致重复执行
      • 2. 空补偿:正向请求没到,补偿先到了
      • 3. 悬挂:补偿都执行完了,正向请求才姗姗来迟
  • 二、用状态机做流程编排,别写满屏if-else
    • 状态模型定义
    • 状态机核心实现
  • 三、编排式SAGA在跨系统API场景的工程落地
    • 编排器完整实现
    • 使用示例
    • 生产落地的几个补充建议
  • 最后
领券
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档