Saga分布式事务模式详解:事件编排与命令编排的工程实践

📅 发布时间:2026/9/17 2:37:35
Saga分布式事务模式详解:事件编排与命令编排的工程实践
做微服务的人大概率都被分布式事务折磨过。Saga 分布式事务模式这几个字听起来像是教科书里的概念但真到订单、库存、支付之间来回折腾的时候你才会意识到它到底在解决什么问题。我自己第一次接触 Saga是几年前把单体订单系统拆成订单、库存、账户三个服务的时候当时用 2PC 被锁表锁到怀疑人生后来换成 Saga 才把线上那些超时和死锁问题压住。这篇文章不是给你念论文而是把我对 Saga 的理解、两种实现方式、真实业务里的代码样例、以及踩过的坑一次性讲明白适合刚接触分布式事务的开发者也适合已经被 Seata 折磨过但没搞懂底层逻辑的老哥。1. 为什么需要 Saga分布式事务的痛点1.1 从本地事务到跨服务事务问题到底出在哪刚开始做服务拆分的时候很多人会觉得“不过是从一个大事务变成几个小事务”。但真相没那么简单。在单体应用里一个订单操作只需要在一个数据库里执行BEGIN TRANSACTION更新订单表、扣减库存表、记录账户流水任何一个操作失败都可以直接ROLLBACK数据最终是原子的。微服务拆分之后订单、库存、账户各占一个数据库甚至部署在不同机器上。这时候你没法用一句BEGIN去控制所有库因为没有一个连接能同时操作多个数据库。更重要的是有些操作根本不是数据库操作比如调用第三方支付网关、发送短信通知、扣减外部积分这些动作根本无法纳入数据库事务。所以你会面临一个现实问题订单创建成功但库存扣减失败或者支付完成但订单状态更新失败。系统不会崩溃只是数据最终变得不一致这比崩溃更难排查。分布式事务的本质不是“尽量让所有操作同时成功”而是“在无法保证同时成功的前提下设计一种机制让最终结果一致”。Saga 就是这类机制中非常典型的一种。它的核心思路是把一个大事务拆成一组有顺序的本地事务每个本地事务都对应一个补偿操作一旦某个环节失败就反向执行已经完成事务对应的补偿操作。简单说就是不追求“原子性”而是追求“可回滚”。1.2 为什么不用 2PC两阶段提交的痛很多人在想分布式事务时第一反应是 2PC两阶段提交或者 XA 协议。2PC 的思路是通过一个协调者分两步来统一提交第一阶段叫准备阶段所有参与者把资源锁住并告诉协调者“我能提交”第二阶段协调者根据所有参与者反馈决定提交或回滚。这套机制理论上保证了强一致性但实际落到微服务里问题一个接一个阻塞性太强从准备阶段开始所有涉及的资源都要被锁住事务时间越长锁范围越大高并发时数据库连接池很快就耗尽。单点风险协调者一旦挂了整个事务悬在半空所有参与者都不知道该提交还是该回滚。不适用长事务和外部调用如果某个参与者是第三方支付支付网关根本不会配合你玩“二阶段准备”收到扣款请求后钱就出去了你没法让支付网关“先锁住这笔钱等我广播再扣”。我印象很深的是最早用 Atomikos 配 XA 连接连接池里专门开一堆预留连接结果压测一上来业务线程全卡在prepare阶段数据库服务器 CPU 飙升最后只能把事务方案推倒重来。所以后来在业务允许最终一致性的场景下我几乎不再考虑 2PC而是优先看 Saga 和 TCC。TCCTry-Confirm-Cancel是另一个思路它要求业务提供确认和取消操作适合需要更强隔离性的场景但实现成本高Saga 则弱化了对“锁定资源”的要求靠反向补偿解决失败回滚适合大多数许最终一致的业务链路。2. Saga 模式的两种经典实现方式Saga 从 1987 年提出到现在落地形态基本分成两类Choreography事件编排和 Orchestration命令编排。这两个词中文很容易翻译混我按自己的理解给你讲清楚Choreography 就是没有总指挥每个服务通过监听事件自发协作像个即兴舞会Orchestration 则是有一个中心指挥者像交响乐团的指挥统一告诉每个服务该做什么。两者都有真实的生产案例没有绝对好坏只有适合不适合。2.1 Choreography事件编排模式详解Choreography 模式下系统没有中心化的协调器每个服务在本地事务完成后发布一个领域事件其他服务监听这个事件来决定下一步动作。比如一个经典的订单创建链路订单服务、库存服务、支付服务、积分服务各干各的通过事件串联。给你一个流程草图不画图用文字描述订单服务创建“新建订单”事件状态为“待支付”库存服务监听“新建订单”事件执行库存扣减成功后发布“库存已扣减”事件支付服务监听“库存已扣减”事件执行扣款如果扣款失败发布“扣款失败”事件订单服务监听“扣款失败”事件把订单状态改为“已取消”库存服务监听“订单已取消”事件执行库存回补。这种模式的优点很直观业务服务之间没有强制依赖任何人只需要关心自己监听的事件和自己发布的事件加一个新的参与者比如新增一个优惠券服务只需要关心“新建订单”和“订单已取消”事件不用改其他服务。代码层面也很干净没有庞大的协调类。缺点也非常致命当链路变长以后整个流程被事件“摊平”了。你要想知道“这个订单现在到底处于什么状态”得把所有服务的事件日志拉出来拼一下。而且事件风暴里如果某个服务没有收到消息、或者补偿操作反复重试排查链路非常痛苦。还有一个问题是循环依赖的风险服务 A 发事件触发服务 BB 回应事件触发 A一旦事件设计得不好可能出现事件风暴和死循环。2.2 Orchestration命令编排模式详解Orchestration 模式里有一个显式的 Saga 协调器Orchestrator它负责定义整个事务流程保存当前 Saga 的状态然后主动调用各个参与者的接口。参与者只需要提供“执行”和“补偿”两个方法不用关心整体流程。还是用订单链路举例协调器的逻辑大致是调用订单服务创建订单调用库存服务扣库存调用支付服务扣款如果扣款失败调用库存服务回补库存再调用订单服务将订单取消。在这种模式下协调器通常以状态机的形式存在。每个 Saga 实例从开始到结束都有确定的流程、状态、前置条件和补偿动作。状态机存储在数据库里比如用一个saga_instance表记录每个流程的状态这样即使协调器本身宕机重启后也能从上次暂停的状态继续推进。优点显而易见状态集中、流程清晰开发人员可以直接查看 Saga 状态机来确定业务走到哪一步卡住了调试时只需要看协调器的日志新增阶段也只需要修改状态机配置。缺点是协调器本身可能成为性能和可用性的瓶颈也存在额外开发成本。2.3 两种方式的实战对比与选型建议根据我自己的项目经验业务链路简单比如只有两三个服务且事件语义清晰的时候用 Choreography 很舒服开发效率高。但一旦链路超过四个节点或者涉及的团队比较多我会强烈建议切到 Orchestration。原因很简单人脑处理线性的状态机比处理一张复杂的事件网络容易得多线上问题排查时间基本能砍半。列一个对比表方便你直接参考维度Choreography事件编排Orchestration命令编排中心协调器无有负责状态推进与补偿服务间耦合通过事件解耦通过接口/命令耦合可跟踪性差需要拼接事件日志好协调器就是天然日志流程修改可能需要改多个服务的订阅逻辑只改状态机配置性能瓶颈无单点但事件量大会有压力协调器可能成为瓶颈可水平扩展实现成本低依赖消息中间件中高需要设计和实现状态机适合场景链路短、团队少、事件语义清晰链路长、跨团队、需要严格监控我自己的习惯是从一开始做 Saga 就尽量选 Orchestration。因为分布式事务本身就是个“高风险”环节风险控制比“少写几行代码”重要得多。你不想半夜爬起来为了排查一个订单状态翻十几个服务的日志找哪条事件没发出去。3. 核心概念与关键技术要点3.1 Saga Log事务状态机如何设计不管采用哪种实现方式Saga 都需要一个“账本”来记录每个实例的实时状态。这个账本在 Choreography 里往往依赖消息中间件的消息日志在 Orchestration 里则是专门的状态机表。我曾经踩过一个坑早期做 Choreography没有额外的状态记录某次Kafka消费者批量重平衡消息重复消费补偿操作执行了两次账目全乱。后来才意识到Saga 事务必须有可查询的状态记录不能只靠日志来还原现场。比较推荐的做法是设计两张核心表saga_instance记录每个 Saga 实例的全局状态字段包括saga_id、saga_type、statusPENDING / COMPENSATING / COMPENSATED / SUCCEEDED / FAILED、current_step、payload、create_time、update_time。saga_step_log记录每个步骤的执行状态字段包括step_id、saga_id、step_name、statusEXECUTING / COMPENSATING / COMPENSATED / SUCCEEDED / FAILED、error_message、retry_count。每次状态变化都要写入这两张表而且更新操作最好加上乐观锁比如version字段避免并发重试导致状态被覆盖。这里其实和数据库本地事务处理方式相似唯一的差别是“本地事务”变成了“一个 Saga 实例”的流转。状态机设计时要考虑清楚主流程状态和补偿流程状态的区别。比如主流程是ORDER_CREATED - INVENTORY_DEDUCTED - PAYMENT_PAID补偿流程是PAYMENT_COMPENSATED - INVENTORY_COMPENSATED - ORDER_CANCELLED。每个状态变换要有明确的前置条件保证状态机不会出现“死循环”或“跳变”。在实践中我一般会用枚举或者状态机框架来约束合法迁移路径防止代码写多了以后随意改状态带来的隐患。3.2 补偿事务的幂等与空补偿问题Saga 里最容易被低估的两个问题是幂等和空补偿。先讲幂等这是所有分布式系统里必须面对的问题。由于网络超时、消息重复投递、消费者重试等原因一个补偿操作完全可能被触发两次。比如库存回补这个动作如果同一个 Saga 实例在回补库存时第一次执行超时协调器重试调用了第二次结果库存被多回补一次账面库存数量就错了。解决思路其实不复杂每个补偿操作需要带上唯一的业务流水号比如saga_id step_id目标服务的处理逻辑里先查流水表如果已经处理过直接返回成功不再重复执行。这个“防重表”或者“去重索引”是必须做的别指望消息中间件的“恰好一次投递”实际基础设施很难保证端到端幂等。再说空补偿。这个概念听起来很绕明明该补偿了但你会发现当前步骤并没有真正执行成功过却触发了补偿。典型场景是协调器调用库存服务扣减库存因为网络超时协调器认为失败决定进入补偿流程开始调用回补库存但那个“超时”的扣减请求实际上也被库存服务成功执行了只是响应丢失。如果不做处理就会出现“根本没有扣减成功却执行了回补”最终数据比原来还多。解决空补偿的方法是引入一个“事务记录表”任何参与者在执行正向操作之前先在本地记录一个待执行标记正向操作真正成功后再把它更新为“已执行”。补偿处理前先检查该事务记录是否存在且为“已执行”如果不是说明是空补偿直接跳过补偿逻辑。这个操作和本地消息表模式很像本质上都是在“没有全局事务”的前提下靠本地记录来保持一致。3.3 隔离性不要以为 Saga 有 ACID 的“I”Saga 最大的软肋是没有隔离性。这句话怎么理解传统事务里A 事务还没提交B 事务是看不到 A 修改的数据的这是隔离性在起作用。但 Saga 中每个子事务都是直接提交的中间状态对其他人可见。比如订单扣减库存成功、支付尚未完成时其他消费者去查询库存会看到库存“少了一件”但实际上订单并未完成支付这可能造成超卖或业务误判。有几种常用手段做“业务级隔离”状态标记 语义锁在业务表里加一个status字段把正在参与 Saga 的数据标记为“待确认”或“锁定中”。其他业务逻辑读取时可以判断这个标记把中间态的数据排除掉或者只允许特定操作访问。乐观锁版本号在涉及数据更新的每条记录上加version字段每次更新前比较版本冲突则拒绝或重试。这种方式可以在一定程度上避免“多人同时操作同一条记录”导致的数据错乱。读时校验对于无法做状态标记的外部系统可以做一个“校验步骤”比如在支付或发货前重新核对订单和库存状态不一致就拒绝继续。我之前做过一个订单系统在订单表里增加了一个txn_status字段来表示“Saga进行中”所有面向用户的查询都会过滤掉txn_status ! CONFIRMED的订单。这样虽然底层数据短暂“脏”但业务层给用户展示的视图是干净的。这是一种很实用的“脏读隔离”策略。Saga 本身不保证隔离但业务上必须自己补隔离方案否则后期会冒出一堆奇怪的并发问题。4. 一个完整的 Saga 实现案例订单与库存分布式事务4.1 业务场景与事务边界拆分用一个最常见的订单与库存场景用户下单 - 扣库存 - 扣款 - 加积分。这个链路有四个参与者四个服务。单体时代一条 SQL 事务就结束了微服务时代你没法用一个数据库事务罩住所有表只能靠 Saga 串起来。我们按事务边界拆一下每个服务对外暴露两个方法一个是正向方法一个是补偿方法。订单服务创建订单正向取消订单补偿库存服务扣减库存正向回补库存补偿账户服务扣减余额正向回补余额补偿积分服务增加积分正向扣减积分补偿。这里的顺序不是固定的。你可以先扣库存再扣款也可以先扣款再扣库存。我的建议是把“风险较高的外部依赖”放在后面。因为越往后失败时触发补偿的范围越小。比如支付外呼最可能失败就应该把支付放在库存之后而不是第一个否则支付一旦失败你还要补偿订单创建和库存扣减徒增复杂度。4.2 事件编排的代码级实现伪代码假设我们采用 Kafka 作为事件总线事件编排实现大概是这个样子。订单服务发布创建订单事件public void createOrder(OrderDTO order) { try { // 本地业务操作保存订单状态为待支付 orderRepository.save(order); // 发布领域事件 eventPublisher.publish(new OrderCreatedEvent(order.getId(), order.getItems())); } catch (Exception e) { // 本地事务失败直接抛出 throw new BusinessException(订单创建失败); } }库存服务监听事件并扣减库存KafkaListener(topics order_event) public void handleOrderCreated(OrderCreatedEvent event) { try { stockService.deduct(event.getItems()); eventPublisher.publish(new InventoryDeductedEvent(event.getOrderId())); } catch (Exception e) { // 记录错误根据重试策略后续处理 eventPublisher.publish(new InventoryDeductFailedEvent(event.getOrderId(), e.getMessage())); } }支付服务监听库存扣减成功、执行扣款KafkaListener(topics inventory_event) public void handleInventoryDeducted(InventoryDeductedEvent event) { try { paymentService.pay(event.getOrderId()); eventPublisher.publish(new PaymentPaidEvent(event.getOrderId())); } catch (Exception e) { eventPublisher.publish(new PaymentFailedEvent(event.getOrderId(), e.getMessage())); } }订单服务监听支付失败进入补偿KafkaListener(topics payment_event) public void handlePaymentFailed(PaymentFailedEvent event) { orderService.cancelOrder(event.getOrderId()); eventPublisher.publish(new OrderCancelledEvent(event.getOrderId())); }库存服务监听订单取消回补库存KafkaListener(topics order_event) public void handleOrderCancelled(OrderCancelledEvent event) { stockService.compensate(event.getOrderId()); }这套代码很容易读但实际运行起来你会发现几个问题。第一如果库存服务一直处于不稳定状态事件被反复重试消费端需要考虑幂等和死信队列。第二事件之间的顺序需要设计好比如PaymentFailedEvent和OrderCancelledEvent如果重复发布库存回补会不会被执行多次所以每个消费者都要做去重。第三没人统一记录 Saga 实例的状态出了问题只能靠消息中间件的消费位置去回放非常痛苦。这就是为什么到了生产级别我更倾向用编排模式。4.3 命令编排的实现方式对比同样这个场景用 Orchestration 来实现就会清晰很多。定义一个 Saga 协调器它就是一个简单的状态机引擎负责顺序调用参与者。伪代码如下public class CreateOrderSagaOrchestrator { private final OrderClient orderClient; private final InventoryClient inventoryClient; private final PaymentClient paymentClient; private final SagaStateRepository sagaStateRepository; public void executeCreateOrder(CreateOrderRequest request) { String sagaId UUID.randomUUID().toString(); sagaStateRepository.save(new SagaState(sagaId, STARTED)); // 1. 创建订单 CreateOrderResponse response orderClient.createOrder(request); sagaStateRepository.updateStep(sagaId, ORDER_CREATED); // 2. 扣减库存 try { inventoryClient.deductInventory(request); sagaStateRepository.updateStep(sagaId, INVENTORY_DEDUCTED); } catch (Exception e) { // 扣库存失败补偿已经创建的订单 orderClient.cancelOrder(response.getOrderId()); sagaStateRepository.updateStep(sagaId, FAILED, e.getMessage()); return; } // 3. 执行扣款 try { paymentClient.pay(request); sagaStateRepository.updateStep(sagaId, PAYMENT_PAID); } catch (Exception e) { // 支付失败补偿库存和订单 inventoryClient.compensateInventory(request); orderClient.cancelOrder(response.getOrderId()); sagaStateRepository.updateStep(sagaId, FAILED, e.getMessage()); return; } // 4. 增加积分若失败通知但不阻断主流程 loyaltyClient.addPoints(request); sagaStateRepository.updateStep(sagaId, SUCCEEDED); } }这个写法比事件编排直观得多。你一眼就能看到先做什么失败后补偿什么。不过这个版本只是一个入门示例实际生产环境我不会把 Saga 状态机写死在 Java 代码里而是使用配置化状态机比如 Seata Saga 的 JSON DSL。这样以后加一个“发送通知”的步骤不用改代码只改配置。4.4 关于 Seata Saga 的快速上手国内做 Java 后端的大多绕过不了 Seata。Seata 提供了 AT、TCC、SAGA、XA 四种模式其中 SAGA 模式就是基于状态机引擎实现的编排模式。Seata 的 Saga 实际上借鉴了轻量级状态机思路用 JSON 描述状态和动作。一个简单的 Seata Saga 状态机配置大概长这样{ Name: createOrderSaga, Start: create_order, States: { create_order: { Type: ServiceTask, DurableRequest: false, ServiceName: order-service, ServiceMethod: createOrder, Next: deduct_inventory, Compensate: cancelOrder }, deduct_inventory: { Type: ServiceTask, ServiceName: inventory-service, ServiceMethod: deduct, Next: pay_order, Compensate: compensateInventory }, pay_order: { Type: ServiceTask, ServiceName: pay-service, ServiceMethod: pay, Next: success, Compensate: refund }, success: { Type: Succeed, Status: SUCCEEDED } } }用 Seata 的好处是框架帮你处理了状态持久化、Saga 实例注册、补偿调用、重试等通用逻辑你只需要实现参与者接口。不过坏处是引入了一个较重的中间件依赖需要 Seata ServerTC配合 Nacos/Eureka 使用。如果团队对 Seata 已经比较熟用它做 Saga 是省力的选择但如果只是在一个很小的系统里用我也不建议为了一个分布式事务专门搭建 Seata直接用消息队列做事件编排就够了。5. 常见问题与排查技巧实录5.1 补偿事务执行失败怎么办Saga 最大的噩梦不是主流程失败而是失败后的补偿操作也失败了。比如支付成功后需要回补库存结果回补库存的接口又因为数据库连接池异常挂掉。这时候 Saga 实例会一直停留在COMPENSATING状态。处理方案一般分三层重试给补偿操作设置合理的重试策略比如 3 次每次间隔指数退避重试要注意幂等防止重复补偿导致数据错乱。死信处理重试超过最大次数把 Saga 实例标记为FAILED同时把失败消息投递到死信队列触发告警通知开发人员人工介入。人工兜底如果业务链路没有复杂到无法修复可以写一个定时任务扫描长时间处于COMPENSATING状态的 Saga针对特定状态执行人工补单或者对账脚本。我的经验是补偿逻辑一定要比正向逻辑更健壮。不能觉得补偿操作只在异常时触发一次就粗心大意地不写幂等、不写重试。实际上线上大部分灾难都是“补偿又失败”导致的连环噩耗。5.2 事务悬挂如何规避前面提到的空补偿是“补偿跑了但正向没执行成功”。还有一个更隐蔽的问题是事务悬挂正向操作还没到达补偿操作先到了等正向操作到达时事务已经处于回滚状态这会导致状态机混乱。举一个具体场景协调器调用库存服务扣减库存由于网络超时协调器进入补偿流程马上调用回补接口。但扣减请求在网络上被延迟了很久直到协调器已经回补完成之后才到达库存服务执行了一次扣减。最终库存被扣了但没有补偿可以撤销这次扣减。规避手段可以从两个方向入手协调器侧每个 Saga 实例在进入补偿前必须确认当前状态是可补偿的如果状态机设计成INVENTORY_DEDUCTED状态才能补偿那么即使扣减请求还没落库协调器也不能盲目补偿。参与者侧在处理正向操作时先查询当前 Saga 的状态是否已经是COMPENSATED如果是直接拒绝执行正向操作或者写入一条悬挂记录等待调度器处理。我在实际项目里更多依赖协调器侧的状态管理因为参与者侧不好判断全局状态容易被其他流程干扰。协调器把状态流转做成有方向的迁移例如没有INVENTORY_DEDUCTED就不能进入INVENTORY_COMPENSATING基本可以把悬挂问题拦在门口。5.3 状态机状态流转异常如何排查如果你用了编排模式状态机状态异常应该很好排查。我经常遇到的现象是某个 Saga 实例长时间停在INVENTORY_DEDUCTED说明下一步PAYMENT_PAID没有推进。这时候从saga_step_log表能看到pay_order的步骤状态是EXECUTING而且没有更新说明支付调用一直没有返回。再查协调器日志和支付服务日志通常能定位到是网络超时还是支付接口自身的问题。但如果是事件编排模式排查就要复杂得多。建议一定要给每个事件补全traceId和sagaId在消息头和日志里都带上否则你无法串联一条事件的完整链路。同时消息消费端要记录消费日志和消费异常日志这样出了问题才能根据 eventId 反向查找是哪个消费者漏消费了或者哪段重试死循环了。我还有个习惯为 Saga 实例提供一张“状态机可视化”页面。不需要很复杂就是从表里读数据画个流程进度图运维和开发能直接看到每个订单停留在哪个节点。这套页面对排查线上问题帮助极大可以少开很多“语音会议”。5.4 性能与超时调优经验Saga 本身强调的是可用性和最终一致性但它也不是没有性能瓶颈。尤其是在纯编排模式下所有流程都由协调器串行驱动协调器的 QPS 上限就决定了整个事务链路的创建上限。要提升性能可以有这么几个方向协调器无状态化把 Saga 状态存到 Redis 或数据库协调器多个节点可以水平扩展每个请求落到任意节点都能接着旧状态走。异步化不必每一步都同步等返回。比如创建订单后可以异步通知扣库存协调器只在状态回调时感知结果。这里的异步化不是放弃一致性而是减少同步阻塞时间。超时设置要有余量扣库存接口设置 2 秒超时但支付接口可能因为第三方响应慢至少设置 10 秒超时。超时时间如果设置太短容易误判失败触发不必要的补偿太长又会让 Saga 实例一直卡住。最好为每个步骤单独配置超时。批量补偿如果大量 Saga 实例同时进入补偿状态不要每个实例单独开线程补偿可以使用批量任务把需要补偿的步骤捞出来统一处理减轻下游压力。最后再分享一个我个人的经验不要在 Saga 状态机里塞进“查询”操作。尽管有时你很想在补偿前确认一下当前业务状态但查询操作会让状态机依赖读接口的可用性一旦读接口挂了连补偿都没法执行。真要校验可以让补偿操作本身把“状态”作为参数传递由参与者在自己数据库里校验这样更可靠。用 Saga 这么多年我没有哪一次事故是“状态机设计得太简单”导致的倒是有好几次因为“想一下子解决所有问题”把状态机搞得太复杂反而出了岔子。分布式事务没有银弹Saga 也一样它就是让你在强一致和可用性之间找到一个让业务能继续跑下去的平衡点。