编排器配置手册 (Configuration)

每个 Worker 服务的应答均通过基于领域的专用*回调主题 (Callback Topic)* 流回编排器,并由 Spring for Apache Kafka 监听器容器 (Listener Container) 消费。 这些容器的分配与调优方式 — 连同声明式配置属性与可重写的 Provider Bean — 决定了编排器的回调吞吐量、响应延迟、系统资源占用以及故障隔离边界。 本页汇集了编排器端可配置的所有核心要素。

默认情况下,每个 SHARED_GLOBAL 事件管理器的回调主题都由*单个*全局回调容器统一消费,这能保持极低的消费者线程开销,并适用于大多数微服务。 但对于高吞吐或对延迟敏感的关键 Saga,可以为其分配独立的专属容器以避免发生队头阻塞 (Head-of-Line Blocking);而消费组协议 (groupType) 则决定了并行度是否受限于分区数量。 调优的核心在于将每个事件管理器的*隔离级别 (Isolation)、*通信协议 (Protocol) 和*并发度 (Concurrency)* 与其实际业务负载精确匹配。

本页内容从顶层核心决策逐步递进到细节配置:

  1. 回调主题与监听器模型 — 基于领域的回调主题、listenerScope 容器分配策略 (SHARED_GLOBAL、SHARED_GROUP、ISOLATED) 以及 groupType 协议。 详见 回调主题与监听器模型 (Callback Topic & Listener Models)。

  2. 配置属性参考 — 声明式的 stacksaga.kafka.orchestrator. 与 stacksaga.instance. 配置项。 详见 stacksaga-kafka-orchestrator-spring-boot-starter 配置属性参考。

  3. 编程式配置与 Bean 重写 — 当单纯属性配置无法满足需求时,替换框架底层 Provider Bean(如自定义 ConsumerFactory)。 详见 编程式配置(Bean 重写)。

回调主题与监听器模型 (Callback Topic & Listener Models)

基于领域的专用回调主题 (Per-Domain Callback Topic)

在 stacksaga-kafka-orchestrator 中,为每个 DomainEntity 对应的 EventManager 均会创建一个专用主题,用于接收来自 Kafka Worker 端点的响应应答消息。

主题名称遵循如下命名规范:saga.callback.{appName}.{domainCallbackTopicSuffix},其中:

  • {appName} 是应用程序的服务名,全部小写并将任何非字母数字字符合并替换为 -。

  • {domainCallbackTopicSuffix} 从 @SagaEventManager 注解中原样追加。

例如,在应用名为 order-service 的微服务中,一个声明了 domainCallbackTopicSuffix = "place-order" 的 SHARED_GLOBAL 事件管理器将解析为真实的 Kafka 主题名:saga.callback.order-service.place-order。

domainCallbackTopicSuffix 在 EventManager 的 @SagaEventManager 注解中配置,如下所示:

@SagaEventManager(
        domainCallbackTopicSuffix = "place-order"
)

监听器容器作用域 (Listener Container Scopes)

尽管为每个 EventManager 创建了独立的回调主题,但实际*消费*这些回调主题的底层容器则是根据 @SagaEventManager 注解中配置的 listenerScope 进行分配的。 共有三种作用域:

独立于作用域之外,每个容器均采用通过 groupType 属性指定的 Kafka 消费组协议进行消息拉取 — GroupType.CONSUMER(经典消费者组,底层由 ConcurrentMessageListenerContainer 支撑)或 GroupType.SHARE(Kafka 4 共享组 / KIP-932,队列式消费,底层由 ShareKafkaMessageListenerContainer 支撑,其并行度不受分区数限制)。

单个监听器容器只能运行一种消费协议,因此共享同一个容器的所有事件管理器(无论是全部 SHARED_GLOBAL 管理器,还是指定 SHARED_GROUP 的所有成员)都必须声明*相同*的 groupType。 若协议不一致,应用将在启动时快速失败报错。启用 GroupType.SHARE 必须使用开启了共享组特性的 Kafka 4.x Broker。

