Saga 执行器 (Saga Executors)

概述 (Overview)

在 StackSaga 中,Saga 执行器 (Saga Executor) 是一个专用组件,负责在分布式 Saga 工作流中封装并执行单个原子事务 (Atomic Transaction)。 它充当业务逻辑与 Saga 编排引擎之间的桥梁,确保长事务 (LRT - Long-Running Transaction) 中的每个步骤都能可靠且一致地执行。 Saga 执行器分为两类:命令执行器 (Command Executors,处理修改状态的操作,同时支持主执行与补偿执行) 与查询执行器 (Query Executors,处理只读操作)。 通过将每个原子事务隔离在独立的执行器中,StackSaga 实现了精确的控制、重试与补偿机制,降低了数据异常的风险,并保证了跨微服务的最终一致性 (Eventual Consistency)。

如果您刚接触 Saga 模式,建议先阅读 长事务中的原子事务与幂等性 (Atomic Transactions & Idempotency in LRT),了解分布式系统中的幂等性 (Idempotency) 和原子操作概念。

通常,一个原子执行包含两个阶段:主执行 (Primary Execution) 与 补偿执行 (Compensating Execution)。

  1. 主执行 (Primary Execution / Main Execution)

    • 主执行代表长事务中的核心业务操作。 每个主执行都是整体业务工作流中的独立步骤,推动流程向最终完成演进。 这些操作通常是幂等且隔离的,确保它们可以独立执行而不会产生意外的副作用。

      例如,在下单工作流中,主执行可能包括获取用户详情、初始化订单、进行预授权、扣减库存以及处理最终支付。

    • 这些主执行推动流程向前演进。[]

  2. 补偿执行 (Compensating Execution / Revert Execution)

    • 补偿执行旨在当后续步骤发生失败时,撤销先前已成功完成的主执行所产生的影响。 每个主执行都有一个对应的补偿动作来逆转其影响,确保系统在发生故障时能够恢复到一致状态。

    • 补偿执行推动流程向后回滚。[]

执行器类型 (Executor Types)

基于上述分类,Saga 执行器主要分为两种类型。但除了命令执行器与查询执行器之外,StackSaga 还提供了一种特殊的执行器类型——子执行器 (Sub Executor),用于处理额外扩展的补偿执行。

下图展示了 StackSaga 中的执行器类型及其包含的方法:

Stacksaga 执行器类型

上述每种执行器类型均提供了阻塞式 (Blocking) 与非阻塞响应式 (Reactive) 两种实现。您可以根据技术栈需求自由选择。

命令执行器 (Command Executors)

如果某个原子执行同时包含主执行与补偿执行,该类原子事务必须在命令执行器 (Command Executor) 中实现。
命令执行器 包含两个核心方法:分别用于执行主流程与执行补偿回滚流程。

命令执行器的典型示例:

  • 初始化订单 (Initialize order)

  • 预留库存商品 (Reserve the items)

  • 处理支付扣款 (Make the payment)

阻塞式命令执行器 (Blocking Command-Executor)

@SagaExecutor(executeFor = "order-service", value = "initializeOrderExecutor") (1)
@AllArgsConstructor
public class InitializeOrderExecutor implements CommandExecutor<OrderDomainEntity> { (2)

    private final OrderService orderService;

    @Override (3)
    public ProcessStepManager<OrderDomainEntity> doProcess(
            OrderDomainEntity currentDomainEntity,
            ProcessStepManagerUtil<OrderDomainEntity> stepManager,
            String idempotencyKey
    ) throws RetryableExecutorException, NonRetryableExecutorException {

        try {
            (4)
            String orderId = this.orderService.initializeOrder(
                    currentDomainEntity.getUsername(),
                    currentDomainEntity.getProductItems(),
                    currentDomainEntity.getTotalAmount(),
                    idempotencyKey
            );
            currentDomainEntity.setOrderId(orderId);
            return stepManager.next(RserveItemsExecutor.class, "INITIALIZED_ORDER"); (5)
        } catch (OperationAlreadyExecutedException alreadyExecutedException) {
            (6)
            return stepManager.next(RserveItemsExecutor.class, "INITIALIZED_ORDER");
        } catch (FeignException.ServiceUnavailable unavailableException) {
            (7)
            throw RetryableExecutorException.of(unavailableException);
        } catch (FeignException.BadRequest badRequestException) {
            (8)
            throw NonRetryableExecutorException
                    .buildWith(badRequestException)
                    .put("time", LocalDateTime.now())
                    .put("reason", "BadRequest")
                    .build();
        }
    }

