端点监听器配置手册 (Endpoint Listener Configuration)

编排器向 Worker 发送的每条命令消息,都是通过 Spring for Apache Kafka 监听器容器 (Listener Container) 投递至您的 @SagaEndpoint 端点的。 因此,这些容器的分配与调优方式直接决定了 Worker 的吞吐量、响应延迟、系统资源消耗以及故障隔离边界 — 这是 Worker 端最具架构影响力的配置决策之一。

默认情况下,所有端点共享*单个*全局监听器容器。 这种模式能将消费者线程与连接开销降至最低,适合绝大多数常规微服务。 但由于该容器的消费者线程是全端点共享的,若单个端点调用耗时过长或流量过大,可能会垄断工作线程并延迟无关端点的消息处理(即队头阻塞 / Head-of-Line Blocking);反之,如果为每个端点都盲目创建专属容器,在流量不高时又会造成线程与网络连接的严重浪费。 同时,消费组协议带来了第二个维度的考量:传统消费者组协议将有效并行度死死限制在主题的分区数量内,而全新的 Kafka 4 共享组协议突破了该上限,允许纯粹通过提高并发度进行横向扩容。 调优监听器容器的核心在于将每个端点的*隔离级别 (Isolation)、*组协议 (Protocol) 和*并发度 (Concurrency)* 精确匹配至其实际业务负载 — 既不饿死核心高优端点,也不超配浪费硬件资源。

本页自顶向下全面解析这些调优杠杆:

  1. 监听器容器作用域 (listenerScope) — SHARED_GLOBAL、SHARED_GROUP 与 ISOLATED:端点如何在容器间分组聚合,以及各自选型时机。 详见 端点主题监听器模型 (Endpoint Topic Listener Models)。

  2. 消费组协议 (groupType) — 传统消费者组与 Kafka 4 共享组深度对比,及其对并行度与确认机制的影响。 详见 消费组协议 (groupType)。

  3. 配置属性参考 — 全局监听器与即时重试行为的声明式 stacksaga.kafka.worker.* 配置项。 详见 stacksaga-kafka-worker 配置属性参考。

  4. 编程式配置与 Bean 重写 — 当属性配置不足以满足需求时,如何替换框架 Provider Bean(如自定义 ConsumerFactory)。 详见 编程式配置(Bean 重写)。

端点主题监听器模型 (Endpoint Topic Listener Models)

初识 Kafka?一分钟速览核心术语

本页基于若干 Kafka 基础概念展开。如对以下术语感到生疏,请先浏览本摘要:

分区 (Partition)

主题被水平划分为一个或多个*分区*,Kafka 将消息分散存储在这些分区中。在传统协议下,每个分区在同一个组内*至多被一个消费者*读取,因此分区数量决定了该主题并行消费者的上限天花板。

消费者组 (Consumer group)

在经典 Kafka 协议下协同分担读取主题任务的一组消费者。Kafka 将每个分区指派给组内单个成员;若消费者数量多于分区数,多余的消费者将完全闲置。

共享组 (Share group / KIP-932)

Kafka 4 引入的*队列式 (Queue-style)* 消费协议。Broker 不再将整个分区绑定给特定消费者,而是按单条记录分发给当前空闲的消费者 — 因此并行度*不再受分区数限制*,只需增加消费者并发即可线性扩容。

偏移量与提交 (Offset & offset commit)

*偏移量*是标识消费组在分区中读取进度的游标位置。*提交*偏移量能持久化记录该进度,以便服务重启后从断点处继续消费。auto.offset.reset=earliest 决定在*尚无*已提交偏移量时从何处开始消费:从最早可用的历史消息开始(而非仅消费全新消息)。

监听器容器 (Listener container)

Spring for Apache Kafka 中运行消费者线程并将收到的消息分发给端点方法的底层组件。StackSaga 会根据您选择的监听器作用域自动构建并管理这些容器。

并发度 (Concurrency)

监听器容器并行运行的消费者线程数量。它能带来多大实际收益取决于底层采用的协议(参见上述共享组与传统消费者组差异)。

