编排器配置手册 (Configuration)
每个 Worker 服务的应答均通过基于领域的专用*回调主题 (Callback Topic)* 流回编排器,并由 Spring for Apache Kafka 监听器容器 (Listener Container) 消费。 这些容器的分配与调优方式 — 连同声明式配置属性与可重写的 Provider Bean — 决定了编排器的回调吞吐量、响应延迟、系统资源占用以及故障隔离边界。 本页汇集了编排器端可配置的所有核心要素。
默认情况下,每个 SHARED_GLOBAL 事件管理器的回调主题都由*单个*全局回调容器统一消费,这能保持极低的消费者线程开销,并适用于大多数微服务。
但对于高吞吐或对延迟敏感的关键 Saga,可以为其分配独立的专属容器以避免发生队头阻塞 (Head-of-Line Blocking);而消费组协议 (groupType) 则决定了并行度是否受限于分区数量。
调优的核心在于将每个事件管理器的*隔离级别 (Isolation)、*通信协议 (Protocol) 和*并发度 (Concurrency)* 与其实际业务负载精确匹配。
本页内容从顶层核心决策逐步递进到细节配置:
-
回调主题与监听器模型 — 基于领域的回调主题、
listenerScope容器分配策略 (SHARED_GLOBAL、SHARED_GROUP、ISOLATED) 以及groupType协议。 详见 回调主题与监听器模型 (Callback Topic & Listener Models)。 -
配置属性参考 — 声明式的
stacksaga.kafka.orchestrator.与stacksaga.instance.配置项。 详见stacksaga-kafka-orchestrator-spring-boot-starter配置属性参考。 -
编程式配置与 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.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 中隔离监听器模型的架构:
该专用容器通过 @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 配置属性参考
| 配置属性 | 数据类型 | 默认值 | 说明描述 |
|---|---|---|---|
|
|
|
扫描标注有 |
|
|
|
消费所有 |
|
|
|
全局回调监听器容器的并发级别。对于 |
|
|
|
全局回调监听器容器是否在应用启动时自动开启消费。参见 [orchestrator_auto_start]。 |
|
|
|
支撑框架共享非阻塞 |
|
|
|
|
|
|
|
当达到 |
|
|
|
未使用的 |
|
|
|
|
|
|
|
应用优雅停机时,在强制销毁之前给予每个调度器处理飞行中/排队任务的等待超时时间(秒)。同时适用于 |
|
|
|
当前实例所属的集群名称。各组件仅在 |
|
|
|
当前实例所属的运维部署区域。用于标识事务发起地、向 Kafka 命令封装头注入区域路由元数据以及协调分布式重试。
+
NOTE: 默认为 |
|
|
|
当前实例所属的可用区 (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 应答载荷通过消息头携带类型标识,只能由该反序列化器进行精准重构。 |
|
注册自定义 |
即使您提供了自定义工厂,框架仍会在*内部强制执行少量关键设置*,以保证消费行为的不变性。这些设置会叠加在您的配置之上,您无需自行设置也无法强行覆盖:
当容器的 groupType 为… |
框架强制约束… |
|---|---|
|
在消费者工厂上强制设置 |
|
共享组消费者工厂*派生自*上述工厂,但移除了仅适用于消费者组的专有键 — 包括 |
相同的 Bean 重写范式同样适用于编排器上的其他 StackSaga-Kafka Provider Bean — 包括生产者工厂 (OrchestratorPayloadProducerFactoryProvider)、KafkaTemplate (OrchestratorKafkaTemplateProvider) 以及 Saga 执行调度器 (AbstractSchedulerProvider) — 每个 Bean 均通过 @ConditionalOnMissingBean 注册并支持相同方式的重写。
后续章节将提供专属指南。
|