主题与事件管理器 (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 类中声明每个主题时,请遵循以下命名规则:

  1. 以其针对的服务和动作命名主题 — 采用 {service-name}.{action} 形式(如 user-service.fetch-user-details、payment-service.make-payment),而非以 Saga 领域命名,因为该端点是 Worker 服务的核心能力,其他 Saga 领域完全可以复用它。参见 端点主题命名。

  2. 使用英文句点 . 作为分层分段分隔符(例如 order-service.initialize-order)。连字符 (-) 允许在分段*内部*使用(例如 user-service、fetch-user-details),但不能*作为*分段分隔符;严禁在任何位置使用下划线 (_)。

  3. 对于补偿主题,复用其所逆转的主执行主题的*相同*名称 — 例如 UNDO_INITIALIZE_ORDER 复用 order-service.initialize-order。它们通过 SagaEventType 以及 主题键 的正负符号进行严格区分。

  4. 切勿自行书写 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 主题名称

获取用户详情

QUERY_DO_ACTION

user-service.fetch-user-details

saga.user-service.fetch-user-details

初始化订单

COMMAND_DO_ACTION

order-service.initialize-order

saga.order-service.initialize-order

执行支付扣款

COMMAND_DO_ACTION

payment-service.make-payment

saga.payment-service.make-payment

库存更新

COMMAND_DO_ACTION

inventory-service.update-inventory

saga.inventory-service.update-inventory

取消订单

COMMAND_UNDO_ACTION

order-service.initialize-order

saga.order-service.initialize-order

退还款项

COMMAND_UNDO_ACTION

payment-service.make-payment

saga.payment-service.make-payment

释放库存

COMMAND_UNDO_ACTION

inventory-service.update-inventory

saga.inventory-service.update-inventory

主执行主题与其补偿主题解析为*相同*的实际 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 工厂参数返回路由裁决:
  • actionUtil.next(topic) 将工作流路由推进至下一步/主题。

  • actionUtil.complete() 当所有步骤均成功处理完毕时,结束长事务 (LRT)。

  • actionUtil.error(exception) 中断前向执行并启动补偿(回滚)流程。

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() 运行在回调监听器容器的处理线程上。在此处执行耗时操作会占用这些有限的线程,引发排队队头阻塞并降低全系统吞吐量。
Table 1. 速查表:onNext() 与 onNextRevert() 中的异常行为机制
抛出 / 返回自 即时效应 产生的状态流转

onNext() — 任何异常(包括 actionUtil.error(…​))

启动补偿回滚

IN_PROGRESS → FAILED → COMPENSATING

onNextRevert() — 任何异常

终止补偿流程

COMPENSATING → 补偿失败 (Compensation Failed)

从 onNext() 抛出的任何异常 — 或通过 actionUtil.error(…​) 返回的异常 — 均被框架视为不可重试故障,并立即将 Saga 状态置为 FAILED,触发逆向补偿序列。 这意味着 onNext() 可以被*主动用于*基于编排器端条件触发补偿:例如,在收到 Worker 应答后评估的业务规则判定该 Saga 不应继续向前推进。
类似地,从 onNextRevert() 抛出的任何异常都会导致补偿序列立即终止,事务被永久标记为补偿失败。 关于所有方法的完整异常行为明细,请参阅 异常处理参考手册。
每个 EventManager 均通过专用的*基于领域的专属回调主题*接收 Worker 应答,消费该主题的监听器容器由 @SagaEventManager 上声明的 listenerScope(SHARED_GLOBAL、SHARED_GROUP 或 ISOLATED)与 groupType 共同分配。关于回调主题命名规范、三种监听器作用域及消费组协议行为的完整说明,请参阅 回调主题与监听器模型。