容器是否在应用启动时自动启动由 SagaEventManagerListener 的 autoStart 属性控制(通过 isolatedExecutionListener/sharedGroupExecutionListener 声明;boolean 类型,默认为 true)。

单个容器具有统一的生命周期,因此共享容器的所有事件管理器(指定 SHARED_GROUP 的全部成员)必须声明*相同*的 autoStart。 若出现冲突,应用会在启动时抛出 ValidationException 快速失败。此限制不适用于 SHARED_GLOBAL,后者的自启行为由全局配置属性 stacksaga.kafka.orchestrator.global-callback-listener.auto-start 统一管控(参见 stacksaga-kafka-orchestrator-spring-boot-starter 配置属性参考)。

以下深入剖析各个监听器作用域的架构细节:

SHARED_GLOBAL(全局共享)

若 @SagaEventManager 配置为 listenerScope = OrchestratorListenerScope.SHARED_GLOBAL,框架会将该事件管理器的回调主题绑定到一个由*所有* SHARED_GLOBAL 事件管理器共享的全局单一回调监听器容器上。 这是默认行为,也是绝大多数 Saga 场景的推荐选择:它将所有回调主题汇聚到单一容器中,而非为每个编排器分别启动容器,从而大幅节约消费者线程与内存开销。

下图展示了 stacksaga-kafka-orchestrator 中共享监听器模型的架构流转:

stacksaga-kafka 编排器中的共享监听器容器
  • 全局容器的*消费组协议*与*并发度*通过 stacksaga.kafka.orchestrator.global-callback-listener.group-type 与 stacksaga.kafka.orchestrator.global-callback-listener.concurrency 全局统一配置(参见 stacksaga-kafka-orchestrator-spring-boot-starter 配置属性参考)。默认组协议为 SHARE。

  • 消费者组遵循命名约定:saga-os-{serviceName}-…,其中 os 代表*编排服务 (Orchestration Service)*。

  • 对于 GroupType.CONSUMER,监听器配置了 auto.offset.reset=earliest,以便服务重启且没有可用已提交偏移量时仍能消费历史消息。对于 GroupType.SHARE,该属性不适用 — 共享组消费者在 Broker 端管理偏移量 (KIP-932),而非传统的单消费者 auto.offset.reset。

  • 无论采用何种 groupType,容器在监听器回调返回后立即确认/提交应答记录 — EventPreparationManager.prepare() 在内部处理所有状态流转(成功推进、暂停/待重试或已终止事务),绝不会将异常泄露至 Kafka 容器,因此回调始终正常返回。 因此,可重试的业务失败本身不会促使 Kafka 重新投递该应答。 对于 GroupType.SHARE,容器使用 Spring 的 ShareAckMode.IMPLICIT 隐式确认,从不调用共享组的 RELEASE/REJECT 动作 — 框架利用 GroupType.SHARE 的唯一诉求是突破分区上限的队列式并发,而非其消息重试语义。

  • Saga 步骤的“至少一次 (At-Least-Once)”处理完全由框架自身的重试系统提供担保,与 Kafka 是否重投递原始应答无关。当 Worker 报告 RetryableExecutorException 时,编排器将事务持久化为已暂停状态;随后完全独立的模块 — 环形协调器 (Ring Coordinator) — 会在持有该事务令牌区间的编排器实例上重新调用停滞的跨度。 在发生真实 Kafka 级别重投递时(例如监听器回调返回前容器崩溃),框架在每条消息上附带的基于跨度的*幂等键 (idempotency key)* 会确保安全防御。框架*不会*自动丢弃重复消息,利用该幂等键防止重复执行是业务开发者的责任。

SHARED_GROUP(分组共享)

