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 引擎后,事务会在后台推进执行,并在每次步骤执行(主流程与补偿流程)完成后适时更新状态。
如果您希望持续监听事务状态的流转变化并据此执行伴随动作,框架提供了监听器机制。
根据运行环境提供两种监听接口:
-
面向非响应式环境的
KafkaTransactionEventListener。 -
面向响应式环境的
ReactiveKafkaTransactionEventListener。
获取事务状态变更通知在框架中有两种方式:
-
使用
StackSagaKafkaTemplate.getCurrentState手动轮询。 -
使用
KafkaTransactionEventListener或ReactiveKafkaTransactionEventListener实时被动接收回调。
|
实时回调与主动轮询的对比:
通过 Kafka 编排器中的状态获取开销:
由于 Kafka 编排器中的执行采用无状态处理(每个事件消息独立处理,跨度之间不保留常驻内存事务上下文),引擎在调用每个监听器回调之前,会在内部从事件存储库中调用 顺序执行与延迟影响:
至少一次投递与幂等性:
在重试或超时恢复期间,引擎可能会重新调用受影响的跨度,导致该步骤的 |
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();
});
}
}
推荐的 Handler 模式 (Recommended Handler Pattern)
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 方法,以集中监听事务生命周期的状态变迁事件。 |