SagaTemplate 与事件监听器 (SagaTemplate and Event Listeners)

SagaTemplate

SagaTemplate 是在编排服务中与 StackSaga 引擎 (SEC) 进行交互的核心入口点。 它提供了流畅的构建器 API (Fluent Builder API),用于初始化事务、配置执行模式(即发即弃 Fire-and-Forget 或触发并观察 Fire-and-Watch),并查询任何事务的当前状态与执行历史。

SagaTemplate 无缝支持命令式(阻塞,Blocking)与响应式(非阻塞,Non-Blocking)两种应用程序架构。

非响应式环境中的 SagaTemplate (SagaTemplate in a Non-Reactive Environment)

在传统的 Spring MVC 或阻塞式环境中,使用 SagaTemplate 可以同步发起 Saga 事务并同步检索事务快照:

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate; (1)

    public String handle(String username, double amount, List<String> items) {
        String transactionId = this.sagaTemplate
                (2)
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    {// 初始化用于 Saga 执行的领域实体及其属性
                        orderDomainEntity.setUsername(username);
                        orderDomainEntity.setTotalAmount(amount);
                        //...
                    }
                    return orderDomainEntity;
                })
                (3)
                .peek(orderDomainEntity -> {
                    log.info("transactionId {}:", orderDomainEntity.getTransactionId());
                })
                (4)
                .startWith(ValidateUserExecutor.class)
                (5)
                .fireAndForget()
                (6)
                .execute();
        log.info("已启动下单 Saga 事务,transactionId: {}", transactionId);
        return transactionId;
    }

    public void printCurrentState(String transactionId) {
        TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> state = this
                .sagaTemplate
                (7)
                .getCurrentState(transactionId)
                .fetch();

        TransactionCompleteStatus currentStatus = state.getCurrentStatus(); (8)
        log.info("事务 {} 当前状态: {}", transactionId, currentStatus);
        NavigableMap<Integer, ? extends SyncExecutionEvent<OrderDomainEntity>> executionHistory = state.getExecutionHistory(); (9)
        log.info("事务 {} 执行历史: {}", transactionId, executionHistory);
        OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity(); (10)
        log.info("当前领域实体状态: {}", currentDomainEntity); (11)
        ZonedDateTime startedDateTime = state.getStartedDateTime(); (12)
        log.info("事务 {} 启动于: {}", transactionId, startedDateTime.toLocalDateTime());
    }
}
1 注入针对自定义 DomainEntity 类型配置的 SagaTemplate。
2 使用 init() 提供用于创建初始 DomainEntity 的 Supplier。框架在初始化期间生成并赋予唯一的 transactionId。
3 使用 peek()(可选)在执行开始前检查或丰富已初始化的 DomainEntity。
4 使用 startWith() 指定启动事务流的初始 执行器 (Executor)。
5 使用 fireAndForget() 选择即发即弃执行模式。
6 调用 execute() 以阻塞方式发起事务。初始化后返回生成的 transactionId。
7 调用 getCurrentState(transactionId).fetch() 从事件存储中同步检索事务状态快照。
8 获取事务的总体完成状态 (TransactionCompleteStatus)。
9 获取按时间顺序排列的执行历史记录(NavigableMap<Integer, ? extends SyncExecutionEvent<DE>>),包含每个已执行跨度(主执行和补偿执行)的元数据与状态。
10 获取表示截至当前应用的所有状态变更的领域实体快照。
11 记录领域实体快照日志。
12 获取事务的启动时间戳。
要在不进行轮询的情况下获取实时状态更新,请实现 状态变更监听器。

响应式环境中的 SagaTemplate (SagaTemplate in a Reactive Environment)

在响应式环境(例如 Spring WebFlux)中,SagaTemplate 在整个生命周期阶段提供完全非阻塞的操作符。