至少一次投递 (At-Least-Once delivery)

确保每条消息被投递*至少*一次的投递担保。消息仅在处理*完成后*才被确认,若在确认前进程崩溃,消息将被再次投递 — 这就是为什么业务处理逻辑必须保证*幂等性*(允许重复执行而无副作用)。

框架提供了三种监听器容器*作用域 (Scopes)*,用于在 Worker 应用中消费来自端点主题的消息,可通过 @SagaEndpoint 注解的 listenerScope 属性在每个端点上独立配置:

独立于作用域之外,每个容器均采用通过 groupType 属性指定的 Kafka 消费组协议进行消息拉取 — 参见下文的 消费组协议 (groupType)。

消费组协议 (groupType)

每个监听器容器严格运行一种 Kafka 消费协议,通过 @SagaEndpoint 的 groupType 属性指定:

  • GroupType.CONSUMER — 传统的 Kafka 消费者组 (Consumer-Group) 协议,底层由 Spring 的 ConcurrentMessageListenerContainer 支撑。 有效并行度受限于所消费主题的分区总数(一个分区最多由组内一个消费者线程处理,超出分区数量的线程将被闲置)。

  • GroupType.SHARE — Kafka 4 共享组 (Share-Group) 协议 (KIP-932),底层由 ShareKafkaMessageListenerContainer 支撑。 记录由 Broker 端按队列模式协同分发,因此消费者并行度*不再受分区数量限制* — 您可以将 concurrency 调高到远超分区数的数值以加速并发处理。 启用 GroupType.SHARE 必须使用开启了共享组特性的 Kafka 4.x Broker。

单个容器只能运行一种协议。 因此,共享同一容器的所有端点(无论是全部 SHARED_GLOBAL 端点,还是同一 SHARED_GROUP (containerId) 的所有成员)都必须声明*相同*的 groupType。 若出现协议不一致,应用在启动时将快速失败。

自动启动 (autoStart)

端点的监听器容器是否在应用启动时自动开启消费由 @SagaEndpoint 的 autoStart 属性控制(boolean 类型,默认为 true)。 若设置为 false,框架会注册该容器但不会自动启动它,消费只有在显式调用启动方法后才会开始(例如注入容器 Bean 并调用 start())。

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

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

SHARED_GLOBAL(全局共享)

若 @SagaEndpoint 配置为 listenerScope = WorkerListenerScope.SHARED_GLOBAL,框架会将该端点的主题绑定到由*所有*声明了 SHARED_GLOBAL 的端点共享的全局单一监听器容器上。 这是默认行为,也是绝大多数端点的理想之选:通过将所有端点汇聚在单个容器内,避免为每个端点创建单独容器,从而大幅节约消费者线程与内存开销。

下图展示了 stacksaga-kafka-worker 应用中端点主题的共享监听器模型架构:

