Kafka Streams Rebalance Protocol(KIP-1071)详解:Broker 驱动的流式任务再均衡与配置迁移指南
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
导读
本指南围绕 Apache Kafka 4.2 引入的Streams Rebalance Protocol(KIP-1071)展开,系统讲解 Kafka Streams 应用如何从"客户端自行计算任务分配"切换到"由 Broker 上的 Group Coordinator 统一协调"的新模型。通过阅读本文,你将掌握该协议的启用方式、Broker/分组/客户端三级配置体系、内置与自定义 TaskAssignor 的注册与选择、拓扑校验与 NOT_READY 状态的行为,以及从经典协议离线迁移的完整操作步骤,并了解当前版本(4.2.x)的支持范围与已知限制。
背景:从客户端协调到 Broker 驱动的再均衡
在 Kafka Streams 的经典(classic)再均衡协议中,当组成员发生变化时,所有成员都要参与一次同步的全局再均衡:客户端计算新分配、互相协商,这会形成一个全局同步点,延长应用停机窗口。
Streams Rebalance Protocol 延续 KIP-848 的思想(KIP-848 将普通消费者的再均衡协调从客户端移到 Broker),通过 KIP-1071 将这一模型推广到 Kafka Streams 工作负载:任务分配不再由客户端在再均衡期间计算,而是由 Broker 持续计算。应用程序不再注册为普通消费者组,而是向 Broker 注册一个专用的streams group(流式组),由 Broker 管理并暴露协调流式应用实例所需的全部元数据,从而让 Kafka Streams 的协调与 KIP-848 引入的现代 Broker 驱动再均衡模型对齐,并提供带流式特有语义和元数据管理的专用组类型。
佐证:客户端侧存在独立的
StreamsGroupHeartbeatRequest/StreamsGroupHeartbeatResponse实现(见 StreamsGroupHeartbeatRequest.java 与 StreamsGroupHeartbeatRequest.json),请求发送逻辑封装在 StreamsGroupHeartbeatRequestManager.java,其测试见 StreamsGroupHeartbeatRequestManagerTest.java。
当前版本支持的能力
在当前发布版本中,Streams Rebalance Protocol 提供以下能力:
- 核心 Streams 组再均衡协议:通过
group.protocol=streams配置启用专用的流式再均衡协议。它把 streams group 与 consumer group 分离,并在 Broker 上提供流式特有的组成员生命周期与元数据管理。 - Sticky Task Assignor:内置一种在再均衡期间最小化任务迁移的基础任务分配策略,注册名为
sticky,是默认分配器。其实现位于 StickyTaskAssignor.java,name()返回"sticky",通过assign(GroupSpec, TopologyDescriber)计算分配,并在分配活跃任务与备援任务时优先保留各成员先前持有的任务(见该文件assignActive/assignStandby逻辑)。 - 自定义 TaskAssignor:Broker 可通过
group.streams.assignors配置自定义任务分配器,该配置接受内置分配器名称与自定义TaskAssignor实现的完全限定类名列表,第一项为默认分配器。单个组通过分组级配置streams.assignor.name按名称选择一个已注册的分配器;未设置时使用group.streams.assignors的第一项。自定义实现必须是线程安全的,因为单个实例会在一个 Broker 上的所有组之间共享。 - 交互式查询(Interactive Query, IQ)支持:IQ 操作与新的流式协议兼容。
- 新的 Admin RPC:
StreamsGroupDescribeRPC 提供与消费者组信息分离的流式特有元数据,可通过Admin接口访问(请求/响应实现见 StreamsGroupDescribeRequest.java 与 StreamsGroupDescribeResponse.java)。 - CLI 集成:可通过 bin/kafka-streams-groups.sh 脚本列出、描述和删除 streams group。该脚本本质上是启动
org.apache.kafka.tools.streams.StreamsGroupCommand的包装。 - Topology Description Plugin:Broker 可通过可插拔后端记录每个 streams group 处理拓扑的可读描述,用
group.streams.topology.description.plugin.class配置。记录的拓扑可通过Admin接口或kafka-streams-groups.sh --describe --topology查看。相关实现见 StreamsGroupTopologyDescriptionManager.java。 - 离线迁移:关闭所有成员并等待其
session.timeout.ms过期(或强制显式离开组)后,可将经典组转换为 streams group,也可将 streams group 转换回经典组。Broker 侧唯一保留的组数据是已提交的偏移量(committed offsets)。内部主题(changelog 与 repartition topic)会继续作为普通 Kafka 主题存在。 - 静态成员(Static Membership):使用
group.protocol=streams时,流式应用可配置group.instance.id(见 config-streams)。但对于没有持久化状态存储的拓扑,Kafka Streams 每次重启都会生成新的进程 ID(process ID),导致 Broker 重新计算组分配,从而在跨重启场景下实际上抵消了静态成员的收益。
当前版本不支持的能力
使用新协议时应避免以下尚未就绪的特性:
- 拓扑更新:如果拓扑发生显著变化(例如新增源主题或改变子拓扑数量),必须创建一个新的 streams group。
- 高可用分配器(High Availability Assignor):sticky 是唯一的内置分配器,且不支持机架感知分配;但通过
group.streams.assignors注册的自定义分配器可以实现机架感知。 - 预热任务(Warmup Tasks):与经典再均衡协议不同,预热任务不是分配器特性,而是由 Group Coordinator 注入到分配结果中。其好处是预热任务可以独立于所配置的分配器使用;但当前版本尚未实现预热任务支持。
- 正则表达式:不支持基于模式的(pattern-based)主题订阅。
- 在线迁移:经典协议与新流式协议之间不支持应用运行期间的在线组迁移。
- 自定义客户端提供者(Custom Client Supplier):使用自定义
KafkaClientSupplier时,只能提供 restore/global consumer、producer 和 admin client;启用 "streams" 组后无法提供 "main" consumer。
为什么使用 Streams Rebalance Protocol
与经典客户端驱动协议相比,Streams Rebalance Protocol 的关键优势:
- Broker 驱动协调:将任务分配逻辑集中到 Broker 而非客户端。这提供了来自单一协调点的、一致且权威的任务分配决策,降低了脑裂(split-brain)场景的可能性。
- 更快、更稳定的再均衡:通过移除全局同步点来缩短再均衡持续时间并减小其影响,最小化成员变更或故障期间的应用停机时间。
- 更好的可观测性:提供专门的指标和管理接口,将 streams 与 consumer groups 区分开,借助 Broker 侧的可观测性实现更清晰的故障排查。指标细节见 监控文档 的 group coordinator 监控一节。
启用协议
从 Apache Kafka 4.2 开始,新集群默认启用 Streams Rebalance Protocol。要使用该协议,Broker 与客户端都必须运行 Apache Kafka 4.2 或更高版本。
Broker 配置
协议在新 Apache Kafka 4.2 集群上默认启用。要在既有集群(升级到 4.2 之后)上启用该特性,或希望显式控制时,使用kafka-features.sh操作streams特性版本:
启用特性(streams.version=1):
bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature streams.version=1禁用特性(streams.version=0):
bin/kafka-features.sh --bootstrap-server localhost:9092 downgrade --feature streams.version=0客户端配置
在 Kafka Streams 应用配置中设置:
group.protocol=streams配置体系
Broker 配置
以下 Broker 配置控制 streams group 的行为。完整细节见 broker-configs。
group.coordinator.rebalance.protocols:已启用的再均衡协议列表。要启用 streams group,需在协议列表中包含"streams"。group.streams.session.timeout.ms:所有 streams group 的默认超时时间(若某个特定 streams group 未单独覆盖),用于在使用 streams group 协议时检测客户端故障。group.streams.min.session.timeout.ms:最小会话超时时间。group.streams.max.session.timeout.ms:最大会话超时时间。group.streams.heartbeat.interval.ms:提供给成员的心跳间隔默认值。group.streams.min.heartbeat.interval.ms:最小心跳间隔。group.streams.max.heartbeat.interval.ms:最大心跳间隔。group.streams.max.size:单个 streams group 可容纳的流式客户端最大数量。group.streams.num.standby.replicas:每个任务的备援副本默认数量。group.streams.max.standby.replicas:备援副本配置动态调整时的最大值上限。group.streams.initial.rebalance.delay.ms:新组(即此前为空的组)的首次再均衡延迟该时长,以便更多成员加入。group.streams.topology.description.plugin.class:StreamsGroupTopologyDescriptionPlugin实现的完全限定类名。未设置时,拓扑描述特性 处于禁用状态。group.streams.assignors:可供 streams group 使用的任务分配器列表,为内置分配器名称与自定义TaskAssignor实现完全限定类名的列表。第一项是未通过streams.assignor.name选择分配器的组的默认分配器。
分组级配置(Group Configuration)
资源类型为GROUP的配置可在DescribeConfigs与IncrementalAlterConfigsRPC 中使用,以动态覆盖特定组的默认 Broker 配置。可通过AdminJava 接口或bin/kafka-configs.sh工具设置。完整细节见 group-configs。
以下分组级配置可用于 streams group(源码定义见 GroupConfig.java):
streams.session.timeout.ms:使用 streams group 协议时检测客户端故障的超时时间。streams.heartbeat.interval.ms:提供给成员的心跳间隔。streams.num.standby.replicas:每个任务的备援副本数量。streams.initial.rebalance.delay.ms:组首次再均衡的延迟时长,以便更多成员加入。streams.assignor.name:本组使用的任务分配器名称,必须是 Broker 通过group.streams.assignors注册的分配器之一。未设置时,组使用group.streams.assignors的第一项。
示例:设置分组级配置
bin/kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type groups --entity-name wordcount \ --add-config streams.num.standby.replicas=1注意:在 streams 再均衡协议中,
session.timeout.ms、heartbeat.interval.ms和num.standby.replicas是分组级配置,在客户端设置时会被忽略。请使用如上所示的bin/kafka-configs.sh工具进行设置。
客户端(Streams)配置
所有 Kafka Streams 配置的完整细节见 kafka-streams-configs。
启用流式再均衡协议的客户端配置:
group.protocol:指示是否使用流式再均衡协议的标志。设置为streams以启用(默认值为classic)。
被忽略的配置
启用流式再均衡协议后,以下配置会被忽略:
acceptable.recovery.lagmax.warmup.replicasnum.standby.replicas(改用分组级配置)probing.rebalance.interval.msrack.aware.assignment.tagsrack.aware.assignment.strategyrack.aware.assignment.traffic_costrack.aware.assignment.non_overlap_costtask.assignor.class(分配发生在 Broker 侧;改用group.streams.assignors与streams.assignor.name)session.timeout.ms(改用分组级配置)heartbeat.interval.ms(改用分组级配置)
从源码结构看,分组级配置解析集中在 GroupConfig.java,它通过
ConfigDef.define(...)定义了上述各项(含streams.assignor.name、streams.num.standby.replicas等),并将未显式覆盖的项回落到GroupCoordinatorConfig中对应的 Broker 默认值。这意味着同一份配置代码同时支撑kafka-configs.sh的 GROUP 实体操作与 Broker 动态配置校验。
管理运维(Administration)
Admin API
使用Admin接口中的 "streams groups" 方法以编程方式管理 streams group。这些 API 大部分基于与消费者组 API 相同的实现。
与消费者组 API 的主要区别:
describeStreamsGroups使用DescribeStreamsGroupRPC,包含的信息与消费者组不同。- streams group 有一个额外状态
NOT_READY,且没有经典协议的遗留状态。 removeMembersFromConsumerGroup在本版本没有对应 API,因为它使用仅适用于经典消费者组的LeaveGroupRPC,该 RPC 不适用于 KIP-848 风格的组。
kafka-streams-groups.sh
新增的bin/kafka-streams-groups.sh工具用于操作 streams group。它取代了 streams group 场景下的bin/kafka-streams-application-reset.sh,可用于列出、描述和删除 streams group。详细用法见 kafka-streams-groups.sh 文档。该脚本通过kafka-run-class.sh启动org.apache.kafka.tools.streams.StreamsGroupCommand(见 kafka-streams-groups.sh)。
拓扑描述(Topology Description)
当在 Broker 上通过group.streams.topology.description.plugin.class配置了拓扑描述插件后,Group Coordinator 会记录每个 streams group 处理拓扑的可读描述,该描述由 Kafka Streams 客户端自动推送。可通过Admin#describeStreamsGroups或kafka-streams-groups.sh --describe --topology检索。推送/描述工作流、插件实现指南与故障排查见 Topology Description Plugin。
架构与工作原理
Streams Groups
协议引入了与消费者组并行的streams group概念。流式客户端使用专用的心跳 RPCStreamsGroupHeartbeat加入组、离开组,并向 Group Coordinator 更新其当前拥有的任务及客户端特有元数据。
Group Coordinator 对 streams group 的管理方式与消费者组类似:通过心跳响应持续更新组成员元数据,并在检测到变化时运行分配逻辑。Group Coordinator 引入了名为streams的新组类型,并为组元数据、拓扑元数据、组成员元数据引入了新的记录键与值类型。这些记录持久化在__consumer_offsets主题中。
一个组要么是 streams group,要么是 share group,要么是 consumer group,由使用对应 GroupId 发出的第一个心跳请求决定。
佐证:Group Coordinator 侧的实现分散在 group-coordinator/src/main/java/org/apache/kafka/coordinator/group/streams/ 目录,包括
StreamsGroup.java(组状态机)、StreamsGroupMember.java(成员元数据)、StreamsCoordinatorRecordHelpers.java(记录编解码,对应写入__consumer_offsets的记录键/值类型)、CurrentAssignmentBuilder.java与TargetAssignmentBuilder.java(当前/目标分配构建)等。
拓扑配置与校验
为了在流式客户端之间分配任务,Group Coordinator 使用拓扑元数据:该元数据在成员加入组时初始化,并持久化在 consumer offsets 主题中。
每当成员加入 streams group 时,第一个心跳请求包含拓扑元数据。该元数据将拓扑描述为一组子拓扑(subtopologies),每个子拓扑由唯一字符串标识符识别,并包含创建内部主题和分配任务所需的元数据。
拓扑校验与 NOT_READY 状态
在处理 streams group 心跳期间,Group Coordinator 可能检测到拓扑所需的 source/sink 主题或内部主题不存在,或它们的配置与拓扑成功执行所需配置不一致。这会触发 "topology configuration"(拓扑配置)过程,Group Coordinator 执行以下步骤:
- 检查所有已配置的 source 主题是否存在。
- 检查 "copartition groups"(协同分区组)是否满足——即所有应当协同分区的 source 主题确实协同分区。
- 根据 source 主题配置推导所有内部主题所需的分区数。
- 检查所有内部主题是否存在且配置正确。
如果任何 source 主题或内部主题缺失,组进入NOT_READY状态。在NOT_READY状态下,所有心跳照常处理(因此通常不应失败),但心跳响应中的状态会指示存在何种问题。组处于NOT_READY状态时,所有成员都会获得空分配。
佐证:拓扑元数据与内部主题管理在源码中有对应模块——TopologyMetadata.java 与 topics/ 目录下的
CopartitionedTopicsEnforcer.java(协同分区校验)、InternalTopicManager.java、ConfiguredTopology.java/ConfiguredSubtopology.java、ChangelogTopics.java/RepartitionTopics.java(changelog 与 repartition 内部主题推导)等类。
集中的分配配置
核心分配选项集中在 Broker 上配置,不依赖每个客户端的配置。这允许在不重新部署流式应用的情况下调优 streams group。Broker 侧引入的核心分配选项是num.standby.replicas,既可以全局配置在 Broker 上,也可以通过IncrementalAlterConfigs与DescribeConfigsRPC 动态配置到特定 streams group。
最近一次使用的分配配置存储在 Broker 上的组元数据中。这样,当分配配置被动态更改时,可以立即触发重新分配。
监控与指标
现有组指标已扩展,以区分 streams group 与 consumer group,并覆盖 streams group 的状态。完整细节见 streams groups metrics 的 group coordinator 监控一节。
按协议统计的组数量
按协议类型统计的组数量,其中协议列表通过protocol=streams变体扩展:
kafka.server:type=group-coordinator-metrics,name=group-count,protocol={consumer|classic|streams}按状态统计的 Streams 组数量
按状态统计的 streams group 数量:
kafka.server:type=group-coordinator-metrics,name=streams-group-count,state={empty|not_ready|assigning|reconciling|stable|dead}Streams 组再均衡
Streams group 再均衡传感器:
kafka.server:type=group-coordinator-metrics,name=streams-group-rebalance-rate kafka.server:type=group-coordinator-metrics,name=streams-group-rebalance-count从经典协议迁移
当前仅支持离线迁移。将 Kafka Streams 应用从经典协议迁移到流式再均衡协议:
- 关闭所有应用实例。
- 等待
session.timeout.ms过期,使组变为空(或强制显式离开组)。 - 更新应用配置,设置
group.protocol=streams。 - 重启应用实例。
Broker 侧唯一保留的组数据是已提交的偏移量。其他所有组元数据将在应用以新协议启动时重新创建。内部主题(changelog 与 repartition topic)将继续作为普通 Kafka 主题存在。
类似地,可以按相同流程将 streams group 转回经典组,只需设置group.protocol=classic。
警告:本版本不支持在线迁移(应用运行期间迁移)。在协议之间迁移时请规划维护窗口。
警告:由于离线迁移代码中存在一个严重的 Broker 侧缺陷(KAFKA-20254),官方建议在 4.2.0 中避免从经典组迁移到 streams group。新建的 streams group 不受影响。修复在 4.2.1 中可用。
小结
Streams Rebalance Protocol(KIP-1071)将 Kafka Streams 的任务协调从客户端迁至 Broker:以streams组类型、StreamsGroupHeartbeat专用 RPC、Broker 侧集中式分配配置(group.streams.assignors+streams.assignor.name)、NOT_READY拓扑校验状态以及独立的指标与管理接口为支柱。若你的集群与客户端均已升级到 Apache Kafka 4.2+,可在维护窗口内按离线迁移流程切换到该协议;在 4.2.0 上请留意 KAFKA-20254 的迁移缺陷,优先使用 4.2.1 或更高版本。结合本文引用的 group-coordinator streams 实现、StickyTaskAssignor 与 GroupConfig.java,可以进一步深入理解其底层原理。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考