以响应式方式执行 Saga 事务主要有两种模式:

  1. 即发即弃 (Fire-and-Forget) —— 启动事务并在 Saga 被接受后立即接收发出事务 ID 的 Mono<String>。Saga 在后台持续执行。当您只需要事务 ID 并依赖 getCurrentState(transactionId).fetchAsync() 或 EventListener 追踪进度时使用该模式。

  2. 触发并观察 (Fire-and-Watch) —— 启动事务并将实时状态转换观察为 Flux<TransactionState<DE, SyncExecutionEvent<DE>>>。非常适合向客户端流式传输实时执行事件(例如通过 Server-Sent Events 或 WebSockets)。

即发即弃 (Fire-and-Forget)

调用 .fireAndForget().executeAsync() 将事务提交给 StackSaga 引擎,并返回发出生成事务 ID 的 Mono<String>:

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate;

    public Mono<String> handle(String username, double amount, List<String> items) {
        return this.sagaTemplate
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    {// 初始化用于 Saga 执行的领域实体及其属性
                        orderDomainEntity.setUsername(username);
                        orderDomainEntity.setTotalAmount(amount);
                        //...
                    }
                    return orderDomainEntity;
                })
                .peek(orderDomainEntity -> {
                    log.info("transactionId {}:", orderDomainEntity.getTransactionId());
                })
                .startWith(ValidateUserExecutor.class)
                .fireAndForget()
                (1)
                .executeAsync()
                .map(transactionId -> "下单成功,事务 ID 为: " + transactionId);
    }
}
1 executeAsync() 返回一个 Mono<String>,一旦 Saga 初始化完成即发出事务 ID。事务继续在后台异步运行;使用 getCurrentState(transactionId).fetchAsync() 或 EventListener 跟踪其执行进度。

触发并观察 (Fire-and-Watch)

调用 .fireAndWatch(Duration waitTimeout).toFlux() 启动事务并返回 Flux<TransactionState<DE, SyncExecutionEvent<DE>>>,在每个执行器跨度执行完成后发出更新的 TransactionState 快照:

@GetMapping(value = "/order/watch", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>>> createOrderAndWatch() {
    return sagaTemplate
            .init(OrderDomainEntity::new)
            .peek(orderDomainEntity -> {
                log.info("已启动事务 ID: {}", orderDomainEntity.getTransactionId());
            })
            .startWith(ValidateUserExecutor.class)
            .fireAndWatch(Duration.ofSeconds(40)) (1)
            .toFlux()
            .doOnNext(state -> {
                log.info("收到状态变更事件: {}", state);
            });
}
1 waitTimeout 定义了订阅者在管道发出 WaitTimeExceededException 之前等待执行结果的最长等待持续时间。
关键行为与错误处理 (Key Behaviors and Error Handling)
  • 超时异常 (WaitTimeExceededException):如果执行耗时超过指定的 Duration waitTimeout,流式管道会通过 onError 以 WaitTimeExceededException 终止。这仅代表订阅流的结束 —— Saga 事务本身仍在后台继续异步执行。在订阅结束后,您仍可通过两种方式观察其进展:

    • 如果您注册了 状态变更监听器(TransactionEventListener 或 ReactiveTransactionEventListener),它将继续接收事件回调。

    • 或者,随时通过 getCurrentState(transactionId).fetchAsync() 主动查询最新状态。

  • 执行暂停异常 (TransactionPausedException):如果某一步骤遇到可重试异常(例如瞬时数据库异常或自定义 RetryableExecutorException),执行将被暂停并重新调度。由于 TransactionState 仅代表已完成的步骤状态,因此实时流订阅会暂停并通过 onError 发出 TransactionPausedException。这代表临时重试暂停,并非致命故障。一旦重试重新调度执行,状态通知将继续通过监听器推送,或者您可通过 getCurrentState(transactionId).fetchAsync() 查询最新状态。

恢复暂停的事务(可暂停长事务,Pausable LRT)

当 Saga 执行器参与 可暂停长事务 (Pausable LRT) 并将工作委托给外部异步服务(例如异步支付网关或后厨备餐人工确认)时,它通过返回 stepManager.pause(…​) 故意识别并暂停正向执行。 类似地,当补偿步骤必须通过异步网关或银行转账撤销事务时,它通过返回 stepManager.pause(…​) 暂停补偿流。有关概念背景,请参阅 双向流可暂停工作流。

在暂停的那一刻,事务进入故意的业务等待状态: * SagaDomainEntity 状态被持久化到事件存储中。 * es_transaction 将标记更新为 is_paused = 1 并记录 current_resume_key = <idempotencyKey>。 * CPU 工作线程立即被释放回到线程池。

一旦外部系统完成其异步任务,为了唤醒并继续执行事务,SagaTemplate 提供了 resume(…​) 操作符。

Fluent ResumeSpec 构建器 API (The Fluent ResumeSpec Builder API)

调用 sagaTemplate.resume(…​) 返回 SagaTemplate.ResumeSpec<DE> 流畅构建器,用于配置回调元数据并触发事务恢复:

// 使用流畅键值上下文条目恢复执行:
this.sagaTemplate
        .resume(transactionId, correlationKey)
        .put("status", "SUCCESS")
        .put("gatewayReference", "REF-987654")
        .execute(); // 阻塞式(Spring MVC)

// 响应式非阻塞恢复:
Mono<Void> resumeMono = this.sagaTemplate
        .resume(transactionId, correlationKey)
        .put("status", "SUCCESS")
        .put("gatewayReference", "REF-987654")
        .executeAsync(); // 响应式(Spring WebFlux)

ResumeSpec<DE> 暴露以下方法:

  • .put(String key, String value):向 ResumeContext 添加单个键值对。

  • .putAll(Map<String, String> map):将指定 Map 中的所有条目批量添加到 ResumeContext。

  • .withContext(ResumeContext resumeContext):提供或替换预构建的 ResumeContext(例如来自 DefaultResumeContext.of(map))。

  • .execute():同步阻塞,直到引擎恢复该事务(适用于 Spring MVC / 命令式环境)。

  • .executeAsync():返回一个 Mono<Void>,一旦恢复被处理即异步完成(适用于 Spring WebFlux / 响应式环境)。

还提供了一个便捷的重载快捷方式:sagaTemplate.resume(transactionId, correlationKey, resumeContext).execute()。

统一恢复:无缝处理正向与补偿流 (Uniform Resumption)

StackSaga 架构的一大突出优势在于:从编排器的视角来看,恢复操作是完全统一的:

  • 无论 Saga 是在正向执行 (doProcess) 期间暂停,还是在逆向补偿 (doRevert) 期间暂停,您调用的都是完全相同的方法 —— sagaTemplate.resume(transactionId, correlationKey)。

  • StackSaga 执行协调器 (SEC) 会自动查询事件存储,确定事务当前所处的阶段(PROCESSING 与 REVERTING),将领域实体快照还原到暂停检查点,重置 is_paused = 0 和 current_resume_key = NULL,并调度恢复执行:

    • 如果在正向流中暂停:推进到在 stepManager.pause(NextExecutor.class, …​) 中指定的目标 NextExecutor,携带传入的 ResumeContext 调用其 doProcess(…​)。

    • 如果在补偿流中暂停:重新调用已暂停 CommandExecutor 的 doRevert(…​),传入提供的 ResumeContext,使其完成冲正撤销操作。参见 实现可暂停补偿。

Webhook / 回调控制器模式 (Webhook / Callback Controller Pattern)

以下示例演示 Spring REST 控制器如何处理外部回调,涵盖正向业务确认与补偿性退款确认两种场景:

命令式回调控制器 (Spring MVC)
@Slf4j
@RestController
@RequestMapping("/orders/callbacks")
@RequiredArgsConstructor
public class OrderCallbackController {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate;

    /**
     * 恢复正向执行(例如:后厨备餐完成或 3D-Secure 扣款确认)。
     */
    @PostMapping("/order-ready")
    public ResponseEntity<Void> onOrderReadyCallback(@RequestBody OrderReadyCallbackDto callbackDto) {
        log.info("收到备餐完成回调,事务 ID: {} correlationKey: {}",
                callbackDto.getTransactionId(), callbackDto.getCorrelationKey());

        this.sagaTemplate
                .resume(callbackDto.getTransactionId(), callbackDto.getCorrelationKey())
                .put("kitchenStatus", callbackDto.getKitchenStatus())
                .put("readyTimestamp", callbackDto.getReadyTimestamp().toString())
                .execute();

        return ResponseEntity.ok().build();
    }

    /**
     * 恢复补偿执行(例如:支付网关异步退款确认)。
     */
    @PostMapping("/refund-confirmed")
    public ResponseEntity<Void> onRefundConfirmed(@RequestBody RefundWebhookDto refundDto) {
        log.info("收到退款 Webhook 回调,事务 ID: {} correlationKey: {}",
                refundDto.getTransactionId(), refundDto.getCorrelationKey());

        this.sagaTemplate
                .resume(refundDto.getTransactionId(), refundDto.getCorrelationKey())
                .put("refundStatus", refundDto.getRefundStatus())
                .put("gatewayRefundRef", refundDto.getGatewayRefundReference())
                .execute();

        return ResponseEntity.ok().build();
    }
}
响应式回调控制器 (Spring WebFlux)
@Slf4j
@RestController
@RequestMapping("/orders/callbacks")
@RequiredArgsConstructor
public class ReactiveOrderCallbackController {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate;

    @PostMapping("/refund-confirmed")
    public Mono<ResponseEntity<Void>> onRefundConfirmed(@RequestBody Mono<RefundWebhookDto> refundDtoMono) {
        return refundDtoMono
                .flatMap(dto -> {
                    log.info("收到响应式退款 Webhook 回调,事务 ID: {} key: {}",
                            dto.getTransactionId(), dto.getCorrelationKey());

                    return this.sagaTemplate
                            .resume(dto.getTransactionId(), dto.getCorrelationKey())
                            .put("refundStatus", dto.getRefundStatus())
                            .put("gatewayRefundRef", dto.getGatewayRefundReference())
                            .executeAsync()
                            .thenReturn(ResponseEntity.ok().<Void>build());
                });
    }
}

当调用 resume(…​) 时: . StackSaga 验证传入的 correlationKey 与 es_transaction 上注册的 current_resume_key 是否一致。 . 将序列化后的 ResumeContext JSON 存储到 es_transaction_execution_tryout.resume_context 中。 . 重置 is_paused = 0 并清空 current_resume_key = NULL。 . 还原带有纯净检查点状态的 OrderDomainEntity。 . 调用目标执行器的 doProcess(…​) 或 doRevert(…​),注入填充好的 ResumeContext。有关执行器如何使用此数据,请参阅 在执行器中消费 ResumeContext。

ResumeContext API 与规范契约

ResumeContext (org.stacksaga.api.ResumeContext) 是外部回调数据流入执行器的标准化契约:

  • 键值存储:实现 java.util.Map<String, String> 与 Serializable。

  • 安全提取 (get):提供 Optional<String> get(String key) 进行安全提取:

    String refundStatus = resumeContext.get("refundStatus").orElse("FAILED");
  • 不可变性保证:一旦传入执行器(doProcess 或 doRevert),ResumeContext 会被标记为只读。调用修改操作(put、remove、clear)将抛出 UnsupportedOperationException。这防止执行器意外修改回调审计记录。

  • 序列化与事件溯源:使用框架配置的 JSON 映射器 Bean 提供 prettyPrintAsString(ApplicationContext) 和 printAsString(ApplicationContext)。

  • 深度对比:提供 deepEquals(ResumeContext obj, ApplicationContext applicationContext) 进行结构化 JSON 相等性比对。

  • 事件存储持久性:被序列化到 es_transaction_execution_tryout 的 resume_context JSON 列中,为系统接收到的每个回调载荷提供完整的可审计性。参见 解耦外部回调载荷。

超时重新执行机制(未收到回调时)

当执行器带着最长持续时间暂停时(例如 stepManager.pause(DeliveryDispatchExecutor.class, "ORDER_PREPARING", Duration.ofHours(5))),应用程序期望外部回调在 5 小时到期之前触发 sagaTemplate.resume(…​)。

然而,在外部系统崩溃、网络中断或在配置的超时时间内未送达回调的最坏情况下:

  • 为什么框架会重新执行已暂停的执行器: 框架无法强行推进到下一个执行器 (DeliveryDispatchExecutor),因为下游业务操作从未获得确认。 此外,调用外部目标服务、核对订单或支付状态以及重新触发回调所需的逻辑完全位于该暂停执行器内部。 因此,StackSaga 的重试协调器会直接从暂停的执行器 (MakePaymentExecutor) 重新执行事务,即使其主执行逻辑在进入暂停前已经成功过。

  • 重新调用流转: 超时到期后,重试协调器接管停滞事务的所有权,安全地将 SagaDomainEntity 还原到进入该执行器之前的检查点,并再次调用 doProcess(…​)。这为执行器提供了查询外部系统状态(例如检查扣款是否成功)、重新发起外部调用或抛出 NonRetryableExecutorException 启动补偿的机会。

  • 对状态监听器的影响: 由于已暂停的执行器被再次执行,这会直接影响已注册的状态变更监听器。请参阅下文 可暂停长事务中的至少一次通知 中的详细指导。

观察事务状态变更 (Observing Transaction State Changes)

在通过 execute() 或 executeAsync() 发起事务后,StackSaga 引擎会顺序执行每个步骤(正向命令/查询与逆向补偿)。

有两种不同的机制来观察状态变化:

  1. 状态变更监听器 (State Change Listeners) —— 基于推送 (Push-based):在每个执行器步骤完成后,引擎自动通知您的监听器回调方法。最适合实时反应性动作,如发送用户通知、更新缓存或发布事件。

  2. SagaTemplate.getCurrentState() —— 基于拉取 (Pull-based):按需查询事件存储,以获取给定事务 ID 的最新聚合状态快照。最适合面向用户的状态查询和完成后的离线审计。

使用状态变更监听器(实时推流,Real-time)

状态变更监听器在每个执行器跨度完成后(包括正向执行与补偿),直接从引擎接收推送通知。

选择与您的应用程序架构匹配的监听器接口:

  1. TransactionEventListener —— 适用于命令式 / 阻塞环境。

  2. ReactiveTransactionEventListener —— 适用于响应式 / 非阻塞环境。

顺序执行与延迟影响 (Sequential Execution and Latency Impact):

onStateChanged 回调作为事务生命周期的一部分顺序执行。 具体而言,Saga 流程中的下一个跨度/执行器步骤仅在当前 onStateChanged 执行完成后才会被触发。

在 onStateChanged 内部执行重型计算、缓慢的 I/O 调用或耗时操作会直接阻塞事务的推进,增加事务的端到端延迟。

最佳实践: 保持 onStateChanged 轻量化(例如快速发布通知或记录状态日志)。如果需要执行耗时任务,请将工作转交至异步线程池(必要时使用 TransactionState 的快照或深拷贝),从而允许引擎立即推进事务的下一个跨度。

可暂停长事务 (Pausable LRT) 中的至少一次通知(重复调用风险)

在 可暂停长事务 (Pausable LRT) 中,如果执行器配置了超时时长 (Duration)(例如 5 小时)并暂停等待外部回调,监听器回调 (onStateChanged) 可能会为同一个执行器被调用多次:

  1. 第一次调用(初始成功): 当执行器成功完成其初始正向逻辑(例如向外部支付服务派发异步请求)并通过 stepManager.pause(…​) 暂停时,引擎调用 onStateChanged。

  2. 第二次调用(超时重新执行): 如果在配置的超时时间到期之前没有外部回调到达以恢复 Saga,框架将重新执行该暂停执行器。在重新执行成功后,引擎将为同一个执行器再次调用 onStateChanged。

至关重要的设计原则:

  • 切勿在监听器中放置关键业务逻辑或状态变更: 因为监听器回调遵循至少一次投递语义,并且在重试和超时恢复期间可以被重新触发,所以将非幂等业务逻辑、状态变更或外部扣款放置在 onStateChanged 内部会导致重复副作用。监听器应严格保留用于观察性任务(日志记录、指标采集或推送用户通知)。

  • 推荐模式 —— 使用 QueryExecutor: 如果您的工作流要求在某个执行器之后立即执行某项操作(例如生成内部收据、更新读取视图或记录审计日志)且该操作不需要补偿,切勿将其放置在 TransactionEventListener 中。 相反,请在 Saga 流程中将其建模为合法的 QueryExecutor(例如通过 stepManager.next(RecordAuditQueryExecutor.class, …​) 链接)。 作为一等 Saga 参与者,QueryExecutor 可保证:

    • 在事件存储中持久记录与跟踪。

    • 与 Saga 的正向推进保持确定性的一致执行。

    • 享有完整的幂等性保护以及在事务执行历史中的全链路可见性。

TransactionEventListener

实现 TransactionEventListener 以在命令式环境中监听事务事件:

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler implements TransactionEventListener<OrderDomainEntity> {

    @Override
    public void onStateChanged(TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> state) {
        TransactionCompleteStatus currentStatus = state.getCurrentStatus();
        log.info("当前事务状态: {}", currentStatus);
        NavigableMap<Integer, ? extends SyncExecutionEvent<OrderDomainEntity>> executionHistory = state.getExecutionHistory();
        log.info("事务执行历史: {}", executionHistory);
        OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity();
        log.info("当前领域实体状态: {}", currentDomainEntity);
        ZonedDateTime startedDateTime = state.getStartedDateTime();
        log.info("事务启动时间: {}", startedDateTime.toLocalDateTime());
    }
}
由于 onStateChanged 是同步阻塞的用户代码,引擎会在由 AbstractSchedulerProvider 提供的专用有界弹性调度器 (blockingExecutionScheduler) 上调用它。这隔离了监听器执行与引擎的核心响应式线程。线程上限、队列容量和 TTL 可通过 stacksaga.scheduler.blocking-execution.* 配置属性进行调优。参见 配置属性。

ReactiveTransactionEventListener

实现 ReactiveTransactionEventListener 以在响应式环境中监听事务事件:

@Slf4j
@Component
@RequiredArgsConstructor
public class ReactivePlaceOrderHandler implements ReactiveTransactionEventListener<OrderDomainEntity> {

    @Override
    public Mono<Void> onStateChanged(TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> transactionState) {
        return Mono
                .just(transactionState)
                .doOnNext(state -> {
                    TransactionCompleteStatus currentStatus = state.getCurrentStatus();
                    log.info("当前事务状态: {}", currentStatus);
                    NavigableMap<Integer, ? extends SyncExecutionEvent<OrderDomainEntity>> executionHistory = state.getExecutionHistory();
                    log.info("事务执行历史: {}", executionHistory);
                    OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity();
                    log.info("当前领域实体状态: {}", currentDomainEntity);
                    ZonedDateTime startedDateTime = state.getStartedDateTime();
                    log.info("事务启动时间: {}", startedDateTime.toLocalDateTime());
                })
                .flatMap(state -> {
                    // 模拟异步非阻塞处理
                    return Mono.delay(Duration.ofSeconds(1)).then();
                });
    }
}
与 TransactionEventListener 不同,ReactiveTransactionEventListener#onStateChanged 运行在共享的非阻塞 executionScheduler 上。返回的 Mono<Void> 必须完全非阻塞;任何阻塞调用都会使工作线程发生饥饿,并影响其他并发事务。并发度可通过 stacksaga.scheduler.execution.parallelism 进行调优。参见 配置属性。

使用 SagaTemplate.getCurrentState()(按需拉取,On-demand)

SagaTemplate.getCurrentState(String transactionId) 提供基于拉取的途径,通过事务 ID 访问任何事务的完整状态快照和执行历史。

调用 getCurrentState(transactionId) 返回一个 SagaTemplate.StateSpec<DE> 规范,该规范提供两个方法:

  • fetch() —— 同步检索,返回 TransactionState<DE, SyncExecutionEvent<DE>>。

    切勿在非阻塞 Reactor 或 Netty 事件循环线程上调用 fetch()。该方法会执行线程检查 (Schedulers.isInNonBlockingThread()),如果从非阻塞线程调用会抛出 IllegalStateException。在响应式应用中,请始终使用 fetchAsync()。

  • fetchAsync() —— 异步检索,返回 Mono<TransactionState<DE, SyncExecutionEvent<DE>>>。在 Spring WebFlux 和响应式事件循环线程中使用是完全安全的。

阻塞式状态检索 (Spring MVC / 命令式)

TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> state =
        this.sagaTemplate
                .getCurrentState(transactionId)
                .fetch();

TransactionCompleteStatus status = state.getCurrentStatus();
NavigableMap<Integer, ? extends SyncExecutionEvent<OrderDomainEntity>> history = state.getExecutionHistory();
OrderDomainEntity entity = state.getCurrentDomainEntity();

响应式状态检索 (Spring WebFlux / 响应式)

Mono<TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>>> stateMono =
        this.sagaTemplate
                .getCurrentState(transactionId)
                .fetchAsync();

stateMono.subscribe(state -> {
    log.info("事务状态: {}", state.getCurrentStatus());
});
与在执行推进过程中在内存中接收推送状态事件的状态变更监听器不同,getCurrentState(transactionId) 在每次调用时都会查询底层的事件存储。它专为按需查询而设计,而不是在活跃执行期间进行高频轮询。

StackSaga 团队建议在专用的 Handler 组件中组织 Saga 事务逻辑,该组件封装了事务发起、状态查询以及生命周期事件观察。

命令式处理程序模式 (Imperative Handler Pattern)

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler implements TransactionEventListener<OrderDomainEntity> {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate;

    public String handle(String username, double amount, List<String> items) { (1)
        return this.sagaTemplate
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    orderDomainEntity.setUsername(username);
                    orderDomainEntity.setTotalAmount(amount);
                    return orderDomainEntity;
                })
                .startWith(ValidateUserExecutor.class)
                .fireAndForget()
                .execute();
    }

    public CustomOrderStatusView getCurrentState(String orderId) { (2)
        TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> state = this
                .sagaTemplate
                .getCurrentState(orderId)
                .fetch();
        return new CustomOrderStatusView(state);
    }

    @Override
    public void onStateChanged(TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> transactionState) { (3)
        log.info("事务 {} 状态已更新为 {}",
                transactionState.getTransactionId(),
                transactionState.getCurrentStatus());
    }
}
1 handle 方法使用 SagaTemplate 发起 Saga 事务并返回生成的事务 ID。
2 使用 .getCurrentState(orderId).fetch() 按需检索当前事务状态。
3 实现 TransactionEventListener 的 onStateChanged 以处理实时生命周期通知。