    @Override
    @RevertBefore(startFrom = OrderInitializeSubBeforeExecutor.class) (12)
    @RevertAfter(startFrom = OrderInitializeSubAfterExecutor.class) (13)
    public RevertStepManager doRevert(
            NonRetryableExecutorException processException,
            OrderDomainEntity finalDomainEntityState,
            RevertHintStore revertHintStore,
            String idempotencyKey,
            RevertStepManagerUtil stepManager
    ) throws RetryableExecutorException {
        try {
            (9)
            this.orderService.cancelOrder(finalDomainEntityState.getOrderId(), idempotencyKey);
            return stepManager.done("CANCELLED_ORDER");
        } catch (OperationAlreadyExecutedException alreadyExecutedException) {
            (10)
            return stepManager.done("CANCELLED_ORDER");
        } catch (FeignException.ServiceUnavailable unavailableException) {
            (11)
            throw RetryableExecutorException.of(unavailableException);
        } catch (FeignException.BadRequest badRequestException) {
            revertHintStore.put("InitializeOrderExecutor:FAILED", processException.getMessage());
            throw new RuntimeException();
        }
    }
}
1 执行器使用 @SagaExecutor 注解声明为 Spring Bean,并提供必要的元数据(如该执行器所属的目标服务名称以及执行器的全局唯一标识名称)。
2 执行器实现 CommandExecutor 接口以声明为命令执行器。
3 重写 doProcess 方法以实现主执行逻辑。
StackSaga 会自动将 idempotencyKey(幂等键)作为第三个参数传入。该键在重试该原子执行时保持确定且稳定。您应将此键透传给下游服务(或放入 HTTP/消息头中),当下游服务收到重试请求时便能识别出重复请求。
NOTE: 了解 StackSaga 如何生成此键或如何自定义生成规则,请参阅 领域实体中的幂等键生成。
4 在 doProcess 方法内实现主执行业务逻辑。调用下游服务时,连同业务参数一并透传 idempotencyKey。
5 通过 stepManager 导航、暂停或完成事务:
stepManager 是传入 doProcess 的核心流程控制工具类。它提供了三个关键方法来控制 Saga 的生命周期流转:
  • stepManager.next(…​)(连续向前导航):
    用于 连续型长事务 (Continuous LRT)。在当前步骤成功完成后,立即将 Saga 推进到下一个执行器跨度 (Span):

    return stepManager.next(ReserveItemsExecutor.class, "INITIALIZED_ORDER");
  • stepManager.pause(…​)(面向可暂停 LRT 的业务等待状态):
    用于 可暂停长事务 (Pausable LRT)。进入预期的业务等待状态,此时执行挂起,领域状态安全持久化至事件存储库,工作线程被释放,直到外部回调 (Callback) 或 Webhook 恢复该事务:

    // 无限期等待外部回调(不暴露给自动重试引擎):
    return stepManager.pause(DeliveryExecutor.class, "ORDER_PREPARING", null);
    
    // 或者在重新调用已暂停执行器之前等待指定的超时时长:
    return stepManager.pause(DeliveryExecutor.class, "ORDER_PREPARING", Duration.ofHours(5));
    关联键与 pause(…​) 中的超时处理:
    • idempotencyKey 作为 correlationKey(关联键):在可暂停步骤中,将 idempotencyKey(作为 doProcess 的第 3 个参数提供)传递给下游异步服务(放在请求体、查询参数或 HTTP 标头中)。当下游服务稍后调用您的 Webhook/回调端点时,必须回传此完全相同的键。然后您将其作为 correlationKey 传入 sagaTemplate.resume(transactionId, correlationKey) 以安全恢复执行。详见 可暂停 LRT 架构。

    • 超时时长 (Timeout Duration):第三个参数 (Duration) 指定引擎在将已暂停的执行器暴露给重试重新调用之前,应等待外部回调的时长。如果传入 null,事务将无限期保持暂停,不会暴露给自动重试。如果指定了时长(例如 Duration.ofHours(5)),若 5 小时内未收到回调,框架将重新调用该暂停的执行器跨度,允许应用程序重新查询外部系统或触发补偿。

    • 通过事件溯源恢复原始干净状态 (Pristine State Restoration):超时重新调用执行器时,StackSaga 的事件溯源引擎会从事件存储库中重构 currentDomainEntity,将其精确恢复到进入该执行器之前的初始状态。上一次暂停尝试期间所做的任何内存中修改都将被丢弃。执行器将作为一次全新调用针对原始状态执行,彻底防止跨重试的脏状态泄漏。

  • stepManager.complete(…​)(终态完成):
    在执行完当前执行器后,成功将整个长事务标记为终态完成:

    return stepManager.complete("INITIALIZED_ORDER");

stepManager 方法与事件动作最佳实践:

stepManager 为导航 (next)、暂停 (pause) 和完成 (complete) 事务提供了重载方法签名:

  1. 直接以 String 传递事件名称(直接便捷):
    您可以直接将事件动作名称作为 String 传递(例如 stepManager.next(ReserveItemsExecutor.class, "INITIALIZED_ORDER")、stepManager.pause(DeliveryExecutor.class, "ORDER_PREPARING", null) 或 stepManager.complete("INITIALIZED_ORDER"))。底层 StackSaga 会自动将此事件动作持久化到事件存储库中。

  2. 以 SagaExecutionEventName 传递事件名称(集中式枚举 - 强烈推荐):
    相比于在各个执行器中散落硬编码字符串,框架推荐的最佳实践是为每个自定义 LRT 领域实体(如 OrderDomainEntity)维护一个实现 SagaExecutionEventName 的集中式枚举 (enum):

    public enum OrderExecutionEvent implements SagaExecutionEventName {
        INITIALIZED_ORDER,
        CANCELLED_ORDER,
        FETCHED_USER_DETAILS,
        ORDER_PREPARING,
        MADE_PAYMENT,
        REFUNDED_PAYMENT,
    }

    这在编译期提供了类型安全以及跨整个 Saga 的集中可追溯性:

    // doProcess 中的向前导航(连续型 LRT):
    return stepManager.next(ReserveItemsExecutor.class, OrderExecutionEvent.INITIALIZED_ORDER);
    
    // doProcess 中的预期业务暂停(可暂停 LRT):
    return stepManager.pause(DeliveryExecutor.class, OrderExecutionEvent.ORDER_PREPARING, Duration.ofHours(5));
    
    // doProcess 中完成整个事务:
    return stepManager.complete(OrderExecutionEvent.INITIALIZED_ORDER);
    
    // doRevert 中的补偿完成:
    return stepManager.done(OrderExecutionEvent.CANCELLED_ORDER);
  3. 动作事件命名规范 (Naming Convention):

    • 首词:必须是过去分词动词(例如 INITIALIZED, CANCELLED, FETCHED, MADE, REFUNDED, REVERTED, EXECUTED)。

    • 次词:代表真实的业务主体或领域实体(例如 ORDER, USER_DETAILS, PAYMENT)。

      示例:INITIALIZED_ORDER, CANCELLED_ORDER, FETCHED_USER_DETAILS, MADE_PAYMENT, REFUNDED_PAYMENT。

  4. 通过 RevertStepManagerUtil 在 doRevert(…​) 中返回 RevertStepManager:
    在返回 RevertStepManager 的补偿方法(如 doRevert(…​))中,您使用 stepManager.done(…​) 来完成当前补偿步骤(例如 return stepManager.done(OrderExecutionEvent.CANCELLED_ORDER); 或 return stepManager.done("CANCELLED_ORDER");),或者使用 stepManager.pause(…​) 进入等待异步回调的业务等待状态(例如 return stepManager.pause("REFUND_INITIATED", Duration.ofHours(24));)。详见 doRevert 中的可暂停补偿。

+ <6> 处理已执行的幂等响应(开发者职责):

+

