分布式事务:基于 Saga 补偿的状态机实践
背景
上一篇"分布式事务:本地消息表"里,重点在"任务怎么可靠地重试"。但真实的下游链路远比"调一个下游、重试到成功"复杂。
以BNPL支用业务为例:一笔 BNPL 支用,从收单到最终成功,中间要串起多个下游动作:过风控、冻优惠券、冻额度、创建借据、签合同、扣减额度、确认借据、核销券、更新业务订单……每一步都是独立的远程调用,每一步都可能失败。
问题在于:这些步骤不是彼此独立的。如券冻结后,冻结额度节点失败了,券得解冻这种逆向补偿动作。 单纯的"重试到成功"解决不了这个问题——有些节点根本不该重试,而是应该往回走,把前面已经产生的副作用抵消掉。
这就是 Saga:把一个长事务拆成若干个带补偿动作的本地事务,正向一路往前推,某一步失败了,就沿着补偿链往回冲销。 不追求 ACID 的强一致,只保证最终一致——要么全部成功,要么把已经做的事情全部回滚干净。
一、整体思路
1.1 为什么是状态机
本地消息表回答的是"一个任务怎么可靠执行",Saga 要回答的是"一串有先后依赖、且需要补偿的任务怎么按顺序推进、失败怎么回退"。
把每个下游动作抽象成一个状态,状态与状态之间的跳转规则用一张转移表固定下来,整条链路就变成了一台状态机:
- 每个状态对应一个执行单元(
StateUnit),干一件具体的事——比如冻额度、建合同; - 每个状态执行完,根据成功 / 失败 / 处理中三种结果,查转移表决定下一步跳到哪;
- 正向链把业务一路推到成功终态;
- 某一步返回失败,状态机不是直接停,而是跳到该步预先配置好的补偿状态,沿补偿链往回走,直到失败终态。
所以这里的状态机有两条路:一条正向推进链,一条反向补偿链。它们不是两套代码,而是同一张转移表里,每个状态各自配了"成功去哪、失败去哪"。
1.2 三个结果驱动流转
整个状态机只认三种执行结果:
| 结果 | 含义 | 状态机动作 |
|---|---|---|
SUCCESS |
本节点业务成功 | 跳到该节点配置的成功状态,继续往前 |
FAIL |
明确的业务失败(如额度不足) | 跳到该节点配置的失败状态,进入补偿链 |
PROCESSING |
不确定 / 可重试(超时、未知异常) | 停在原地,等待下一次调度重试 |
这三个结果的边界是整套机制的关键,尤其是 FAIL 和 PROCESSING 的区分。
二、状态机框架
框架由五个角色构成,职责非常清晰:
TradeState —— 状态定义(id + 名字),一个 enum 常量就是一个状态
Transition —— 一条转移规则:当前状态 →(成功)去哪 /(失败)去哪 + 用哪个执行单元
TransitionManager —— 转移表容器,启动时把 Transition[] 处理成 Map,运行时 O(1) 查询
StateUnit —— 状态执行单元,干具体业务的地方
AbstractJump —— 一条业务线的状态机门面,聚合上面这些,对外提供"下一步是什么"
2.1 Transition:一条转移规则
Transition 把"一个状态该怎么走"标准化:当前状态、成功跳哪、失败跳哪、用哪个执行单元。
public interface Transition {
TradeState getCurrent(); // 当前状态
TradeState getSuccessState(); // 成功后跳转到
TradeState getFailState(); // 失败后跳转到
Class<? extends StateUnit> getStateUnit(); // 该状态用哪个执行单元
}
一条 Transition 就是转移表里的一行——这个状态干什么、成功往哪走、失败往哪走,三件事绑死在一起。补偿链的走向,就藏在每个状态的 getFailState() 里。
2.2 TransitionManager:把转移表编译成 Map
启动时把 Transition[] 数组处理翻译成查询表,运行时 O(1) 拿"执行单元 / 下一状态"。核心是这段逻辑:
@Override
public void setApplicationContext(ApplicationContext ctx) {
for (Transition t : transitions) {
StateUnit unit = ctx.getBean(t.getStateUnit()); // Class 引用 → 容器 Bean
Map<Boolean, Integer> nextMap = new HashMap<>();
nextMap.put(true, t.getSuccessState().getId()); // 成功态
nextMap.put(false, t.getFailState().getId()); // 失败态
transitionMap.put(t.getCurrent().getId(), new Pair<>(unit, nextMap));
}
}
每条业务编排一个 TransitionManager 实例,用配置类显式声明。 支用、退款等各自一张独立的转移表,互不干扰。新增一条业务编排,就是加一个 enum + 加一个 Bean,核心调度代码一行都不用动。
2.3 AbstractJump:状态机门面
AbstractJump 对外只暴露三个问题的答案:从哪开始、下一步去哪、现在是不是终态。 getNextTradeStatusByResult——把三种结果映射成三种走向:
public int getNextTradeStatusByResult(TransactionEventDTO event, MachineStateEnum result) {
int state = event.getTradeState();
if (result == MachineStateEnum.FAIL) {
return getTransitionManager().getFail(state); // 失败 → 补偿链
} else if (result == MachineStateEnum.SUCCESS) {
return getTransitionManager().getSuccess(state); // 成功 → 正向链
} else if (result == MachineStateEnum.PROCESSING) {
return state; // 处理中 → 原地不动
}
throw new RuntimeException("get next error");
}
补偿的语义完全落在 getFail(state) 上——它返回的下一个状态,可能就是一个补偿动作。
2.4 具体业务线:PayJump
PayJump 是支用状态机。它做两件事:定义所有状态、定义转移表。转移表节选,重点看正向和补偿两条链是怎么在一张表里同时表达的:
public enum PayJumpTransition implements Transition {
// 当前状态 成功去 → 失败去 → 执行单元
FREEZE_COUPON(FREEZE_COUPON, FREEZE_LIMIT, UPDATE_TRADE_FAIL, FreezeCouponUnit.class),
FREEZE_LIMIT (FREEZE_LIMIT, CREATE_LOAN, UNFREEZE_COUPON, FreezeLimitUnit.class),
// …… 正向链一路到 终态
// 补偿链
UNFREEZE_COUPON(UNFREEZE_COUPON, UPDATE_TRADE_FAIL, UNFREEZE_COUPON, UnfreezeCouponUnit.class),
UPDATE_TRADE_FAIL(UPDATE_TRADE_FAIL, SEND_FAIL_MQ, UPDATE_TRADE_FAIL, UpdateTradeFailUnit.class);
}
看 FREEZE_LIMIT 这条:冻额度成功去建借据(正向),失败去 UNFREEZE_COUPON(解冻优惠券)——额度冻不上,得先把上一步冻的券退回来,再走失败收尾。这就是 Saga 补偿的精髓:每个可能产生副作用的正向节点,都在 getFailState() 里指向它的补偿路径。
2.5 状态流转图
把 PayJump 转移表画出来,正向推进链和补偿回退链一目了然:
stateDiagram-v2
direction TB
[*] --> 支用风控
支用风控 --> 授信开通状态: 成功
授信开通状态 --> 冻结优惠券: 成功
冻结优惠券 --> 冻结额度: 成功
冻结额度 --> 创建借据: 成功
创建借据 --> 创建合同: 成功
创建合同 --> 占用额度: 成功
占用额度 --> 更新借据成功: 成功
更新借据成功 --> 优惠券核销: 成功
优惠券核销 --> 更新交易单成功: 成功
更新交易单成功 --> 发支用成功MQ: 成功
发支用成功MQ --> 查询合同签署: 成功
查询合同签署 --> FINISH: 成功
FINISH --> [*]
支用风控 --> 更新交易单失败: 失败
授信开通状态 --> 更新交易单失败: 失败
冻结优惠券 --> 更新交易单失败: 失败
冻结额度 --> 解冻优惠券: 失败(补偿)
解冻优惠券 --> 更新交易单失败: 成功
更新交易单失败 --> 发支用失败MQ: 成功
发支用失败MQ --> GIVEUP: 成功
GIVEUP --> [*]
note right of 冻结额度
1. SUCCESS 走正向链
2. FAIL 走补偿链
3. PROCESSING 原地重试
end note
说明:正向链上任一节点返回
FAIL,都会跳到它配置的失败态。风控/查授信/冻券这些尚未产生需回滚副作用的节点,失败直接去"更新交易单失败"收尾;而冻额度这种已经冻了券的节点,失败必须先经过"解冻优惠券"补偿,再收尾。
三、ER / 存储结构
Saga 状态机任务表,一个主单表 main_transaction和一个日志表main_transaction_log。
CREATE TABLE `main_transaction` (
`id` bigint NOT NULL AUTO_INCREMENT COMMENT '主键ID',
`instruction_id` bigint NOT NULL COMMENT '业务指令ID,一笔支用/退款的唯一标识',
`trade_code` int NOT NULL COMMENT '交易类型:支用/退款/预冻结支用,路由到不同状态机',
`trade_state` int NOT NULL COMMENT '当前状态机状态(TradeState.id),Saga 的游标',
`fail_state` int NOT NULL DEFAULT 0 COMMENT '首次失败发生在哪个状态,用于排查/补偿定位',
`run_times` int NOT NULL DEFAULT 0 COMMENT '状态机累计执行次数,用于限制重试上限',
`product_id` varchar(64) NOT NULL COMMENT '产品ID',
`context` text COMMENT '原始请求参数(静态)',
`biz_data` text COMMENT '状态机执行过程中产生的业务数据',
`create_time` datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
`update_time` datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3),
PRIMARY KEY (`id`),
UNIQUE KEY `uk_instruction` (`instruction_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='主交易单';
配套一张状态流转日志表 main_transaction_log,每跨一个状态就记一条,把整条 Saga 的推进 / 回退轨迹完整落盘:
CREATE TABLE `main_transaction_log` (
`id` bigint NOT NULL AUTO_INCREMENT,
`instruction_id` bigint NOT NULL COMMENT '业务指令ID',
`trade_code` int NOT NULL COMMENT '交易类型',
`before_trade_state` int NOT NULL COMMENT '跳转前状态',
`after_trade_state` int NOT NULL COMMENT '跳转后状态',
`trace_id` varchar(64) COMMENT '链路追踪ID',
`operator` varchar(64) COMMENT '操作人',
`create_time` datetime(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3),
PRIMARY KEY (`id`),
KEY `idx_instruction` (`instruction_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='状态流转日志';
字段设计直接对应 Saga 的需求:
trade_state就是 Saga 的执行游标。 每推进一步就更新它,进程随时崩溃、重启,拿instruction_id把主单捞出来读trade_state就知道从哪接着跑——状态天然持久化,不依赖内存。fail_state记录首次失败点。 补偿链一旦启动,trade_state会变成补偿态,原始失败点就丢了;单独存它是为了事后一眼看出"最初卡在哪个正向节点"。run_times卡执行次数上限。 防止PROCESSING的节点无限重试。biz_data存过程数据。 正向节点产生的、补偿节点要用的中间结果存这里,补偿时读出来用于冲销。
四、执行引擎:machineExecute
状态机跑起来的引擎:一个循环 + 乐观锁推进的结构:一次调用尽可能把状态机往前推到底,推不动了就交给延迟线程池稍后再来。
入口套了一层 Redis 分布式锁,同一笔单不允许两个线程同时推进:
machineExecute 的骨架:
private MachineStateEnum machineExecute(Long instructionId) {
MachineStateEnum result = MachineStateEnum.PROCESSING;
TransactionEventDTO event = transactionService.getByInstructionId(instructionId);
if (isFinal(event)) {
return MachineStateEnum.SUCCESS; // 已终态,幂等直接返回
}
while (!isFinal(event)) { // 只要没到终态,一直往前推
int stepBefore = event.getTradeState();
StateUnit stateUnit = getStateUnit(event); // 查转移表拿执行单元
try {
result = stateUnit.execute(event); // 执行业务
} catch (Exception e) {
result = MachineStateEnum.PROCESSING; // 异常当作可重试
}
event.setRunTimes(event.getRunTimes() + 1);
if (MachineStateEnum.FAIL.equals(result) && event.getFailState() == 0) {
event.setFailState(stepBefore); // 记录首次失败点
transactionService.updateTransaction(event, stepBefore);
}
int stepAfter = getNextTradeStatusByResult(event, result); // 查表算下一步
if (stepAfter == stepBefore) { // 推不动了(PROCESSING/自重试)
transactionService.updateTransaction(event, stepBefore);
result = MachineStateEnum.PROCESSING;
break;
}
event.setTradeState(stepAfter); // 推进游标
if (transactionService.updateTransaction(event, stepBefore) != 1) { // 乐观锁落库
throw new RuntimeException(UPDATE_ERROR); // 被别人改过,放手
}
}
if (!isFinal(event)) { // 没推到终态 → 自驱动,延迟后再来一次
if (event.getRunTimes() > maxRetryTimes) {
return result; // 超自驱动上限,交给外部定时兜底
}
scheduledExecutor.schedule(
() -> machineStateExecutor.execute(() -> run(event.getInstructionId())),
delayRunSecond, TimeUnit.SECONDS);
}
return result;
}
四个设计点是整套机制的关键:
4.1 一次调用,推到推不动为止
while (!isFinal(event)) 意味着:一次进来,只要状态机能往前跳,就连着一路跳下去——冻券成功立刻冻额度,冻额度成功立刻建借据……直到遇到一个返回 PROCESSING 的节点(推不动),或跑到终态,更能满足低延迟。
4.2 FAIL vs PROCESSING:补偿还是重试,全看这一步
这是整个 Saga 最需要拿捏的地方。执行单元返回什么,直接决定状态机是往回补偿还是原地重试。
- 明确业务失败(和下游约定错误码) → 走补偿
- 其它业务异常,可能抖动 → 原地重试
- 未知异常(超时等) → 原地重试
判断标准很清晰:
- 只有能明确判定"这事就是办不成"(业务规则拒绝、特定错误码),才返回
FAIL,触发补偿。 - 凡是"可能是临时故障"(超时、连接失败、未知异常),一律返回
PROCESSING,原地重试。
4.3 补偿节点的幂等与自重试
补偿动作本身也可能失败,所以补偿节点在转移表里的失败去向,通常指向它自己(自重试)。同时补偿单元必须幂等——因为它可能被重试多次。
整条 Saga 的最终一致,依赖的就是每个节点(无论正向还是补偿)都幂等,配合"至少执行一次"的重试,收敛到一致状态。
4.4 乐观锁 + 分布式锁,双重防并发
状态推进落库带乐观锁,WHERE trade_state = stepBefore,本质 SQL:
UPDATE main_transaction
SET trade_state = #{stepAfter}, run_times = #{runTimes}, fail_state = #{failState}
WHERE instruction_id = #{instructionId}
AND trade_state = #{stepBefore}; -- 乐观锁:状态被改过就更新 0 行
更新影响行数不等于 1,machineExecute 直接抛异常中断——说明这笔单已被另一个执行者推进,当前线程主动放手,避免重复推进。
为什么外层已有 Redis 锁,这里还要乐观锁? 两道防线各管一段:Redis 锁防"同一时刻两个线程同时进来推";乐观锁防"锁超时失效、锁误删、极端并发"下状态被并发修改。
4.5 自驱动 + 外部兜底
一次 machineExecute 推不到终态(遇到 PROCESSING),会通过延迟线程池 schedule 一个延迟任务,过 delayRunSecond 秒再 run 一次——状态机自己驱动自己往前。
但自驱动有次数上限 maxRetryTimes:超过后不再自 schedule,把这笔单交给外部兜底(MQ 延时消息或者外部MIS触发)。
五、关键设计小结
5.1 状态机与 Saga 的对应关系
| Saga 概念 | 本方案实现 |
|---|---|
| 长事务 | 一笔支用/退款的完整链路 |
| 子事务(正向) | 正向链上的 StateUnit(冻额度、建借据……) |
| 补偿动作 | 补偿链上的 StateUnit(解冻优惠券……) |
| 补偿触发 | 执行单元返回 FAIL,getFailState() 指向补偿态 |
| 事务游标 | 主单 trade_state 字段 |
| 编排定义 | Transition 转移表(声明式,一张表说清所有走向) |
| 最终一致 | 每节点幂等 + 至少执行一次 + 重试收敛 |
5.2 声明式编排
整条链路的走向全在 stateTransition 这张 enum 表里,没有一行 if-else 判断"下一步该干嘛"。加节点、改补偿路径,就是改表;新增业务线,就是加一张表 + 一个 TransitionManager Bean。调度引擎 machineExecute 完全通用,对具体业务零感知。
5.3 下游返回三个结果的边界
SUCCESS / FAIL / PROCESSING 的划分,直接决定状态机是前进、补偿还是重试。这个判断只能由最懂业务语义的执行单元来做——框架只负责"根据结果查表跳转",不会替业务判断"这算不算失败"。
5.4 持久化即状态
状态全落在主单表,进程完全无状态。任何实例、任何时刻,拿 instruction_id 把主单捞出来就能接着推。崩溃、重启、扩缩容,对一笔进行中的 Saga 毫无影响——这是它高可用的根本原因。
六、总结
Saga 状态机解决了两个问题:有依赖的多步骤怎么按序推进,以及某步失败怎么把副作用回滚干净。
落地方案上,用几个小接口搭起一个声明式状态机,把 Saga 的正向链和补偿链都编码进一张转移表;用主单表的一个 trade_state 字段做持久化游标;用"循环推进 + 乐观锁 + Redis 锁 + 自驱动 + MQ延时消息兜底"把它可靠地跑起来。