详解:Broker 驱动的流式任务再均衡与配置迁移指南)
Kafka Streams Rebalance ProtocolKIP-1071详解Broker 驱动的流式任务再均衡与配置迁移指南【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka导读本指南围绕 Apache Kafka 4.2 引入的Streams Rebalance ProtocolKIP-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.protocolstreams配置启用专用的流式再均衡协议。它把 streams group 与 consumer group 分离并在 Broker 上提供流式特有的组成员生命周期与元数据管理。Sticky Task Assignor内置一种在再均衡期间最小化任务迁移的基础任务分配策略注册名为sticky是默认分配器。其实现位于 StickyTaskAssignor.javaname()返回sticky通过assign(GroupSpec, TopologyDescriber)计算分配并在分配活跃任务与备援任务时优先保留各成员先前持有的任务见该文件assignActive/assignStandby逻辑。自定义 TaskAssignorBroker 可通过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 PluginBroker 可通过可插拔后端记录每个 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.protocolstreams时流式应用可配置group.instance.id见 config-streams。但对于没有持久化状态存储的拓扑Kafka Streams 每次重启都会生成新的进程 IDprocess ID导致 Broker 重新计算组分配从而在跨重启场景下实际上抵消了静态成员的收益。当前版本不支持的能力使用新协议时应避免以下尚未就绪的特性拓扑更新如果拓扑发生显著变化例如新增源主题或改变子拓扑数量必须创建一个新的 streams group。高可用分配器High Availability Assignorsticky 是唯一的内置分配器且不支持机架感知分配但通过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.version1bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature streams.version1禁用特性streams.version0bin/kafka-features.sh --bootstrap-server localhost:9092 downgrade --feature streams.version0客户端配置在 Kafka Streams 应用配置中设置group.protocolstreams配置体系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的配置可在DescribeConfigs与IncrementalAlterConfigsRPC 中使用以动态覆盖特定组的默认 Broker 配置。可通过AdminJava 接口或bin/kafka-configs.sh工具设置。完整细节见 group-configs。以下分组级配置可用于 streams group源码定义见 GroupConfig.javastreams.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.replicas1注意在 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.namesession.timeout.ms改用分组级配置heartbeat.interval.ms改用分组级配置从源码结构看分组级配置解析集中在 GroupConfig.java它通过ConfigDef.define(...)定义了上述各项含streams.assignor.name、streams.num.standby.replicas等并将未显式覆盖的项回落到GroupCoordinatorConfig中对应的 Broker 默认值。这意味着同一份配置代码同时支撑kafka-configs.sh的 GROUP 实体操作与 Broker 动态配置校验。管理运维AdministrationAdmin 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.javachangelog 与 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 监控一节。按协议统计的组数量按协议类型统计的组数量其中协议列表通过protocolstreams变体扩展kafka.server:typegroup-coordinator-metrics,namegroup-count,protocol{consumer|classic|streams}按状态统计的 Streams 组数量按状态统计的 streams group 数量kafka.server:typegroup-coordinator-metrics,namestreams-group-count,state{empty|not_ready|assigning|reconciling|stable|dead}Streams 组再均衡Streams group 再均衡传感器kafka.server:typegroup-coordinator-metrics,namestreams-group-rebalance-rate kafka.server:typegroup-coordinator-metrics,namestreams-group-rebalance-count从经典协议迁移当前仅支持离线迁移。将 Kafka Streams 应用从经典协议迁移到流式再均衡协议关闭所有应用实例。等待session.timeout.ms过期使组变为空或强制显式离开组。更新应用配置设置group.protocolstreams。重启应用实例。Broker 侧唯一保留的组数据是已提交的偏移量。其他所有组元数据将在应用以新协议启动时重新创建。内部主题changelog 与 repartition topic将继续作为普通 Kafka 主题存在。类似地可以按相同流程将 streams group 转回经典组只需设置group.protocolclassic。警告本版本不支持在线迁移应用运行期间迁移。在协议之间迁移时请规划维护窗口。警告由于离线迁移代码中存在一个严重的 Broker 侧缺陷KAFKA-20254官方建议在 4.2.0 中避免从经典组迁移到 streams group。新建的 streams group 不受影响。修复在 4.2.1 中可用。小结Streams Rebalance ProtocolKIP-1071将 Kafka Streams 的任务协调从客户端迁至 Broker以streams组类型、StreamsGroupHeartbeat专用 RPC、Broker 侧集中式分配配置group.streams.assignorsstreams.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),仅供参考