news 2026/9/10 18:26:11

Kafka Streams Rebalance Protocol(KIP-1071)详解:Broker 驱动的流式任务再均衡与配置迁移指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Kafka Streams Rebalance Protocol(KIP-1071)详解:Broker 驱动的流式任务再均衡与配置迁移指南

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 RPCStreamsGroupDescribeRPC 提供与消费者组信息分离的流式特有元数据,可通过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.classStreamsGroupTopologyDescriptionPlugin实现的完全限定类名。未设置时,拓扑描述特性 处于禁用状态。
  • group.streams.assignors:可供 streams group 使用的任务分配器列表,为内置分配器名称与自定义TaskAssignor实现完全限定类名的列表。第一项是未通过streams.assignor.name选择分配器的组的默认分配器。

分组级配置(Group Configuration)

资源类型为GROUP的配置可在DescribeConfigsIncrementalAlterConfigsRPC 中使用,以动态覆盖特定组的默认 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.msheartbeat.interval.msnum.standby.replicas分组级配置,在客户端设置时会被忽略。请使用如上所示的bin/kafka-configs.sh工具进行设置。

客户端(Streams)配置

所有 Kafka Streams 配置的完整细节见 kafka-streams-configs。

启用流式再均衡协议的客户端配置:

  • group.protocol:指示是否使用流式再均衡协议的标志。设置为streams以启用(默认值为classic)。

被忽略的配置

启用流式再均衡协议后,以下配置会被忽略:

  • acceptable.recovery.lag
  • max.warmup.replicas
  • num.standby.replicas(改用分组级配置)
  • probing.rebalance.interval.ms
  • rack.aware.assignment.tags
  • rack.aware.assignment.strategy
  • rack.aware.assignment.traffic_cost
  • rack.aware.assignment.non_overlap_cost
  • task.assignor.class(分配发生在 Broker 侧;改用group.streams.assignorsstreams.assignor.name
  • session.timeout.ms(改用分组级配置)
  • heartbeat.interval.ms(改用分组级配置)

从源码结构看,分组级配置解析集中在 GroupConfig.java,它通过ConfigDef.define(...)定义了上述各项(含streams.assignor.namestreams.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#describeStreamsGroupskafka-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.javaTargetAssignmentBuilder.java(当前/目标分配构建)等。

拓扑配置与校验

为了在流式客户端之间分配任务,Group Coordinator 使用拓扑元数据:该元数据在成员加入组时初始化,并持久化在 consumer offsets 主题中。

每当成员加入 streams group 时,第一个心跳请求包含拓扑元数据。该元数据将拓扑描述为一组子拓扑(subtopologies),每个子拓扑由唯一字符串标识符识别,并包含创建内部主题和分配任务所需的元数据。

拓扑校验与 NOT_READY 状态

在处理 streams group 心跳期间,Group Coordinator 可能检测到拓扑所需的 source/sink 主题或内部主题不存在,或它们的配置与拓扑成功执行所需配置不一致。这会触发 "topology configuration"(拓扑配置)过程,Group Coordinator 执行以下步骤:

  1. 检查所有已配置的 source 主题是否存在。
  2. 检查 "copartition groups"(协同分区组)是否满足——即所有应当协同分区的 source 主题确实协同分区。
  3. 根据 source 主题配置推导所有内部主题所需的分区数。
  4. 检查所有内部主题是否存在且配置正确。

如果任何 source 主题或内部主题缺失,组进入NOT_READY状态。在NOT_READY状态下,所有心跳照常处理(因此通常不应失败),但心跳响应中的状态会指示存在何种问题。组处于NOT_READY状态时,所有成员都会获得空分配。

佐证:拓扑元数据与内部主题管理在源码中有对应模块——TopologyMetadata.java 与 topics/ 目录下的CopartitionedTopicsEnforcer.java(协同分区校验)、InternalTopicManager.javaConfiguredTopology.java/ConfiguredSubtopology.javaChangelogTopics.java/RepartitionTopics.java(changelog 与 repartition 内部主题推导)等类。

集中的分配配置

核心分配选项集中在 Broker 上配置,不依赖每个客户端的配置。这允许在不重新部署流式应用的情况下调优 streams group。Broker 侧引入的核心分配选项是num.standby.replicas,既可以全局配置在 Broker 上,也可以通过IncrementalAlterConfigsDescribeConfigsRPC 动态配置到特定 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 应用从经典协议迁移到流式再均衡协议:

  1. 关闭所有应用实例。
  2. 等待session.timeout.ms过期,使组变为空(或强制显式离开组)。
  3. 更新应用配置,设置group.protocol=streams
  4. 重启应用实例。

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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 18:26:10

ECG心电图5分类实战:TCN+Restormer混合模型与Python信号预处理

简介:本资源是一套面向高校本科生及人工智能初学者的心电图(ECG)信号五分类深度学习完整实践方案,聚焦心血管疾病早期筛查中的心律失常识别问题,适用于期末大作业、毕业设计与课程设计等工程实践场景。资源包共36个文件…

作者头像 李华
网站建设 2026/9/10 18:25:43

JAVA毕设项目:基于 SpringBoot 架构的食品安全信息化平台的设计与研究 基于 SpringBoot 的在线食品安全信息管理平台 (源码+文档,讲解、调试运行,定制等)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

作者头像 李华
网站建设 2026/9/10 18:25:39

JAVA毕设项目:基于 SpringBoot 的健身房课程与教练管理系统的设计与实现 基于 SpringBoot 的健身房日常管理系统(源码+文档,讲解、调试运行,定制等)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于Java、小程序技术领域和毕业项目实战 ✌️技术范围:&am…

作者头像 李华
网站建设 2026/9/10 18:24:42

CANN/GE AttrValue属性值类

简介 【免费下载链接】ge GE(Graph Engine)是面向昇腾的图编译器和执行器,提供了计算图优化、多流并行、内存复用和模型下沉等技术手段,加速模型执行效率,减少模型内存占用。 GE 提供对 PyTorch、TensorFlow 前端的友好…

作者头像 李华
网站建设 2026/9/10 18:21:37

YOLOv5车型识别实战:从模型训练到系统部署全流程

简介:这是基于YOLOv5构建的车型识别系统完整工程包,包含可运行的源码和基于PyQt5开发的图形操作界面。系统支持轿车、SUV、商务车三种车型以及奥迪、宝马、大众、奔驰、丰田五种品牌识别,并集成摄像头实时识别、历史记录保存与识别目标数量统…

作者头像 李华