对已执行过的操作返回成功步骤 — 切勿抛出异常:
请注意,此代码片段中的 OperationAlreadyExecutedException 是纯粹为了演示而创建的自定义异常。在生产实践中,下游微服务如何指示操作已被处理完全取决于具体的架构——例如 Feign 客户端抛出 HTTP 409 Conflict、自定义领域异常、或是 API 响应体中包含 "ALREADY_PROCESSED" 状态码。

作为开发者,您有责任在执行器内部捕获或识别该重复执行信号,并进行恰当处理:

  • 您必须将其视为一次正常的成功执行并返回下一步骤 (stepManager.next(…​)),切勿向外抛出错误。

  • 因为目标微服务在先前的尝试中已经持久化了预期的业务状态,返回下一步骤会通知 StackSaga 该原子工作单元已成功达成,从而允许工作流向前推进。如果抛出异常或任由异常向外冒泡,框架将错误地把该已完成的步骤判定为失败,进而触发非必要且灾难性的补偿回滚!

+ <7> 捕获可重试的资源不可用异常,通知 SEC 将事务保持在重试模式,并按照配置暴露给下一个重试调度周期。
您可以使用 RetryableExecutorException.of(Exception e) 方法包装原始异常,创建 RetryableExecutorException 的新实例。
如果您抛出自定义异常且未包装在 RetryableExecutorException 中,SEC 将其视为不可重试异常,立即停止向前执行,并立即以逆序触发补偿执行。 <8> 捕获不可重试异常,通知 SEC 终止前向执行,并立即以逆序启动补偿回滚流程。
您可以使用 NonRetryableExecutorException.buildWith(Exception e) 包装原始异常,并通过 put(String key, Object value) 将任何业务上下文元数据持久化到事件存储库中。这些元数据稍后可以在补偿执行中通过 RevertHintStore 提取使用。 <9> 重写 doRevert 方法以实现补偿回滚业务逻辑。
该方法接收 idempotencyKey 以透传给下游补偿端点确保逆向动作幂等,并接收 RevertStepManagerUtil stepManager 以导航或暂停回滚流程。恢复执行时还可选接收 ResumeContext。 <10> 补偿执行中的幂等返回(开发者职责):
与 doProcess 完全一致,在补偿期间处理重复执行同样是开发者的核心职责。如果下游服务指示取消或回滚动作在先前的重试尝试中已经执行完毕(此处通过捕获 OperationAlreadyExecutedException 表示),必须返回成功步骤 (stepManager.done("CANCELLED_ORDER"))。这向 SEC 确认该步骤的补偿已圆满完成,允许框架继续按逆序补偿上一个历史步骤(或者在所有补偿完成后圆满结束回滚)。 <11> 捕获补偿过程中可重试的临时资源不可用异常,以保持事务处于重试模式。
如果补偿过程中遇到可以安全忽略的不可重试异常,应捕获该异常并通过 put(String key, Object value) 将元数据存储到 RevertHintStore 中,避免事务异常终止。 <12> @RevertBefore 注解用于指定一个前置补偿子执行器 (Sub-Before-Executor),该子执行器将在当前命令执行器的主补偿执行之前运行。 <13> @RevertAfter 注解用于指定一个后置补偿子执行器 (Sub-After-Executor),该子执行器将在当前命令执行器的主补偿执行完成之后运行。

非阻塞响应式命令执行器 (Non-Blocking Command-Executor)

相比阻塞式示例,此处仅针对响应式特性高亮说明。
@SagaExecutor(executeFor = "order-service", value = "initializeOrderExecutor")
@AllArgsConstructor
public class ReactiveInitializeOrderExecutor implements ReactiveCommandExecutor<OrderDomainEntity> { (1)

    private final OrderService orderService;

    @Override (2)
    @NonNull
    public Mono<ProcessStepManager<OrderDomainEntity>> doProcess(
            OrderDomainEntity currentDomainEntity,
            ProcessStepManagerUtil<OrderDomainEntity> stepManager,
            String idempotencyKey
    ) {
        (3)
        return this.orderService
                .createOrder(
                        currentDomainEntity.getUserData(),
                        currentDomainEntity.getOrderDetails(),
                        idempotencyKey
                )
                .map(orderId -> {
                    currentDomainEntity.setOrderId(orderId);
                    return stepManager.next(ReactiveRserveItemsExecutor.class, "INITIALIZED_ORDER");
                })
                .onErrorResume(throwable -> {
                    if (throwable instanceof OperationAlreadyExecutedException) {
                        // 对于已执行的操作返回正常成功步骤 — 切勿报错
                        return Mono.just(stepManager.next(ReactiveRserveItemsExecutor.class, "INITIALIZED_ORDER"));
                    } else if (throwable instanceof ResourceUnavailableException) {
                        return Mono.error(RetryableExecutorException.of(throwable));
                    } else {
                        return Mono.error(NonRetryableExecutorException.buildWith(throwable)
                                .put("time", String.valueOf(System.currentTimeMillis()))
                                .put("reason", "BadRequest")
                                .build());
                    }
                });
    }


    @Override (4)
    @NonNull
    public Mono<RevertStepManager> doRevert(
            NonRetryableExecutorException processException,
            OrderDomainEntity finalDomainEntityState,
            RevertHintStore revertHintStore,
            String idempotencyKey,
            RevertStepManagerUtil stepManager
    ) {
        (5)
        return this.orderService
                .cancelOrder(finalDomainEntityState.getOrderId(), idempotencyKey)
                .thenReturn(stepManager.done("REVERTED_ORDER_INITIALIZATION"))
                .onErrorResume(throwable -> {
                    if (throwable instanceof OperationAlreadyExecutedException) {
                        // 对于已执行的补偿操作返回成功步骤
                        return Mono.just(stepManager.done("REVERTED_ORDER_INITIALIZATION"));
                    } else if (throwable instanceof ResourceUnavailableException) {
                        return Mono.error(RetryableExecutorException.of(throwable));
                    } else {
                        revertHintStore.put("InitializeOrderExecutor:FAILED", processException.getMessage());
                        return Mono.just(stepManager.done("IGNORED"));
                    }
                });
    }
}
1 执行器实现 ReactiveCommandExecutor<OrderDomainEntity> 接口以声明为响应式命令执行器。泛型类型为您在执行器中使用的领域实体类。
2 重写 doProcess 方法,以响应式管道方式实现主执行逻辑。
3 主执行业务逻辑在 doProcess 内利用响应式管道执行。
请注意 idempotencyKey 被透传给下游服务调用 (createOrder)。在 onErrorResume 中,如果发生重复执行异常(如 OperationAlreadyExecutedException),管道通过返回 Mono.just(stepManager.next(…​)) 恢复并将其视为成功,而不是向外传播错误。
异常不能传递到管道外部(例如命令式抛出),且方法绝对不能返回 null 或 Mono.empty()。否则会被视为非预期主执行异常,终止前向执行并立即启动逆向补偿回滚。
4 重写 doRevert 方法,以响应式方式实现补偿执行逻辑并返回 Mono<RevertStepManager>。
5 补偿执行逻辑在 doRevert 内利用响应式管道实现。
应用相同的幂等策略:将 idempotencyKey 传递给 cancelOrder,并在 onErrorResume 中识别已执行异常,返回包含成功补偿步骤的 Mono.just(stepManager.done(…​)),以便框架确认该步骤已成功补偿。
补偿执行期间绝不能发生未处理的不可重试异常。若存在可安全忽略的异常,必须将其捕获并持久化到 RevertHintStore,然后返回 Mono.just(stepManager.done("IGNORED")) 以避免整个事务非正常终止。

