主题与事件管理器 (Topics and EventManager)
在定义了 领域实体 (Domain-Entity) 之后,下一步是描述 Saga 本身的执行拓扑:命名每个执行步骤的*主题 (Topics)* 以及在步骤之间进行动态路由编排的*事件管理器 (EventManager)*。本页涵盖主题的定义(包括命名规范与主题键)以及编排驱动前向流 (onNext()) 与补偿流 (onNextRevert()) 的 EventManager 实现。关于基于领域的回调主题如何消费与调优,在 配置手册 中单独说明。
StackSaga 主题 (StackSaga Topic)
StackSaga 主题是用于在 StackSaga 中标识执行端点的常量对象。
主题本质上就是用于向 Worker 服务发送命令消息的 Kafka 主题。
StackSaga 主题包含附加元数据,例如主题名称、主题类型(主执行或补偿执行)以及用于标识执行节点的跨度 (Span) 名称。
在 EventManager 中使用 StackSaga 主题来根据执行流转动态决定下一步应触发哪个主题。
以下单业务为例,在主执行流程中包含 4 个原子执行(跨度 / Spans):获取用户详情、初始化订单、执行支付以及更新库存;在补偿执行流程中包含 3 个原子补偿:取消订单、退款以及释放库存。因此共计 7 个跨度。以下是下单示例中使用的自定义主题类实现:
class PlaceOrderTopic extends AbstractTopic<PlaceOrderTopic> {(1)
(2)
protected PlaceOrderTopic(String topicName, float topicKey, SagaEventType sagaEventType, String targetService) {
super(topicName, topicKey, sagaEventType, targetService);
}
(3)
protected PlaceOrderTopic(String topicName, float topicKey, SagaEventType sagaEventType, String targetService, PlaceOrderTopic parent) {
super(topicName, topicKey, sagaEventType, targetService, parent);
}
(4)
//主执行主题 (Primary execution topics)
public static final PlaceOrderTopic DO_FETCH_USER_DETAILS = new PlaceOrderTopic("user-service.fetch-user-details", 1, SagaEventType.QUERY_DO_ACTION, "user-service");
public static final PlaceOrderTopic DO_INITIALIZE_ORDER = new PlaceOrderTopic("order-service.initialize-order", 2, SagaEventType.COMMAND_DO_ACTION, "order-service");
public static final PlaceOrderTopic DO_MAKE_PAYMENT = new PlaceOrderTopic("payment-service.make-payment", 3, SagaEventType.COMMAND_DO_ACTION, "payment-service");
public static final PlaceOrderTopic DO_INVENTORY_UPDATE = new PlaceOrderTopic("inventory-service.update-inventory", 4, SagaEventType.COMMAND_DO_ACTION, "inventory-service");
(5)
//补偿/回滚主题 (Revert/compensation topics)
public static final PlaceOrderTopic UNDO_INITIALIZE_ORDER = new PlaceOrderTopic("order-service.initialize-order", -2, SagaEventType.COMMAND_UNDO_ACTION, "order-service", DO_INITIALIZE_ORDER);
public static final PlaceOrderTopic UNDO_MAKE_PAYMENT = new PlaceOrderTopic("payment-service.make-payment", -3, SagaEventType.COMMAND_UNDO_ACTION, "payment-service", DO_MAKE_PAYMENT);
public static final PlaceOrderTopic UNDO_INVENTORY_UPDATE = new PlaceOrderTopic("inventory-service.update-inventory", -4, SagaEventType.COMMAND_UNDO_ACTION, "inventory-service", DO_INVENTORY_UPDATE);
}
| 1 | 继承 AbstractTopic 类创建自定义主题类。 |
| 2 | 重写用于主执行主题实例化的构造函数。topicName:Kafka 中使用的主题名称。详见 主题命名规范。topicKey:代表该主题的常量且(领域内)唯一的 float 浮点数值。详见 主题键规范。sagaEventType:主题类型。查询执行为 SagaEventType.QUERY_DO_ACTION,命令执行为 SagaEventType.COMMAND_DO_ACTION。targetService:负责执行该命令的目标服务名称,用于日志记录与调试取证。 |
| 3 | 重写用于补偿执行主题实例化的构造函数。 参数与主执行构造函数一致,额外增加了父级主题参数: parent:与该补偿主题相关联的主执行主题。用于指明主执行与补偿执行之间的逆向对应关系。例如在下单示例中,补偿主题 UNDO_INITIALIZE_ORDER 关联至主执行主题 DO_INITIALIZE_ORDER。 |
| 4 | 遵循约定,在自定义主题类中将主执行主题声明为 public static final 静态常量字段。 |
| 5 | 遵循约定,在自定义主题类中将补偿执行主题声明为 public static final 静态常量字段。 |
| 自定义主题类不是 Spring Bean。切勿为其添加任何 Spring 组件注解。 |
StackSaga-Kafka 中的主题命名规范 (Topic Name Specification)
端点主题由 Worker 应用程序拥有并定义,而不是编排器。
主题的存在是为了向 Worker 端点 (@SagaEndpoint) 输送消息,因此 Worker 是其权威事实来源。
在编排器端,在 AbstractTopic 类中声明相同的主题仅仅是为了在路由命令消息时进行*寻址* — 编排器向该主题发布消息,但绝不拥有或消费它。
有关端点主题的完整解析(归属权、服务限定命名最佳实践以及框架自动追加的 saga. 前缀),请参阅 Worker 页面上的 端点主题。
您在 AbstractTopic 类中声明的主题名称必须与 Worker 的 topicNameSuffix 完全匹配。
|
在 AbstractTopic 类中声明每个主题时,请遵循以下命名规则:
-
以其针对的服务和动作命名主题 — 采用
{service-name}.{action}形式(如user-service.fetch-user-details、payment-service.make-payment),而非以 Saga 领域命名,因为该端点是 Worker 服务的核心能力,其他 Saga 领域完全可以复用它。参见 端点主题命名。 -
使用英文句点
.作为分层分段分隔符(例如order-service.initialize-order)。连字符 (-) 允许在分段*内部*使用(例如user-service、fetch-user-details),但不能*作为*分段分隔符;严禁在任何位置使用下划线 (_)。 -
对于补偿主题,复用其所逆转的主执行主题的*相同*名称 — 例如
UNDO_INITIALIZE_ORDER复用order-service.initialize-order。它们通过SagaEventType以及 主题键 的正负符号进行严格区分。 -
切勿自行书写
saga.前缀 — 框架构建实际 Kafka 主题时会自动追加该前缀(例如order-service.initialize-order→saga.order-service.initialize-order)。若您*确实*自行包含了saga.前缀,框架会原样保留,不会重复追加。
StackSaga-Kafka 中的主题键规范 (Topic Key Specification)
主题键用于序列化与反序列化的高性能索引。
Kafka 消息头中不传输冗长的原始主题字符串,而是使用紧凑的主题键代表。因此,键一旦在生产系统中使用即不可变更。强烈建议在同一领域内为每个主题分配常量且唯一的 float 浮点数值。例如在下单示例中,DO_FETCH_USER_DETAILS 的主题键为 1,DO_INITIALIZE_ORDER 为 2,DO_MAKE_PAYMENT 为 3,DO_INVENTORY_UPDATE 为 4。对于补偿主题,推荐使用对应的负浮点数值以与主执行区分,例如 UNDO_INITIALIZE_ORDER 为 -2,UNDO_MAKE_PAYMENT 为 -3,UNDO_INVENTORY_UPDATE 为 -4。
小数主题键值(例如 1.1、1.2、-1.1)为未来的子执行 (Sub-Execution) 特性预留,该特性将允许向主执行或补偿执行挂载额外的前置/后置子步骤。在当前版本中请避免使用小数键,以免与未来能力产生冲突。
|
对照上述自定义主题类,下单示例中使用的主题名称、键值与实际 Kafka 主题映射表如下:
| 执行动作 | 主题类型 | 声明的主题名称 | 真实 Kafka 主题名称 |
|---|---|---|---|
获取用户详情 |
|
|
|
初始化订单 |
|
|
|
执行支付扣款 |
|
|
|
库存更新 |
|
|
|
取消订单 |
|
|
|
退还款项 |
|
|
|
释放库存 |
|
|
|
主执行主题与其补偿主题解析为*相同*的实际 Kafka 主题 — 例如 DO_INITIALIZE_ORDER 与 UNDO_INITIALIZE_ORDER 均映射至 saga.order-service.initialize-order。每个端点对应单一主题;框架通过 SagaEventType 和主题键的正负符号区分主执行命令与补偿命令,而非依赖不同主题名称。
|
事件管理器 (EventManager)
StackSaga-Kafka 支持基于业务条件和事务实时状态的完全运行时动态执行导航。
EventManager 正是承担此职责的核心组件。以下为 OrderDomainEntity 创建自定义 EventManager 的示例:
@SagaEventManager( (1)
value = "placeOrderEventManager", (2)
listenerScope = OrchestratorListenerScope.SHARED_GLOBAL, (3)
groupType = GroupType.SHARE,
domainCallbackTopicSuffix = "place-order" (4)
)
public class PlaceOrderEventManager extends AbstractEventManager<OrderDomainEntity, PlaceOrderTopic> { (5)
@Override
public Supplier<List<PlaceOrderTopic>> registerTopics() { (6)
return () -> List.of(
PlaceOrderTopic.DO_FETCH_USER_DETAILS,
PlaceOrderTopic.DO_INITIALIZE_ORDER,
PlaceOrderTopic.DO_MAKE_PAYMENT,
PlaceOrderTopic.DO_INVENTORY_UPDATE,
PlaceOrderTopic.UNDO_INITIALIZE_ORDER,
PlaceOrderTopic.UNDO_MAKE_PAYMENT,
PlaceOrderTopic.UNDO_INVENTORY_UPDATE
);
}
@Override
public @NonNull SagaPrimaryEventAction<PlaceOrderTopic> onNext( (7)
PlaceOrderTopic recentTopic,
OrderDomainEntity currentDomainEntityState,
SagaPrimaryEventActionUtil<PlaceOrderTopic> actionUtil
) {
if (recentTopic.equals(PlaceOrderTopic.DO_FETCH_USER_DETAILS)) {
currentDomainEntityState.getMetadata().put("navigated-do-initialize-order-at", LocalDateTime.now().toString());
return actionUtil.next(PlaceOrderTopic.DO_INITIALIZE_ORDER);
}
if (recentTopic.equals(PlaceOrderTopic.DO_INITIALIZE_ORDER)) {
currentDomainEntityState.getMetadata().put("navigated-do-make-payment-at", LocalDateTime.now().toString());
return actionUtil.next(PlaceOrderTopic.DO_MAKE_PAYMENT);
}
if (recentTopic.equals(PlaceOrderTopic.DO_MAKE_PAYMENT)) {
currentDomainEntityState.getMetadata().put("navigated-do-inventory-update-at", LocalDateTime.now().toString());
return actionUtil.next(PlaceOrderTopic.DO_INVENTORY_UPDATE);
}
if (recentTopic.equals(PlaceOrderTopic.DO_INVENTORY_UPDATE)) {
return actionUtil.complete();
}
return actionUtil.error(new IllegalStateException("Unexpected topic: " + recentTopic));
}
@Override
public void onNextRevert( (8)
PlaceOrderTopic recentExecutedTopic,
PlaceOrderTopic nextTopic,
OrderDomainEntity lastDomainEntityState,
NonRetryableExecutorException nonRetryableExecutorException,
RevertHintStore revertHintStore,
Supplier<NavigableMap<Integer, PlaceOrderTopic>> remainingReverts
) {
{//onNextRevert 方法中 revertHintStore 与 nextTopic 参数的示例用法
if (nextTopic.equals(PlaceOrderTopic.UNDO_MAKE_PAYMENT)) { (9)
revertHintStore.put("BEFORE_NOTE:UNDO_MAKE_PAYMENT", "Sample value before reverting UNDO_MAKE_PAYMENT");
}
if (nextTopic.equals(PlaceOrderTopic.UNDO_INITIALIZE_ORDER)) {
revertHintStore.put("BEFORE_NOTE:UNDO_INITIALIZE_ORDER", "Sample value before reverting UNDO_INITIALIZE_ORDER");
}
}
{//remainingReverts 与 recentExecutedTopic 的示例用法
if (recentExecutedTopic.equals(PlaceOrderTopic.UNDO_MAKE_PAYMENT)) { (10)
log.info("Remaining reverts after reverting UNDO_MAKE_PAYMENT: {}", remainingReverts);
}
if (recentExecutedTopic.equals(PlaceOrderTopic.UNDO_INITIALIZE_ORDER)) {
log.info("Remaining reverts after reverting UNDO_INITIALIZE_ORDER: {}", remainingReverts);
}
}
}
}
| 1 | 使用 @SagaEventManager 注解标记自定义 EventManager 类,将其注册为 Spring Bean。 |
| 2 | value:自定义 EventManager 的 Spring Bean 名称,用于按名称全局标识。 |
| 3 | listenerScope 与 groupType:配置如何消费该事件管理器的回调(应答)主题。listenerScope 选择容器分配策略 — OrchestratorListenerScope.SHARED_GLOBAL(汇入单个全局共享容器,默认之选)、OrchestratorListenerScope.SHARED_GROUP(与其他管理器共享具名容器)或 OrchestratorListenerScope.ISOLATED(专属独立容器)。对于 SHARED_GROUP 与 ISOLATED,容器的 concurrency 与 autoStart 通过 sharedGroupExecutionListener/isolatedExecutionListener 属性 (@SagaEventManagerListener) 配置;共享 SHARED_GROUP 容器的每个事件管理器必须声明*相同*的 autoStart,否则注册时抛出 ValidationException 快速失败。groupType 选择 Kafka 消费组协议 — GroupType.CONSUMER(传统消费者组)或 GroupType.SHARE(Kafka 4 共享组 / KIP-932,不受分区数限制的队列式消费)。详见 stacksaga-kafka-implementation/orchestrator/properties.adoc#topic_model_stacksaga_kafka_orchestrator。 |
| 4 | domainCallbackTopicSuffix:为同一领域相关主题提供通用后缀。框架内部利用此后缀构建该领域的*回调主题* — 用于接收来自 Kafka Worker 端点的响应应答消息。确切命名规则参见 基于领域的专用回调主题。 |
| 5 | 继承 AbstractEventManager 类创建自定义 EventManager,并提供自定义 DomainEntity 与 Topic 类作为泛型参数。 |
| 6 | 重写 registerTopics() 方法以注册事务中涉及的主题。该方法由框架在启动期调用,以注册与 CustomDomainEntity 相关的全量主题列表。需通过 Supplier 返回主题列表。 |
| 7 | 重写 onNext() 方法,根据最近刚完成的主题、领域实体的当前状态以及提供的 SagaPrimaryEventActionUtil (actionUtil) 决定下一步触发的主题。
该方法在主流程中每次执行成功后由框架自动调用,其首要职责是路由 (Routing) — 决定下一个触发哪个主题。
此处允许执行轻量级的编排器端薄记操作:例如在领域实体元数据中记录导航时间戳(如上所示),这是一项安全且低开销的操作,会被持久化到领域实体快照中以供追溯。
但是,业务处理、数据库调用、外部 HTTP 请求或任何阻塞式 I/O 绝不能在此处执行,因为该方法运行在回调监听器容器的工作线程上;阻塞这些线程将降低回调吞吐量,并可能阻塞共享该容器的其他 Saga。
使用 actionUtil 工厂参数返回路由裁决:
|
| 8 | 重写 onNextRevert()(可选)方法,以便在补偿流程的每个步骤执行附加动作。
调用时机:该方法在 recentExecutedTopic 成功完成其补偿之后、且在 nextTopic 被派发给 Worker 之前被调用。
因此两个参数同时可用 — 您既可以对刚完成的补偿做出响应,也可以为即将运行的下一个补偿准备上下文。 |
| 9 | onNextRevert() 方法中 RevertHintStore 与 nextTopic 参数的使用示例。
利用 nextTopic 判断即将派发的补偿主题,并利用 revertHintStore 存储后续补偿执行所需的任何上下文元数据。 |
| 10 | onNextRevert() 方法中 remainingReverts 与 recentExecutedTopic 参数的使用示例。
利用 recentExecutedTopic 标识刚完成的补偿步骤,利用 remainingReverts 审查仍挂起待执行的完整补偿步骤集合。 |
nextTopic 参数是为了即将推出的子执行 (Sub-Execution) 特性而引入的,该特性将允许向补偿执行挂载额外的前置/后置子步骤。在当前版本中,它反映了下一个顶层补偿主题。
|
作为生产最佳实践,避免在事件导航器内部执行任何 I/O 密集型或高 CPU 消耗操作。这些方法应仅用于根据传入参数计算评估路由条件。onNext() 与 onNextRevert() 运行在回调监听器容器的处理线程上。在此处执行耗时操作会占用这些有限的线程,引发排队队头阻塞并降低全系统吞吐量。
|
| 抛出 / 返回自 | 即时效应 | 产生的状态流转 |
|---|---|---|
|
启动补偿回滚 |
|
|
终止补偿流程 |
|
从 onNext() 抛出的任何异常 — 或通过 actionUtil.error(…) 返回的异常 — 均被框架视为不可重试故障,并立即将 Saga 状态置为 FAILED,触发逆向补偿序列。
这意味着 onNext() 可以被*主动用于*基于编排器端条件触发补偿:例如,在收到 Worker 应答后评估的业务规则判定该 Saga 不应继续向前推进。类似地,从 onNextRevert() 抛出的任何异常都会导致补偿序列立即终止,事务被永久标记为补偿失败。
关于所有方法的完整异常行为明细,请参阅 异常处理参考手册。
|
每个 EventManager 均通过专用的*基于领域的专属回调主题*接收 Worker 应答,消费该主题的监听器容器由 @SagaEventManager 上声明的 listenerScope(SHARED_GLOBAL、SHARED_GROUP 或 ISOLATED)与 groupType 共同分配。关于回调主题命名规范、三种监听器作用域及消费组协议行为的完整说明,请参阅 回调主题与监听器模型。
|