若 @SagaEventManager 配置为 listenerScope = OrchestratorListenerScope.SHARED_GROUP,其回调主题将由仅与声明了相同监听容器名称的*其他*事件管理器共享的专属容器消费。 当一组特定的事件管理器需要专属的高并发度,或需要与全局其他管理器的负载相互隔离时,此模式可将它们归入独立的容器池中。

该容器通过 @SagaEventManager 的 sharedGroupExecutionListener 属性进行配置:

@SagaEventManager(
        value = "placeOrderEventManager",
        domainCallbackTopicSuffix = "place-order",
        listenerScope = OrchestratorListenerScope.SHARED_GROUP,
        groupType = GroupType.SHARE,
        sharedGroupExecutionListener = @SagaEventManagerListener(
                listenerContainerName = "order-group", // 必填:其他管理器声明相同的组名以加入该共享容器
                concurrency = 10,
                autoStart = true
        )
)
  • 此处 listenerContainerName 是*必填项* — 它是其他事件管理器复用以加入该共享容器的组标识键。若留空则注册失败。

  • 所有加入成员必须声明*相同*的 groupType(启动期强制校验)。

  • 所有加入成员必须声明*相同*的 autoStart(启动期强制校验) — 共享容器具有统一生命周期,配置不一致将抛出 ValidationException 快速失败。

  • 共享容器的并发度为该组所有成员中声明的*最大并发度 (Maximum Concurrency)*。

ISOLATED(完全隔离)

若 @SagaEventManager 配置为 listenerScope = OrchestratorListenerScope.ISOLATED,框架将为该事件管理器的回调主题创建独占的专用回调监听器容器,不与其他任何管理器共享。 适用于具有极端严苛性能或高可用要求的核心 Saga(如超高吞吐或对延迟极敏感的支付工作流),避免共享容器发生排队队头阻塞或消费者线程争抢。 它提供了最极致的资源隔离,代价是占用稍多的系统资源。

下图展示了 stacksaga-kafka-orchestrator 中隔离监听器模型的架构:

stacksaga-kafka 编排器中的隔离监听器容器

该专用容器通过 @SagaEventManager 的 isolatedExecutionListener 属性进行配置:

@SagaEventManager(
        value = "placeOrderEventManager",
        domainCallbackTopicSuffix = "place-order",
        listenerScope = OrchestratorListenerScope.ISOLATED,
        groupType = GroupType.SHARE,
        isolatedExecutionListener = @SagaEventManagerListener(
                listenerContainerName = "", // 选填:默认为 {beanName}ListenerContainer
                concurrency = 5,
                autoStart = true
        )
)
  • listenerContainerName 为选填项;若留空,容器命名默认为 {beanName}ListenerContainer,其中 {beanName} 为事件管理器的 value 标识。

  • concurrency 与 autoStart 来自该 @SagaEventManagerListener 声明。

  • 消费者组遵循与 SHARED_GLOBAL 相同的 saga-os-{serviceName}-… 命名约定,并具有相同的 auto.offset.reset=earliest 及至少一次语义。

stacksaga-kafka-orchestrator-spring-boot-starter 配置属性参考

配置属性 数据类型 默认值 说明描述

stacksaga.kafka.orchestrator.domain-entity-scan

String[]

[]

扫描标注有 @SagaDomainEntity 注解的类的包名列表,以逗号分隔。例如:com.example.domain,com.example.anotherdomain。

stacksaga.kafka.orchestrator.global-callback-listener.group-type

GroupType

SHARE

消费所有 SHARED_GLOBAL 事件管理器回调主题的全局监听器容器所使用的 Kafka 消费组协议。SHARE 使用 Kafka 4 共享组协议 (KIP-932,队列式消费);CONSUMER 使用传统消费者组协议。详见 监听器模型。

stacksaga.kafka.orchestrator.global-callback-listener.concurrency

int

20

全局回调监听器容器的并发级别。对于 group-type: SHARE,可根据需要设置任意并发值,因为共享组不受分区数限制。对于 group-type: CONSUMER,该值不应超过绑定到全局容器的回调主题的分区总数。