在前向执行 (doProcess) 中消费 ResumeContext

当执行器暂停等待外部回调(例如 3D-Secure 支付确认、商户接单审批或 KYC 人工审核)时,外部 Webhook 会通过 sagaTemplate.resume(txId, correlationKey).put(…​) 传递载荷数据。

为了安全消费此外部载荷,而不将核心领域聚合暴露给未经校验的直接变更,请实现接收 ResumeContext 的重载 doProcess(…​) 方法:

@Override
public ProcessStepManager<OrderDomainEntity> doProcess(
        OrderDomainEntity currentDomainEntity,
        ProcessStepManagerUtil<OrderDomainEntity> stepManager,
        String idempotencyKey,
        ResumeContext resumeContext (1)
) throws RetryableExecutorException, NonRetryableExecutorException {

    // 标准分支模式:
    if (resumeContext.isEmpty()) { (2)
        // --- 1. 首次执行尝试 ---
        // 向下游系统派发异步请求,透传 idempotencyKey
        this.paymentGateway.initiateAsyncCharge(
                currentDomainEntity.getTotalAmount(),
                idempotencyKey
        );

        // 暂停前向执行,等待外部回调 Webhook:
        return stepManager.pause(
                DispatchDeliveryExecutor.class,
                "PAYMENT_INITIATED",
                Duration.ofHours(2)
        );
    }

    // --- 2. 恢复执行尝试 ---
    // 回调已到达!从只读 ResumeContext 中提取并校验载荷参数
    String paymentStatus = resumeContext.get("paymentStatus").orElse("FAILED"); (3)
    String gatewayRef = resumeContext.get("gatewayTransactionId").orElse(null);

    if ("SUCCESS".equalsIgnoreCase(paymentStatus)) {
        // 基于已验证的回调数据安全更新领域实体
        currentDomainEntity.setPaymentRef(gatewayRef);
        return stepManager.next(DispatchDeliveryExecutor.class, "PAYMENT_CONFIRMED"); (4)
    } else {
        // 前向执行失败并触发补偿回滚
        throw NonRetryableExecutorException
                .buildWith(new PaymentDeclinedException("Payment rejected: " + paymentStatus))
                .put("paymentStatus", paymentStatus)
                .build();
    }
}
1 同时接收 ResumeContext 与 currentDomainEntity、stepManager 和 idempotencyKey。
2 检查 resumeContext.isEmpty() 以区分初始执行(触发异步请求并暂停)与恢复执行(处理回调数据)。
3 通过 resumeContext.get(key) 安全查询上下文条目。该上下文在执行器内严格为只读模式;调用修改方法将抛出 UnsupportedOperationException。详见 ResumeContext 契约。
4 使用 stepManager.next(…​) 推进到下一个执行器 (DispatchDeliveryExecutor)。

完整生命周期请参阅 解耦外部回调载荷 与 回调控制器摄入模式。

doRevert(…​) 中的可暂停补偿(回滚等待状态)

在企业级微服务(尤其在银行、金融科技与电商领域)中,回滚操作通常无法同步完成。 例如,取消已授权交易通常需要向上游收单行发起异步退款,银行将在数小时或数天后通过 Webhook 通知系统结果。

StackSaga 通过 RevertStepManagerUtil.pause(…​) 支持在补偿事务中进行暂停。当暂停时: * 事务在事件存储库中保持 REVERTING 状态且 is_paused = 1。 * 该补偿步骤的 idempotencyKey 被注册为 current_resume_key。 * CPU 线程被释放,直至 Webhook 调用 sagaTemplate.resume(txId, correlationKey)。 * 恢复执行时,StackSaga 会携带填充好的 ResumeContext 重新调用 doRevert(…​)。

生产示例:异步支付退款补偿
@SagaExecutor(executeFor = "payment-service", value = "paymentCommandExecutor")
@AllArgsConstructor
public class PaymentCommandExecutor implements CommandExecutor<OrderDomainEntity> {

    private final PaymentGatewayClient paymentGatewayClient;

    @Override
    public ProcessStepManager<OrderDomainEntity> doProcess(
            OrderDomainEntity currentDomainEntity,
            ProcessStepManagerUtil<OrderDomainEntity> stepManager,
            String idempotencyKey
    ) throws RetryableExecutorException, NonRetryableExecutorException {
        // 前向同步扣款
        String paymentRef = this.paymentGatewayClient.charge(
                currentDomainEntity.getTotalAmount(),
                idempotencyKey
        );
        currentDomainEntity.setPaymentRef(paymentRef);
        return stepManager.next(ReserveStockExecutor.class, "PAYMENT_CHARGED");
    }

