StackSagaKafkaTemplate 与事件监听器 (Event-Listener)

StackSagaKafkaTemplate

StackSagaKafkaTemplate 是在编排器服务中访问 StackSaga-Kafka 引擎 (SEC) 全部核心能力的主入口。 它提供了用于发起全新 Saga 事务以及查询既有事务当前状态的核心 API。

StackSagaKafkaTemplate 同时支持响应式 (Reactive) 与非响应式 (Imperative / Blocking) 编程模型。如果您处于非响应式环境中,可以按如下方式使用:

非响应式环境中的 StackSagaKafkaTemplate

本示例展示了如何使用 StackSagaKafkaTemplate 启动用于处理下单业务的全新 Saga 事务,以及如何查询事务的当前状态:

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler {

    private StackSagaKafkaTemplate<OrderDomainEntity, PlaceOrderTopic> stackSagaKafkaTemplate; (1)

    public String handle(String username, double amount, List<String> items) {
        String transactionId = this.stackSagaKafkaTemplate
                (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(PlaceOrderTopic.DO_FETCH_USER_DETAILS, OrderEventManager.class)
                (5)
                .execute();
        log.info("Started place order saga with transactionId: {}", transactionId);
        return transactionId;
    }

    public void printCurrentState(String transactionId) {
        TransactionState<OrderDomainEntity> state = this
                .stackSagaKafkaTemplate
                (6)
                .getCurrentState(transactionId)
                .execute();

        TransactionCompleteStatus currentStatus = state.getCurrentStatus(); (7)
        log.info("Current status of transaction {}: {}", transactionId, currentStatus);
        NavigableMap<Integer, SagaExecutionEvent> executionHistory = state.getExecutionHistory();(8)
        log.info("Execution history of transaction {}: {}", transactionId, executionHistory);
        OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity(); (9)
        log.info("Current domain entity state: {}", currentDomainEntity); (10)
        ZonedDateTime startedDateTime = state.getStartedDateTime(); (11)
        log.info("Transaction {} started at: {}", transactionId, startedDateTime.toLocalDateTime());
    }
}
1 在类中自动装配 StackSagaKafkaTemplate。必须提供自定义的 DomainEntity 与自定义 Topic 类作为 StackSagaKafkaTemplate 的泛型参数。
2 使用 init 操作符,通过提供自定义领域实体的 Supplier 初始化事务。该方法用于创建事务的初始基准状态,并在领域实体中设置执行 Saga 所需的基础属性。
3 使用 peek 操作符(可选),在启动执行之前对已初始化的领域实体执行附加动作。该操作符的典型用途是查看初始化后生成的 transactionId,因为#transactionId 是由框架在通过 init 传递时自动生成并设置到领域实体中的#。因此,您可以在执行前捕获生成的 transactionId。
4 使用 startWith 操作符指定启动执行时应触发的首个主题,并提供与该事务关联的 EventManager 类。该操作符通过提供首个触发主题以及负责导航执行流的 EventManager 类来指定执行的起始点。
5 使用 execute 操作符以阻塞式方式启动事务执行。
6 使用 getCurrentState 操作符,通过传入 transactionId 获取事务的当前实时状态。无论在任何时间点,只要提供 transactionId 即可检索最新状态。在返回的 TransactionState 对象中可以查看事务状态的所有详细信息,例如当前流转状态、执行历史、领域实体最新快照、启动时间等。
7 从 TransactionState 对象中获取事务的当前终态完成状态。
8 从 TransactionState 对象中获取事务的执行历史记录。执行历史是一个可导航映射表 (NavigableMap),包含事务流程中所有已执行的端点(主流程与补偿流程)及其元数据(执行状态、完成时间、事件名称等)。
9 从 TransactionState 对象中获取领域实体的当前状态。这是应用了事务流中所有已执行端点修改后的领域实体最新快照。
10 记录打印领域实体的当前状态。
11 从 TransactionState 对象中获取事务的启动时间并记录日志。
另一种获取事务状态变更通知的高效途径是使用 事件监听器 (EventListener)。
请查阅 TransactionState 类的 Javadoc,了解事务状态对象中提供的全部高级特性与洞察。

响应式环境中的 StackSagaKafkaTemplate

如果您处于响应式开发环境中,可以利用模板提供的响应式操作符以非阻塞方式使用 StackSagaKafkaTemplate:

@Component
@RequiredArgsConstructor
class ReactivePlaceOrderHandler {
    private StackSagaKafkaTemplate<OrderDomainEntity, PlaceOrderTopic> stackSagaKafkaTemplate;

    public Mono<String> handle(String username, double amount, List<String> items) {
        return this
                .stackSagaKafkaTemplate
                .init(() -> {
                    OrderDomainEntity orderDomainEntity = new OrderDomainEntity();
                    {//使用 Saga 执行所需的属性初始化领域实体。
                        orderDomainEntity.setUsername(username);
                        orderDomainEntity.setTotalAmount(amount);
                        //...
                    }
                    return orderDomainEntity;
                })
                .peek(orderDomainEntity -> {
                    log.info("transactionId {}:", orderDomainEntity.getTransactionId());
                })
                .startWith(PlaceOrderTopic.DO_FETCH_USER_DETAILS, OrderEventManager.class)
                .executeAsync()
                .doOnNext(transactionId -> {
                    log.info("Started place order saga with transactionId: {}", transactionId);
                });
    }
    public Mono<Void> printAndCurrentState(String transactionId) {
        return this
                .stackSagaKafkaTemplate
                .getCurrentState(transactionId)
                .executeAsync()
                .doOnNext(state -> {
                    TransactionCompleteStatus currentStatus = state.getCurrentStatus();
                    log.info("Current status of transaction {}: {}", transactionId, currentStatus);
                    NavigableMap<Integer, SagaExecutionEvent> executionHistory = state.getExecutionHistory();
                    log.info("Execution history of transaction {}: {}", transactionId, executionHistory);
                    OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity();
                    log.info("Current domain entity state: {}", currentDomainEntity);
                    ZonedDateTime startedDateTime = state.getStartedDateTime();
                    log.info("Transaction {} started at: {}", transactionId, startedDateTime.toLocalDateTime());
                })
                .then();

    }
}
除了使用 executeAsync 代替 execute 以非阻塞流方式启动执行与异步获取当前状态之外,其余逻辑与阻塞模型完全一致。关于 StackSagaKafkaTemplate 各操作符的详细说明请参阅上文的 阻塞式示例。

事务状态变更监听器 (Transaction State Change Listeners)

通过调用 StackSagaKafkaTemplate 的 execute 或 executeAsync 方法将事务移交给 StackSaga-Kafka 引擎后,事务会在后台推进执行,并在每次步骤执行(主流程与补偿流程)完成后适时更新状态。 如果您希望持续监听事务状态的流转变化并据此执行伴随动作,框架提供了监听器机制。

根据运行环境提供两种监听接口:

  1. 面向非响应式环境的 KafkaTransactionEventListener。

  2. 面向响应式环境的 ReactiveKafkaTransactionEventListener。

获取事务状态变更通知在框架中有两种方式:

  1. 使用 StackSagaKafkaTemplate.getCurrentState 手动轮询。

  2. 使用 KafkaTransactionEventListener 或 ReactiveKafkaTransactionEventListener 实时被动接收回调。

实时回调与主动轮询的对比: 通过 StackSagaKafkaTemplate.getCurrentState 手动轮询事务更新是非常低效的。相比之下,已注册的监听器会在每个跨度完成后收到实时被动回调,无需任何轮询。例如在下单场景中,可以实时向用户推送订单状态变更通知。若需在事务完全结束后按需查询完整历史详情(例如用户查看历史订单详情页),则使用 StackSagaKafkaTemplate.getCurrentState。

Kafka 编排器中的状态获取开销: 由于 Kafka 编排器中的执行采用无状态处理(每个事件消息独立处理,跨度之间不保留常驻内存事务上下文),引擎在调用每个监听器回调之前,会在内部从事件存储库中调用 StackSagaKafkaTemplate.getCurrentState。该查询会产生 I/O 成本。请避免为同一个领域实体注册多个监听器;推荐每个领域实体仅声明一个统一监听器并在内部做业务状态分支。

顺序执行与延迟影响: onStateChanged 回调在每个跨度执行完毕后按顺序触发。请避免在 onStateChanged 内部执行重度计算或阻塞式 I/O,因为这会直接延缓主事务的向前推进。如果需要执行重度操作,请将其卸载至独立的异步线程池中处理。

至少一次投递与幂等性: 在重试或超时恢复期间,引擎可能会重新调用受影响的跨度,导致该步骤的 onStateChanged 被触发多次。监听器必须保持幂等,且严格作为只读观测机制(日志审计、推送通知、指标度量)。切勿利用监听器来执行关键的业务状态修改。

KafkaTransactionEventListener(非响应式监听器)

在非响应式环境中,可以实现 KafkaTransactionEventListener 接口(来自 org.stacksaga.api.listener 包)以监听事务事件:

@Slf4j
@Component
@RequiredArgsConstructor
public class PlaceOrderHandler implements KafkaTransactionEventListener<OrderDomainEntity> {
    @Override
    public void onStateChanged(TransactionState<OrderDomainEntity, AsyncExecutionEvent> state) {
        TransactionCompleteStatus currentStatus = state.getCurrentStatus();
        log.info("Current status of transaction: {}", currentStatus);
        NavigableMap<Integer, AsyncExecutionEvent> executionHistory = state.getExecutionHistory();
        log.info("Execution history of transaction: {}", executionHistory);
        OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity();
        log.info("Current domain entity state: {}", currentDomainEntity);
        ZonedDateTime startedDateTime = state.getStartedDateTime();
        log.info("Transaction started at: {}", startedDateTime.toLocalDateTime());
    }
}

ReactiveKafkaTransactionEventListener(响应式监听器)

在响应式环境中,可以实现 ReactiveKafkaTransactionEventListener 接口(来自 org.stacksaga.api.listener 包)以监听事务事件:

@Slf4j
@Component
@RequiredArgsConstructor
public class ReactivePlaceOrderHandler implements ReactiveKafkaTransactionEventListener<OrderDomainEntity> {
    @Override
    public Mono<Void> onStateChanged(TransactionState<OrderDomainEntity, AsyncExecutionEvent> transactionState) {
        return Mono
                .just(transactionState)
                .doOnNext(state -> {
                    TransactionCompleteStatus currentStatus = state.getCurrentStatus();
                    log.info("Current status of transaction: {}", currentStatus);
                    NavigableMap<Integer, AsyncExecutionEvent> executionHistory = state.getExecutionHistory();
                    log.info("Execution history of transaction: {}", executionHistory);
                    OrderDomainEntity currentDomainEntity = state.getCurrentDomainEntity();
                    log.info("Current domain entity state: {}", currentDomainEntity);
                    ZonedDateTime startedDateTime = state.getStartedDateTime();
                    log.info("Transaction started at: {}", startedDateTime.toLocalDateTime());
                })
                .flatMap(orderDomainEntityTransactionState -> {
                    //在此处模拟异步处理。
                    return Mono.delay(Duration.ofSeconds(1))
                            .then();
                });
    }
}

StackSaga 团队强烈建议在名为 Handler 的专用类中封装 Saga 事务相关的全部交互操作,集中负责启动、监听和查询 Saga 事务的状态。 以下代码展示了 Handler 的推荐标准架构设计:

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

    private StackSagaKafkaTemplate<OrderDomainEntity, PlaceOrderTopic> stackSagaKafkaTemplate;

    public String handle(String username, double amount, List<String> items) { (1)
        return this
                .stackSagaKafkaTemplate
                .init(() -> {
                  //....
                })
                .....
    }

    public CustomOrderStatusView getCurrentState(String orderId) { (2)
        return this
                .stackSagaKafkaTemplate
                .getCurrentState(orderId)
                ......

    }
    @Override
    public void onStateChanged(TransactionState<OrderDomainEntity, AsyncExecutionEvent> transactionState) { (3)
        //....
    }
}
1 handle 方法负责使用 StackSagaKafkaTemplate 并提供所需业务入参来发起 Saga 事务。
2 通过 StackSagaKafkaTemplate 查询并返回当前事务的定制化视图状态。
3 实现 KafkaTransactionEventListener 接口并重写 onStateChanged 方法,以集中监听事务生命周期的状态变迁事件。