stacksaga.kafka.orchestrator.global-callback-listener.auto-start

boolean

true

全局回调监听器容器是否在应用启动时自动开启消费。参见 [orchestrator_auto_start]。

stacksaga.kafka.orchestrator.scheduler.execution.parallelism

int

Runtime.getRuntime().availableProcessors()

支撑框架共享非阻塞 executionScheduler 的工作线程数 — 即用于驱动日常响应式/Saga 处理(如新事务调度派发和事务状态查询,以及运行 ReactiveKafkaTransactionEventListener#onStateChanged)的并行调度器(通过 Schedulers.newParallel("saga-exe", parallelism) 创建)。参见 事务状态变更监听器。

stacksaga.kafka.orchestrator.scheduler.blocking-execution.thread-cap

int

Runtime.getRuntime().availableProcessors() * 10

blockingExecutionScheduler 允许创建的最大线程上限 — 这是一个专用于调用阻塞式非响应式业务用户代码(如 KafkaTransactionEventListener#onStateChanged)的弹性有界调度器。

stacksaga.kafka.orchestrator.scheduler.blocking-execution.queued-task-cap

int

100000

当达到 thread-cap 线程上限且所有线程繁忙时,blockingExecutionScheduler 在拒绝新任务之前允许排队的最大任务容量。

stacksaga.kafka.orchestrator.scheduler.blocking-execution.ttl-seconds

int

60

未使用的 blockingExecutionScheduler 空闲工作线程在被回收前的存活时间(秒)。

stacksaga.kafka.orchestrator.scheduler.blocking-execution.daemon

boolean

false

blockingExecutionScheduler 的工作线程是否创建为守护线程 (Daemon Thread)。默认为 false,以防 JVM 在监听器执行尚未完成时非预期退出。

stacksaga.kafka.orchestrator.scheduler.graceful-shutdown-timeout-seconds

int

30

应用优雅停机时,在强制销毁之前给予每个调度器处理飞行中/排队任务的等待超时时间(秒)。同时适用于 executionScheduler 与 blockingExecutionScheduler。

stacksaga.instance.cluster

String

-

当前实例所属的集群名称。各组件仅在 cluster 匹配时建立互联通信,因此需要互联的所有 Master、Slave 和编排器节点必须配置完全一致的集群名。

stacksaga.instance.region

String

default

当前实例所属的运维部署区域。用于标识事务发起地、向 Kafka 命令封装头注入区域路由元数据以及协调分布式重试。 + NOTE: 默认为 default(区域代码 0)。如果配置了自定义区域(如 us-east-1),必须显式声明自定义的 SagaRegionResolver Spring Bean;否则应用启动时将抛出 ValidationException 快速失败。

stacksaga.instance.zone

String

-

当前实例所属的可用区 (Zone)。此处无额外功能影响,但建议遵循 StackSaga 规范统一配置。

编程式配置(Bean 重写)

上述配置属性涵盖了常规的声明式调优。 除此之外,StackSaga-Kafka 多个底层基础设施组件均作为*可重写的 Spring Bean* 暴露,供高级代码级定制。 每个组件均提供了带有 @ConditionalOnMissingBean 注解的框架内置默认实现,因此只需在应用中声明您自己的同类型 Provider Bean,即可无缝*替换*默认实现 — 无需额外配置任何属性开关。

本节详细介绍这些关键定制切入点。首先是*响应消费者工厂 (Response Consumer Factory)*:

自定义响应 ConsumerFactory

编排器通过由 OrchestratorResponsePayloadConsumerFactoryProvider Bean 提供的 ConsumerFactory<String, SagaResponsePayload> 消费 Worker 的应答消息。 默认情况下,框架根据您现有的 Spring Boot Kafka 配置(通过 KafkaProperties.buildConsumerProperties() 解析 spring.kafka.*)构建该工厂,并插入 StackSaga 所需的两个核心反序列化器 — 键的 StringDeserializer 与值的 SagaResponsePayloadDeserializer。

若需提供自定义工厂(例如微调抓取批量大小、安全认证或 spring.kafka.* 未暴露的高级客户端属性),继承抽象类 OrchestratorResponsePayloadConsumerFactoryProvider,实现 consumerFactory() 并将其注册为 Spring Bean:

@Component
public class CustomResponseConsumerFactoryProvider extends OrchestratorResponsePayloadConsumerFactoryProvider {

    private final KafkaProperties kafkaProperties;
    private final SagaResponsePayloadDeserializer sagaResponsePayloadDeserializer; (1)

    public CustomResponseConsumerFactoryProvider(
            KafkaProperties kafkaProperties,
            SagaResponsePayloadDeserializer sagaResponsePayloadDeserializer) {
        this.kafkaProperties = kafkaProperties;
        this.sagaResponsePayloadDeserializer = sagaResponsePayloadDeserializer;
    }

    @Override
    protected ConsumerFactory<String, SagaResponsePayload> consumerFactory() {
        Map<String, Object> props = new HashMap<>(kafkaProperties.buildConsumerProperties()); (2)
        // 您的自定义高级调优:
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 250); (3)
        props.put(ConsumerConfig.FETCH_MAX_BYTES_CONFIG, 5 * 1024 * 1024);
        return new DefaultKafkaConsumerFactory<>(
                props,
                new StringDeserializer(),            (4)
                sagaResponsePayloadDeserializer      (5)
        );
    }
}
1 注入框架自带的 SagaResponsePayloadDeserializer Bean,无需自行实例化 — 它已被注册为 @Component 并预先装配好了正确的 JsonMapper。
2 从 KafkaProperties.buildConsumerProperties() 开始构建是可选但推荐的做法,以便使您的 spring.kafka.consumer.* 配置继续生效。您也可以从空 Map 开始并显式设置每一项属性。
3 应用您所需的任何自定义消费者调优参数。
4 键反序列化器*必须*为 StringDeserializer — 消息键为以字符串传输的事务 ID。
5 值反序列化器*必须*为 SagaResponsePayloadDeserializer — StackSaga 应答载荷通过消息头携带类型标识,只能由该反序列化器进行精准重构。

注册自定义 OrchestratorResponsePayloadConsumerFactoryProvider Bean 会*完全替换*框架默认实现(因其标注了 @ConditionalOnMissingBean)。 因此,必须如上所示严格保留这两个必需的反序列化器 — 采用不兼容的反序列化器将导致应答反序列化彻底崩溃。

即使您提供了自定义工厂,框架仍会在*内部强制执行少量关键设置*,以保证消费行为的不变性。这些设置会叠加在您的配置之上,您无需自行设置也无法强行覆盖:

当容器的 groupType 为… 框架强制约束…

GroupType.CONSUMER

在消费者工厂上强制设置 auto.offset.reset=earliest。

GroupType.SHARE

共享组消费者工厂*派生自*上述工厂,但移除了仅适用于消费者组的专有键 — 包括 partition.assignment.strategy、enable.auto.commit、auto.commit.interval.ms、auto.offset.reset 以及 isolation.level(共享组不使用这些属性) — 并将共享确认模式强制设为*隐式确认* (ShareAckMode.IMPLICIT)。含义参见 监听器模型。

相同的 Bean 重写范式同样适用于编排器上的其他 StackSaga-Kafka Provider Bean — 包括生产者工厂 (OrchestratorPayloadProducerFactoryProvider)、KafkaTemplate (OrchestratorKafkaTemplateProvider) 以及 Saga 执行调度器 (AbstractSchedulerProvider) — 每个 Bean 均通过 @ConditionalOnMissingBean 注册并支持相同方式的重写。 后续章节将提供专属指南。