    @Override
    public RevertStepManager doRevert(
            NonRetryableExecutorException primaryExecutionException,
            OrderDomainEntity finalDomainEntityState,
            RevertHintStore revertHintStore,
            String idempotencyKey,
            RevertStepManagerUtil stepManager,
            ResumeContext resumeContext (1)
    ) throws RetryableExecutorException {

        // 检查是首次补偿执行还是恢复后的回调执行
        if (resumeContext.isEmpty()) { (2)
            // --- 1. 首次补偿尝试 ---
            // 向外部支付网关发起异步退款
            String refundSessionId = this.paymentGatewayClient.initiateAsyncRefund(
                    finalDomainEntityState.getPaymentRef(),
                    idempotencyKey (3)
            );
            revertHintStore.put("refundSessionId", refundSessionId);

            // 暂停补偿流程,等待退款确认 Webhook:
            return stepManager.pause("REFUND_INITIATED", Duration.ofHours(24)); (4)
        }

        // --- 2. 恢复补偿尝试 ---
        // Webhook 到达!从 ResumeContext 中摄入退款结果
        String refundStatus = resumeContext.get("refundStatus").orElse("FAILED"); (5)

        if ("SUCCESS".equalsIgnoreCase(refundStatus)) {
            // 逆向回滚完成:完成当前补偿步骤
            return stepManager.done("REFUND_COMPLETED"); (6)
        } else {
            // 临时网关延迟或待结算清算:触发重试
            throw new RetryableExecutorException("Refund processing in banking gateway, retrying verification");
        }
    }
}
1 实现接收 ResumeContext 的 6 参数重载 doRevert。
2 使用 if (resumeContext.isEmpty()) 区分初始退款请求与恢复后的 Webhook 结果处理。
3 将框架提供的 idempotencyKey 透传给外部退款 API。当网关向您的 Webhook 回调时,此完全相同的键将充当 correlationKey。
4 调用 stepManager.pause("REFUND_INITIATED", Duration.ofHours(24)) 将补偿流程置于预期的业务等待状态。
5 当 Webhook 到达并调用 sagaTemplate.resume(txId, correlationKey).put("refundStatus", …​).execute() 时,框架唤醒 Saga 并通过 ResumeContext 传递参数。
6 返回 stepManager.done("REFUND_COMPLETED") 确认该步骤补偿圆满完成。StackSaga 随后自动以逆序推进补偿上一个历史执行器。

查询执行器 (Query Executors)

如果某个原子执行仅包含主执行(无需补偿回滚),此类执行应在查询执行器 (Query Executor) 中实现。

查询执行器 仅包含一个用于执行主流程的方法。

查询执行器的典型示例:

  • 收集用户送货详情
    因为它不会对用户服务的数据库产生任何状态修改,是一项纯只读操作。

阻塞式查询执行器 (Blocking Query Executor)

@SagaExecutor(executeFor = "user-service", value = "chekUserDetailsExecutor") (1)
@AllArgsConstructor
public class ChekUserDetailsExecutor implements QueryExecutor<OrderDomainEntity> { (2)

    private final UserService userService;

    @Override (3)
    public ProcessStepManager<OrderDomainEntity> doProcess(
            OrderDomainEntity currentDomainEntity,
            ProcessStepManagerUtil<OrderDomainEntity> stepManager,
            String idempotencyKey
    ) throws RetryableExecutorException, NonRetryableExecutorException {

        try {
            (4)
            UserDetailDto userDetail = this.userService.getUserDetails(currentDomainEntity.getUsername());
            currentDomainEntity.setUserDetail(userDetail);

            return stepManager.next(InitializeOrderExecutor.class, "FETCHED_USER_DETAILS"); (5)
        } catch (FeignException.ServiceUnavailable unavailableException) {
            (6)
            throw RetryableExecutorException.of(unavailableException);
        } catch (FeignException.BadRequest badRequestException) {
            (7)
            throw NonRetryableExecutorException
                    .buildWith(badRequestException)
                    .put("time", LocalDateTime.now())
                    .put("reason", "BadRequest")
                    .build();
        }
    }
}
1 执行器使用 @SagaExecutor 注解声明为 Spring Bean,并配置目标服务名及全局唯一执行器标识。
2 执行器实现 QueryExecutor 接口以声明为查询执行器。
3 重写 doProcess 方法以实现主执行逻辑。
它接收领域实体的当前状态、步骤管理工具类以及框架生成的 idempotencyKey。由于查询是纯只读操作(天然具备幂等性),因此查询执行器无需专门捕获已执行异常。
4 在 doProcess 方法内实现主执行业务逻辑。在此处您可以按需读取或更新领域实体。它包含了截至目前先前所有执行器所做的全部状态变更。
5 如果执行成功,使用 stepManager.next 导航至下一个执行器。
需传入下一个目标执行器类以及用于在事件存储库中记录当前步骤成功事件的事件名称供应源。
此外还提供了 stepManager.complete 方法,用于在当前执行器执行完毕后成功结束整个长事务:
stepManager.complete("INITIALIZED_ORDER");
6 捕获可重试的临时资源不可用异常,通知 SEC 将事务保持在重试模式,并按调度配置重新执行。
您可以使用 RetryableExecutorException.of(Exception e) 包装原始异常。
如果抛出未包装的自定义异常,SEC 会将其视为不可重试异常,立即终止前向推进并逆向启动补偿。
7 捕获不可重试异常,通知 SEC 终止前向执行并立即启动逆序补偿执行。
您可以使用 NonRetryableExecutorException.buildWith(Exception e) 包装原始异常,并使用 put(String key, Object value) 附加异常元数据到事件存储库中。
这些元数据后续可以在命令执行器的 doRevert 方法中通过 RevertHintStore 提取。
如果在此处直接抛出未经包装的异常,SEC 也会在内部将其视为不可重试异常并触发补偿回滚。

非阻塞响应式查询执行器 (Non-Blocking Query-Executor)

@SagaExecutor(executeFor = "user-service", value = "reactiveChekUserDetailsExecutor")
@AllArgsConstructor
public class ReactiveChekUserDetailsExecutor implements ReactiveQueryExecutor<OrderDomainEntity> { (1)

    private final UserService userService;