响应式处理程序模式 (Reactive Handler Pattern)

@Slf4j
@Component
@RequiredArgsConstructor
public class ReactivePlaceOrderHandler implements ReactiveTransactionEventListener<OrderDomainEntity> {

    private final SagaTemplate<OrderDomainEntity> sagaTemplate;

    public Mono<String> handle(String username, double amount, List<String> items) { (1)
        return this.sagaTemplate
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    orderDomainEntity.setUsername(username);
                    orderDomainEntity.setTotalAmount(amount);
                    return orderDomainEntity;
                })
                .startWith(ValidateUserExecutor.class)
                .fireAndForget()
                .executeAsync();
    }

    public Mono<CustomOrderStatusView> getCurrentState(String orderId) { (2)
        return this.sagaTemplate
                .getCurrentState(orderId)
                .fetchAsync()
                .map(CustomOrderStatusView::new);
    }

    @Override
    public Mono<Void> onStateChanged(TransactionState<OrderDomainEntity, SyncExecutionEvent<OrderDomainEntity>> transactionState) { (3)
        return Mono.fromRunnable(() ->
                log.info("事务 {} 状态已更新为 {}",
                        transactionState.getTransactionId(),
                        transactionState.getCurrentStatus())
        );
    }
}
1 非阻塞事务发起,返回包含事务 ID 的 Mono<String>。
2 非阻塞状态查询,使用 .getCurrentState(orderId).fetchAsync()。
3 实现 ReactiveTransactionEventListener 的 onStateChanged,返回 Mono<Void> 实现非阻塞事件处理。