stacksaga-kafka worker 中的共享监听器容器
  • 全局容器的*消费组协议*与*并发度*通过 stacksaga.kafka.worker.global-endpoint-topic-listener.group-type 与 stacksaga.kafka.worker.global-endpoint-topic-listener.concurrency 全局配置(参见 stacksaga-kafka-worker 配置属性参考)。默认组协议为 SHARE。

  • 全局容器使用从服务名派生的统一消费者组(格式形如 saga-ws-{serviceName}-…),其中 ws 代表 Worker 服务 (Worker Service)。

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

  • 无论采用何种 groupType,容器在监听器回调返回后立即确认/提交记录 — 无论 doProcess()/undoProcess() 内部发生何种情况,该回调始终设计为正常返回。NonRetryableExecutorException、RetryableExecutorException 与 JustRetryableExecutorException 均在内部被捕获并通过回调主题报告给编排器,而绝不会泄露至 Kafka 容器,因此可重试的业务失败本身不会促使 Kafka 重新投递该记录。对于 GroupType.SHARE,容器使用 Spring 的 ShareAckMode.IMPLICIT 隐式确认,从不调用共享组的 RELEASE/REJECT 动作 — 框架利用 GroupType.SHARE 的唯一诉求是突破分区上限的队列式并发,而非其消息重试语义。

  • Saga 步骤的“至少一次”处理完全由框架自身的重试系统保障,与 Kafka 是否重投递原始记录无关。JustRetryableExecutorException 在当前 Worker 进程内立即触发重试(参见 stacksaga-kafka-implementation/worker/worker-endpoints.adoc#exception-handling-reference)。而 RetryableExecutorException 则将失败上报给编排器,编排器将事务持久化为已暂停状态,并在后续通过完全独立的模块 — 环形协调器 (Ring Coordinator) — 重新调用该跨度,而非由 Kafka 重投递本记录。 在发生真实 Kafka 级别重投递时(例如监听器回调返回前容器崩溃),框架在每条消息上附带的基于跨度的*幂等键 (idempotency key)* 会确保安全防御。框架*不会*自动替您去重,利用该幂等键防止重复执行是业务开发者的责任。

SHARED_GROUP(分组共享)

若 @SagaEndpoint 配置为 listenerScope = WorkerListenerScope.SHARED_GROUP,该端点将与声明了相同 containerId 的*其他*端点共享专属容器。 这使您可以将特定的一组端点隔离到独立于全局容器之外的专属池中,通常适用于需要专属高并发或需要与非核心端点负载完全隔离的关键步骤。 例如,将 UserValidateEndpoint 与 InventoryCheckEndpoint 均配置为 containerId = "validation-group",它们将归入同一个专用容器中协同运行。

  • 此处 containerId 是*必填项* — 它是其他端点复用以加入该共享容器的组标识键。

  • 组内所有端点必须声明*相同*的 groupType(启动期强制校验)。

  • 组内所有端点必须声明*相同*的 autoStart(启动期强制校验) — 参见 自动启动 (autoStart)。

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

  • 偏移量管理、消息确认与投递语义遵循前文针对 SHARED_GLOBAL 所述的对应规则。

ISOLATED(完全隔离)

若 @SagaEndpoint 配置为 listenerScope = WorkerListenerScope.ISOLATED,框架将为该端点的主题创建独占的专用监听器容器,不与其他任何端点共享。 这为该特定端点提供了最极致的隔离性与对消费处理的完全自主控制 — 尤其适用于高吞吐或对延迟极敏感的主题,避免共享容器发生排队队头阻塞或消费者线程争抢,代价是稍高的系统资源消耗。

下图展示了 stacksaga-kafka-worker 应用中端点主题的隔离监听器模型架构:

stacksaga-kafka worker 中的隔离监听器容器
  • 专用容器命名为 {beanName}ListenerContainer,其中 {beanName} 为端点的 Spring Bean 名称。

  • 其 groupType、concurrency 与 autoStart 完全继承自该端点自身的 @SagaEndpoint 注解属性。

  • 偏移量管理、消息确认与投递语义同样遵循前文所述规则。

当 groupType = GroupType.CONSUMER 时,任何作用域的有效并行度上限仍受限于所消费主题的分区总数。 当需要将消费并行度扩展至超越分区数时,请选择 groupType = GroupType.SHARE。

stacksaga-kafka-worker 配置属性参考

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

stacksaga.kafka.worker.global-endpoint-topic-listener.group-type

GroupType

SHARE

支撑所有 SHARED_GLOBAL 端点的全局端点监听器容器所使用的 Kafka 消费组协议。SHARE 使用 Kafka 4 共享组协议 (KIP-932,队列式消费);CONSUMER 使用传统消费者组协议。参见 消费组协议 (groupType)。

stacksaga.kafka.worker.global-endpoint-topic-listener.concurrency

int

20

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

stacksaga.kafka.worker.global-endpoint-topic-listener.auto-start

boolean

true

支撑所有 SHARED_GLOBAL 端点的全局端点监听器容器是否在应用启动时自动开启消费。

非响应式端点的重试配置:主执行 (Primary execution)(进程内即时重试)

stacksaga.kafka.worker.retry.primary.max-attempts

int

5

最大执行尝试次数(初始尝试加上后续重试次数)。

stacksaga.kafka.worker.retry.primary.initial-interval

Duration

1s

首次重试前的初始等待时间间隔。

stacksaga.kafka.worker.retry.primary.max-interval

Duration

10s

重试尝试之间的最大等待间隔上限。

stacksaga.kafka.worker.retry.primary.multiplier

double

2.0

每次重试间隔递增的乘数(指数退避)。例如,初始间隔为 1s 且乘数为 2.0 时,等待间隔依次为 1s、2s、4s……直至达到最大间隔。

非响应式端点的重试配置:补偿执行 (Revert execution)(进程内即时重试)

stacksaga.kafka.worker.retry.revert.max-attempts

int

5

补偿执行的最大尝试次数(初始尝试加上后续重试次数)。

stacksaga.kafka.worker.retry.revert.initial-interval

Duration

1s

首次补偿重试前的初始等待时间间隔。

stacksaga.kafka.worker.retry.revert.max-interval

Duration

10s

补偿重试尝试之间的最大等待间隔上限。

stacksaga.kafka.worker.retry.revert.multiplier

double

2.0

每次补偿重试间隔递增的乘数(指数退避)。

上述 retry. 属性仅用于调优非响应式端点的*进程内即时重试(即支撑 JustRetryableExecutorException 的底层 Spring RetryTemplate)。 由 RetryableExecutorException 触发的*延期调度重试*则通过 stacksaga-database-support 单独配置。

编程式配置(Bean 重写)

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

首先是*端点消费者工厂 (Endpoint Consumer Factory)*:

自定义端点 ConsumerFactory

Worker 通过由 WorkerSagaPayloadConsumerFactoryProvider Bean 提供的 ConsumerFactory<String, SagaPayload> 消费来自端点主题的命令消息。 默认情况下,框架根据您现有的 Spring Boot Kafka 配置(通过 KafkaProperties.buildConsumerProperties() 解析 spring.kafka.*)构建该工厂,并插入 StackSaga 所需的两个核心反序列化器 — 键的 StringDeserializer 与值的 SagaPayloadDeserializer。

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

@Component
public class CustomEndpointConsumerFactoryProvider extends WorkerSagaPayloadConsumerFactoryProvider {

    private final KafkaProperties kafkaProperties;
    private final SagaPayloadDeserializer sagaPayloadDeserializer; (1)

    public CustomEndpointConsumerFactoryProvider(
            KafkaProperties kafkaProperties,
            SagaPayloadDeserializer sagaPayloadDeserializer) {
        this.kafkaProperties = kafkaProperties;
        this.sagaPayloadDeserializer = sagaPayloadDeserializer;
    }

    @Override
    protected ConsumerFactory<String, SagaPayload> 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)
                sagaPayloadDeserializer        (5)
        );
    }
}
1 注入框架自带的 SagaPayloadDeserializer Bean,无需自行实例化 — 它已被注册为 @Component(Bean 名称为 clientSagaPayloadDeserializer)并预先装配好了正确的 JsonMapper。
2 从 KafkaProperties.buildConsumerProperties() 开始构建是可选但推荐的做法,以便使您的 spring.kafka.consumer.* 配置继续生效。您也可以从空 Map 开始并显式设置每一项属性。
3 应用您所需的任何自定义消费者调优参数。
4 键反序列化器*必须*为 StringDeserializer — 消息键为以字符串传输的事务 ID。
5 值反序列化器*必须*为 SagaPayloadDeserializer — 除了重构命令载荷外,它还利用消息头提取端点依赖的幂等键、事务 ID 与 Saga 事件类型,替换为其他反序列化器将导致端点处理彻底崩溃。

注册自定义 WorkerSagaPayloadConsumerFactoryProvider 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 重写范式同样适用于 Worker 上的其他 StackSaga-Kafka Provider Bean — 包括响应生产者工厂 (WorkerSagaResponsePayloadProducerFactoryProvider) 与 KafkaTemplate (WorkerKafkaTemplateProvider)。 非响应式执行调度器与重试模板的定制方式完全相同,并在 非响应式端点调度器 与 为非响应式端点配置 RetryTemplate 中详细介绍。