    @Override (2)
    public Mono<ProcessStepManager<OrderDomainEntity>> doProcess(
            OrderDomainEntity currentDomainEntity,
            ProcessStepManagerUtil<OrderDomainEntity> stepManager,
            String idempotencyKey
    ) {
        (3)
        return this
                .userService
                .getUserDetails(currentDomainEntity.getUsername())
                .map(userData -> {
                    currentDomainEntity.setTel(userData.getPhoneNumber());
                    currentDomainEntity.setEmail(userData.getEmail());
                    currentDomainEntity.setAddress(userData.getEmail());
                    return stepManager.next(ReactiveInitializeOrderExecutor.class, "FETCHED_USER_DETAILS");
                })
                .onErrorResume(throwable -> {
                    if (throwable instanceof ResourceUnavailableException) {
                        return Mono.error(RetryableExecutorException.of(throwable));
                    } else {
                        return Mono.error(NonRetryableExecutorException.buildWith(throwable)
                                .put("time", String.valueOf(System.currentTimeMillis()))
                                .put("reason", "BadRequest")
                                .build());
                    }
                });
    }
}
1 执行器实现 ReactiveQueryExecutor<OrderDomainEntity> 接口以声明为响应式查询执行器。泛型参数指定为要在执行器中使用的领域实体类。
2 重写 doProcess 方法,以响应式流方式实现主执行逻辑。
3 主执行业务逻辑在响应式管道中执行。
异常不能传递到管道外部(例如命令式抛出),且方法绝对不能返回 null 或 Mono.empty()。这会导致未捕获的主执行异常,立即终止前向推进并启动逆向补偿回滚。

子执行器 (Sub Executors)

除了查询执行器与命令执行器之外,StackSaga 还提供了一种特殊的执行器类型——子执行器 (Sub Executor)。 它用于在主补偿事务之外,执行额外的辅助原子补偿事务。

如前所述,一个执行器内部只能封装一个原子事务。 该规则对主执行和补偿执行同等适用。 但有时,当某个补偿执行触发时,您可能还需要执行另一个伴随的额外执行动作。

例如,假定业务需求规定:当订单取消执行时,还必须向另一个独立的外部服务同步通知或做状态更新。 根据执行器的原子性原则,您不能在 doRevert 方法内部同时实现取消订单和通知第三方这两个操作,因为它们是两个完全独立的原子操作。 在这种场景下,子执行器正是为此而设计的解耦利器。 根据子执行执行的时机,子执行器分为两类:

  1. 前置补偿子执行器 (Sub-Before-Executors)

    • 如果子执行器必须在主补偿事务执行之前运行,应使用前置补偿子执行器。 根据业务需求,一个命令执行器可以配置任意数量的前置补偿子执行器。 SEC 将依次按序导航并执行它们。 参见代码实现。

  2. 后置补偿子执行器 (Sub-After-Executors)

    • 如果子执行器必须在主补偿事务执行之后运行,应使用后置补偿子执行器。 根据业务需求,一个命令执行器可以配置任意数量的后置补偿子执行器。 SEC 将依次按序导航并执行它们。 参见代码实现。

如果业务需要同时配置前置与后置补偿子执行器,完全支持。 当同时配置了前置和后置子执行器时,整体补偿执行的执行顺序如下:

首先,按配置顺序依次执行该命令执行器的所有前置补偿子执行器;完成前置子执行器后,接着执行命令执行器的默认补偿执行(主补偿); 主补偿执行完成后,最后执行配置的所有后置补偿子执行器。 下图清晰展示了前置子执行器、主补偿执行与后置子执行器之间的执行时序与协作关系:

Stacksaga 执行器

阻塞式前置补偿子执行器 (Blocking Sub-Before-Executor)

@SagaExecutor(executeFor = "order-service", value = "orderInitializeSubBeforeExecutor") (1)
@AllArgsConstructor
public class OrderInitializeSubBeforeExecutor implements RevertBeforeExecutor<OrderDomainEntity, InitializeOrderExecutor> { (2)

    private final OrderService orderService;

    @Override (3)
    public RevertBeforeStepManager<OrderDomainEntity, InitializeOrderExecutor> doProcess(
            OrderDomainEntity finalDomainEntityState,
            NonRetryableExecutorException nonRetryableExecutorException,
            RevertHintStore revertHintStore,
            RevertBeforeStepManagerUtil<OrderDomainEntity, InitializeOrderExecutor> stepManager,
            String idempotencyKey)
    throws RetryableExecutorException {
        try {
            this.orderService.preCancelOrder(finalDomainEntityState.getOrderId(), idempotencyKey);
            return stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_BEFORE"); (4)
        } catch (OperationAlreadyExecutedException alreadyExecutedException) {
            // 如果在先前的重试尝试中已执行,返回成功状态
            return stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_BEFORE");
        } catch (FeignException.ServiceUnavailable unavailableException) {
            throw RetryableExecutorException.of(unavailableException);
        }
    }
}
1 使用 @SagaExecutor 声明为 Spring Bean 并配置元数据(目标服务名与全局唯一标识)。
2 实现 RevertBeforeExecutor<A, C> 接口以声明为前置补偿子执行器。
第一个泛型参数为全局长事务使用的领域实体类,第二个泛型参数为在其主补偿执行之前触发的目标命令执行器类。
3 重写 doProcess 方法,其实现逻辑与命令执行器的 doRevert 方法一致。它接收 idempotencyKey 透传给下游微服务,并捕获重复执行异常以返回成功状态。
4 执行成功后,通过 stepManager.complete 导航至主(父级)doRevert 方法。
需提供用于在事件存储库中记录成功事件的事件名称供应源。
如果配置了多个前置子执行器,可以通过 stepManager.next 链式导航至下一个前置子执行器。

非阻塞响应式前置补偿子执行器 (Non-Blocking Sub-Before-Executor)

@SagaExecutor(executeFor = "order-service", value = "reactiveRevertBeforeExecutor") (1)
@RequiredArgsConstructor
public class ReactiveOrderInitializeSubBeforeExecutor implements ReactiveRevertBeforeExecutor<OrderDomainEntity, ReactiveInitializeOrderExecutor> { (2)

    private final ReactiveOrderService reactiveOrderService;

    @Override (3)
    @NonNull
    public Mono<RevertBeforeStepManager<OrderDomainEntity, ReactiveInitializeOrderExecutor>> doProcess(
            OrderDomainEntity domainEntity,
            NonRetryableExecutorException processException,
            RevertHintStore revertHintStore,
            RevertBeforeStepManagerUtil<OrderDomainEntity, ReactiveInitializeOrderExecutor> stepManager,
            String idempotencyKey
    ) {
        (4)
        return this.reactiveOrderService
                .doSomething(idempotencyKey)
                .thenReturn(stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_BEFORE"))
                .onErrorResume(throwable -> {
                    if (throwable instanceof OperationAlreadyExecutedException) {
                        return Mono.just(stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_BEFORE"));
                    } else if (throwable instanceof ResourceUnavailableException) {
                        return Mono.error(RetryableExecutorException.of(throwable));
                    } else {
                        revertHintStore.put("ReactiveOrderInitializeSubBeforeExecutor:FAILED", processException.getMessage());
                        return Mono.just(stepManager.complete("IGNORED"));
                    }
                });
    }
}
1 使用 @SagaExecutor 声明为 Spring Bean 并配置元数据。
2 实现 ReactiveRevertBeforeExecutor<A, C> 接口以声明为响应式前置补偿子执行器。泛型 A 为长事务领域实体类,泛型 C 为主命令执行器类。
3 重写 doProcess 方法,以响应式流方式实现前置补偿逻辑。
4 在 doProcess 内利用响应式管道执行前置补偿逻辑。
异常不能传递到管道外部,且方法不能返回 null 或 Mono.empty(),否则会被判定为非重试致命错误导致事务终止。
若存在不可重试异常,应捕获并持久化至 RevertHintStore,然后返回虚拟事件名称以避免事务异常终止。

当前置子执行器成功执行完毕且没有更多前置子执行器时,使用 stepManager.complete 导航至父级 doRevert 方法。
如果存在后续的前置子执行器,可以使用 stepManager.next 推进到链中的下一个前置子执行器。

阻塞式后置补偿子执行器 (Blocking Sub-After-Executor)

@SagaExecutor(executeFor = "order-service", value = "orderInitializeSubAfterExecutor") (1)
@AllArgsConstructor
public class OrderInitializeSubAfterExecutor implements RevertAfterExecutor<OrderDomainEntity, InitializeOrderExecutor> { (2)

    private final OrderService orderService;

    @Override (3)
    public RevertAfterStepManager<OrderDomainEntity, InitializeOrderExecutor> doProcess(
            OrderDomainEntity finalDomainEntityState,
            NonRetryableExecutorException processException,
            RevertHintStore revertHintStore,
            RevertAfterStepManagerUtil<OrderDomainEntity, InitializeOrderExecutor> stepManager,
            String idempotencyKey
    ) throws RetryableExecutorException {
        try {
            this.orderService.postCancelOrder(finalDomainEntityState.getOrderId(), idempotencyKey);
            return stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_AFTER"); (4)
        } catch (OperationAlreadyExecutedException alreadyExecutedException) {
            // 如果在先前的重试尝试中已执行,返回成功状态
            return stepManager.complete("REVERTED_ORDER_INITIALIZATION_SUB_AFTER");
        } catch (FeignException.ServiceUnavailable unavailableException) {
            throw RetryableExecutorException.of(unavailableException);
        }
    }
}
1 使用 @SagaExecutor 声明为 Spring Bean 并配置元数据。
2 实现 RevertAfterExecutor<A, C> 接口以声明为后置补偿子执行器。泛型参数分别为领域实体类与对应的主命令执行器类。
3 重写 doProcess 方法,以与命令执行器 doRevert 一致的方式实现后置补偿逻辑。
4 执行成功后,通过 stepManager.complete 结束当前后置补偿步骤。若存在多个后置子执行器,可通过 stepManager.next 推进到下一个。

非阻塞响应式后置补偿子执行器 (Non-Blocking Sub-After-Executor)

@SagaExecutor(executeFor = "order-service", value = "reactiveOrderInitializeSubAfterExecutor") (1)
public class ReactiveOrderInitializeSubAfterExecutor implements ReactiveRevertAfterExecutor<OrderDomainEntity, ReactiveInitializeOrderExecutor> { (2)

    private final ReactiveOrderService reactiveOrderService;

    @Override
    @NonNull (3)
    public Mono<RevertAfterStepManager<OrderDomainEntity, ReactiveInitializeOrderExecutor>> doProcess(
            OrderDomainEntity finalDomainEntityState,
            NonRetryableExecutorException processException,
            RevertHintStore revertHintStore,
            RevertAfterStepManagerUtil<OrderDomainEntity, ReactiveInitializeOrderExecutor> stepManager,
            String idempotencyKey
    ) {
        (4)
        return this.reactiveOrderService
                .doSomething(idempotencyKey)
                .thenReturn(stepManager.complete("EXECUTED_ORDER_INITIALIZATION_SUB_AFTER"))
                .onErrorResume(throwable -> {
                    if (throwable instanceof OperationAlreadyExecutedException) {
                        return Mono.just(stepManager.complete("EXECUTED_ORDER_INITIALIZATION_SUB_AFTER"));
                    } else if (throwable instanceof ResourceUnavailableException) {
                        return Mono.error(RetryableExecutorException.of(throwable));
                    } else {
                        revertHintStore.put("ReactiveOrderInitializeSubAfterExecutor:FAILED", processException.getMessage());
                        return Mono.just(stepManager.complete("IGNORED"));
                    }
                });
    }
}
1 使用 @SagaExecutor 声明为 Spring Bean 并配置元数据。
2 实现 ReactiveRevertAfterExecutor<A, C> 接口以声明为响应式后置补偿子执行器。
3 重写 doProcess 方法以实现响应式后置补偿逻辑。
4 在响应式管道中执行业务逻辑。
异常不能传递到管道外部,且方法不能返回 null 或 Mono.empty()。

如果后置补偿成功且没有更多后置执行器,使用 stepManager.complete 结束该命令执行器的整体补偿流程。
如果存在多个后置执行器,可使用 stepManager.next 导航至下一个。
此外,后置子执行器还支持通过 stepManager.pauseWithNext(…​) 或 stepManager.pauseWithComplete(…​) 暂停等待外部回调。

汇总对比 (Summary)

下列执行器允许抛出 可重试异常 (Retryable Executor Exceptions):

执行器类型 DoProcess() 方法 doRevert() 方法

查询执行器 (Query Executor)

✔

✔

命令执行器 (Command Executor)

✔

✔

前置补偿子执行器 (Revert Before Executor)

✔

后置补偿子执行器 (Revert After Executor)

✔

下列执行器允许抛出 不可重试异常 (Non-Retryable Executor Exceptions):

执行器类型 DoProcess() 方法 doRevert() 方法

查询执行器 (Query Executor)

✔

✖

命令执行器 (Command Executor)

✔

✖

前置补偿子执行器 (Revert Before Executor)

✖

后置补偿子执行器 (Revert After Executor)

✖

步骤管理器导航参考手册 (Step Manager Navigation Reference)

StackSaga 为每个执行阶段提供了专门的步骤管理器工具类 (Step Manager Utilities)。 这些工具类允许开发者以编程式精准控制 Saga 执行协调器 (SEC) 的流转,无需维护集中式静态路由表:

工具类接口 主推进方法 暂停方法 终态 / 链式方法

ProcessStepManagerUtil<DE>
(前向主流程)

next(Class nextStep, event)
推进至主工作流中的下一个执行器跨度。

pause(Class nextStep, event, duration)
在前向流程中进入预期的业务等待状态,等待外部回调。详见 前向暂停。

complete(event)
终态完成:在事件存储库中将整个 Saga 标记为已完成。

RevertStepManagerUtil
(补偿回滚流程)

不适用 (逆序回滚由 SEC 引擎根据历史事件自动判定推进)

pause(event, duration)
在回滚等待状态中暂停补偿,等待外部逆向 Webhook。详见 可暂停补偿。

done(event)
标记当前命令执行器的补偿圆满完成,驱动 SEC 逆向回滚上一历史步骤。

RevertBeforeStepManagerUtil<DE, E>
(前置补偿子流程)

next(Class nextSubBefore, event)
推进至链中的下一个前置补偿子执行器。

pauseWithNext(Class, event, duration) 或 pauseWithComplete(event, duration)
暂停前置补偿回滚,等待异步 Webhook。

complete(event)
完成前置子执行器链,将执行权移交至父级命令执行器的 doRevert()。

RevertAfterStepManagerUtil<DE, E>
(后置补偿子流程)

next(Class nextSubAfter, event)
推进至链中的下一个后置补偿子执行器。

pauseWithNext(Class, event, duration) 或 pauseWithComplete(event, duration)
暂停后置补偿回滚,等待异步 Webhook。

complete(event)
完成整个后置子执行器链,圆满结束该命令执行器的整体补偿。

有关这些工具类如何与生态其他模块集成的详细信息,请参阅 术语表中的步骤管理器、双流向可暂停架构 以及 通过 SagaTemplate 恢复执行。

执行器设计准则 (Guidelines for Creating Executors)

每个 Saga 执行器应当封装且仅封装单个原子事务。 这意味着您绝不能在同一个执行器内部实现多个原子事务。

实施这一限制的核心原因在于:执行器是 Saga 编排引擎 (SEC) 管理的最小重试单元。 如果一个执行器包含了多个原子事务且发生故障,SEC 无法识别具体是哪一个子事务失败了。 例如,如果一个执行器执行了三个原子事务,而在执行第三个事务时发生失败,SEC 重试该执行器将会重新执行第一和第二个事务;如果这些步骤不具备幂等性,将导致重复操作。

这种做法极易引发数据冗余、完整性破坏以及一致性冲突等数据异常。 此外,当需要执行补偿(回滚)时,SEC 缺乏足够的粒度来识别需要逆转哪一个原子事务,因为引擎将该执行器视为不可分割的单一原子单元。

执行类型划分指南 (Executions Classifying Tips)

在创建执行器时,必须决定该执行器是 命令执行器 (Command-Executor) 还是 查询执行器 (Query-Executor)。 以下从数据库状态变更视角整理的图表可作为设计参考:

总结原则:如果该原子执行是*纯只读操作*,应实现为*查询执行器*;如果该原子操作会对任何数据库产生*状态变更*,则必须实现为*命令执行器*。

操作类型 (Operation) 是否存在补偿 (Has a Revert) 执行器类型 (Executor Type)

C - 创建 (Create)

是 (YES)

命令执行器 (Command-Executor)

R - 读取 (Read)

否 (NO)

查询执行器 (Query-Executor)

U - 更新 (Update)

是 (YES)

命令执行器 (Command-Executor)

D - 删除 (Delete)

是 (YES)

命令执行器 (Command-Executor)

例如,让我们对订单结算示例中的各个执行操作进行归类:

执行操作 (Execution) 执行器类型 (Executor Type) 划分依据 (Reason)

收集用户送货详情

查询执行器 (Query-Executor)

提取用户数据不会对用户服务数据库造成任何状态变更。

初始化订单

命令执行器 (Command-Executor)

初始化订单后,如果后续任何原子事务失败,该订单必须被取消回滚。

执行预授权

命令执行器 (Command-Executor)

完成预授权后,如果后续任何原子事务失败,该预授权必须被解冻取消。

扣减更新库存

命令执行器 (Command-Executor)

库存扣减后,如果后续任何原子事务失败,预留库存必须原样恢复。

执行实际支付

命令执行器 (Command-Executor)

支付扣款后,如果后续任何原子事务失败,扣减金额必须执行退款。

在下单示例中,执行实际支付操作之后没有后续原子操作。 但从理论与规范角度出发,执行支付操作仍应在命令执行器中执行。 因为如果未来在支付之后扩展了新的原子流程,您必须具备并执行支付的补偿动作。

合并多个原子执行 (Combine Multiple Atomic Executions)

在 Saga 执行器中有两种合并多个原子操作的实现可能。 我们已知 StackSaga 中有两种原子执行:命令执行与查询执行。 根据具体业务场景,有时可以将多个查询执行合并到同一个执行器中。

在同一个执行器中执行多个只读原子操作可以有效减少事件溯源 (Event Sourcing) 的持久化开销。 因为在每个执行器执行完毕后,Saga 引擎都会将领域实体的最新状态作为一个新事件持久化到数据库中。 例如,如果您在同一个执行器中合并实现 3 个只读原子操作,就可以减少 2 次事件溯源的数据库写入开销。 否则,若将这 3 个操作拆分到 3 个独立执行器中,每次执行后都会更新一次事件存储库(共计 3 次)。

第一种方式 (First way):

stacksaga diagram combine multiple executor option 1.drawio

第二种方式 (Second way):

stacksaga diagram combine multiple executor option 2.drawio