Kafka Streams Groups 工具kafka-streams-groups.sh实战指南基于 KIP-1071 的 Streams 组管理【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafkakafka-streams-groups.sh是 Apache Kafka 提供的 Streams 组管理命令行工具用于基于 Streams Rebalance ProtocolKIP-1071的Streams groups的查看与管理列出/描述组、查看成员与输入主题偏移量lag、重置或删除输入主题偏移量、删除组可连带删除内部主题。本文将基于当前仓库的文档与源码完整讲解该工具的每个命令、全部选项、底层实现原理与安全操作最佳实践帮助你在生产集群上安全、精准地运维 Kafka Streams 应用。Streams Group 是什么与经典 Consumer Group 的区别在 Kafka 引入 KIP-1071Streams Rebalance Protocol之前Kafka Streams 应用底层复用经典的 consumer group 协议来完成成员管理与分区分配。KIP-1071 之后Streams 应用使用由 broker 协调、面向 Streams 专用 RPC 与元数据的组类型即Streams group。它与经典 consumer group 的关键区别在于组状态、分配信息assignment与输入主题偏移量均以 Streams 专用语义存储与暴露组状态包括 Streams 特有状态Empty、NotReady、Assigning、Reconciling、Stable、Dead。这些状态定义在 GroupState.java 中其中groupStatesForType(GroupType.STREAMS)返回STABLE、DEAD、EMPTY、ASSIGNING、RECONCILING、NOT_READY六种见该文件 clients/src/main/java/org/apache/kafka/common/GroupState.java#L79-L87组 id 即 Streams 应用的application.id。kafka-streams-groups.sh正是为这类 Streams 组提供 CLI 视图与运维能力它暴露 Streams 特有的状态、分配与输入主题偏移量使管理员可以像使用 consumer-group 工具一样直观地观察和治理 Streams 应用但语义完全面向 KIP-1071 的 Streams 组。谨慎使用偏移量重置/删除、组删除等变更类操作会影响应用重启后的重处理行为。请始终先用--dry-run预览偏移量重置结果并在执行前确保应用实例已停止/停用、组处于空闲状态Empty。工具能做什么列出集群中的 Streams groups并按组状态Empty、Not Ready、Assigning、Reconciling、Stable、Dead展示或过滤描述某个 Streams group展示组状态、组 epoch、目标分配 epoch配合--state、--verbose查看更多细节每个成员的信息成员 epoch、当前分配与目标分配、该成员是否仍在使用经典协议配合--members、--verbose输入主题的偏移量与 lag配合--offsets了解处理进度落后多少处理拓扑processing topology——由 broker 的 topology description plugin 记录配合--topology输出格式与Topology#describe()一致。该功能要求 broker 运行 Apache Kafka 4.4 或更新版本并配置group.streams.topology.description.plugin.class重置输入主题偏移量用精确的规格earliest、latest、to-offset、to-datetime、by-duration、shift-by、from-file控制重处理边界。需要--dry-run或--execute且要求实例处于停用状态删除输入主题偏移量强制下次启动时重新消费删除 Streams group清理 broker 侧 Streams 元数据偏移量、拓扑、分配。可通过--delete-internal-topic删除指定的内部主题或通过--delete-all-internal-topics删除全部内部主题。使用方式脚本位于bin/kafka-streams-groups.sh仓库根目录 bin/kafka-streams-groups.sh通过--bootstrap-server连接集群对于启用安全认证的集群用--command-config传入 AdminClient 的属性文件。$ kafka-streams-groups.sh --bootstrap-server host:port [COMMAND] [OPTIONS]从源码看脚本本体只是启动入口实际逻辑全部位于org.apache.kafka.tools.streams.StreamsGroupCommandStreamsGroupCommand.java通过kafka-run-class.sh以 AdminClient 方式与 broker 交互exec $(dirname $0)/kafka-run-class.sh org.apache.kafka.tools.streams.StreamsGroupCommand $StreamsGroupCommand.main()在解析参数后强制校验五个动作--list、--describe、--delete、--reset-offsets、--delete-offsets必须且只能指定一个否则抛出IllegalArgumentException见 StreamsGroupCommand.java#L101-L110。说明kafka-streams-groups.sh是 Streams 组管理 Admin API 的 CLI 封装。它提供与 consumer-group 工具在精神上类似的 list/describe/delete 与偏移量管理操作但完全针对 KIP-1071 定义的 Streams 组定制。命令详解列出 Streams groups--list发现集群中的 Streams 组# 列出所有 Streams groups kafka-streams-groups.sh --bootstrap-server localhost:9092 --list--list默认只打印组 id。结合--state可以显示状态列或按指定状态过滤多个状态用逗号分隔# 列出所有组并显示状态 kafka-streams-groups.sh --bootstrap-server localhost:9092 --list --state # 只列出 Stable 状态的组 kafka-streams-groups.sh --bootstrap-server localhost:9092 --list --state Stable # 列出 Stable 与 Assigning 状态的组 kafka-streams-groups.sh --bootstrap-server localhost:9092 --list --state Stable,Assigning--state的合法取值正是 Streams 组的六种状态Empty, NotReady, Stable, Assigning, Reconciling, Dead。从源码看groupStatesFromString()会先按逗号切分解析再用GroupState.groupStatesForType(GroupType.STREAMS)校验合法性非法状态会直接报错并列出合法值StreamsGroupCommand.java#L171-L180。底层通过adminClient.listGroups(new ListGroupsOptions().withTypes(Set.of(GroupType.STREAMS)))实现即只查询 Streams 类型组StreamsGroupCommand.java#L232-L242。描述 Streams groups--describe检查组的状态、成员与 lag# 描述一个组状态 epochs kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --state --verbose # 描述一个组成员当前分配 vs 目标分配classic/streams 协议 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --members --verbose # 描述一个组输入主题偏移量与 lag kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --offsets # 描述一个组处理拓扑 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --topology--describe可以配合--all-groups作用于所有 Streams 组。若不指定--state/--members/--topology默认展示偏移量视图等价于--offsets——这与 StreamsGroupCommandOptions.java 中OFFSETS_DOC的说明一致This is the default sub-action。各子视图的列含义来自 StreamsGroupCommand.java 的打印逻辑状态视图--stateGROUP / COORDINATOR (ID) / ASSIGNOR / STATE / #MEMBERS--verbose时追加GROUP-EPOCH与TARGET-ASSIGNMENT-EPOCH两列StreamsGroupCommand.java#L409-L429成员视图--membersGROUP / MEMBER / PROCESS / CLIENT-ID / ASSIGNMENTS其中 ASSIGNMENTS 按ACTIVE、STANDBY、WARMUP三类任务展示格式如0:[1,2]; 1:[3];--verbose时展示TARGET-ASSIGNMENT-EPOCH / TOPOLOGY-EPOCH / MEMBER / MEMBER-PROTOCOL / MEMBER-EPOCH / PROCESS / CLIENT-ID / ASSIGNMENTS并在 ASSIGNMENTS 中追加TARGET-ACTIVE/TARGET-STANDBY/TARGET-WARMUP同时以member.isClassic() ? classic : streams标明成员使用的协议StreamsGroupCommand.java#L327-L407偏移量视图--offsetsGROUP / TOPIC / PARTITION / OFFSET-LAG--verbose时展示CURRENT-OFFSET / LEADER-EPOCH / LOG-END-OFFSET / OFFSET-LAGStreamsGroupCommand.java#L431-L457。偏移量与 lag 的计算逻辑也值得关注工具会遍历所有成员的activeTasks收集主题分区通过listOffsets分别获取 earliest 与 latest再结合listStreamsGroupOffsets获取已提交偏移量lag log-end-offset − current-offset若从未提交则退化为 latest − earliestStreamsGroupCommand.java#L459-L496。描述处理拓扑--topology--topology打印该组的处理拓扑——由 broker 的 topology description plugin 记录输出格式与Topology#describe()一致Topologies: Sub-topology: 0 Source: KSTREAM-SOURCE-0000000000 (topics: [streams-plaintext-input]) -- KSTREAM-FLATMAPVALUES-0000000001 Processor: KSTREAM-FLATMAPVALUES-0000000001 (stores: []) -- KSTREAM-AGGREGATE-0000000002 -- KSTREAM-SOURCE-0000000000 ...该功能要求 broker 运行Apache Kafka 4.4 或更新版本并配置 broker 参数group.streams.topology.description.plugin.class对旧版本 broker 执行会以UnsupportedVersionException失败。若无可用拓扑描述工具会打印以下消息之一并以非零退出码结束No topology description is stored for streams group id.—— 未记录描述例如 broker 未配置拓扑描述插件或应用尚未推送描述The broker failed to fetch the topology description for streams group id. See the broker logs for details.—— broker 的插件读取已存描述失败。从源码看--topology通过DescribeStreamsGroupsOptions.includeTopologyDescription(true)请求拓扑并根据返回的topologyDescriptionStatus()分发AVAILABLE正常打印、NOT_STORED与ERROR分别对应上述两条错误消息且均返回非零退出码StreamsGroupCommand.java#L308-L325。拓扑描述插件的完整工作机制与故障排查方法见 Topology Description Plugin 文档。重置输入主题偏移量--reset-offsets先预览再执行确保所有应用实例都已停止/停用。执行前始终用--dry-run预览确认影响范围后再用--execute落地# 预览将全部输入主题重置到指定时间戳 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app \ --reset-offsets --all-input-topics --to-datetime 2025-01-31T23:57:00.000 \ --dry-run # 执行将全部输入主题重置到指定时间戳 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app \ --reset-offsets --all-input-topics --to-datetime 2025-01-31T23:57:00.000 \ --execute关键规则与源码实现一致见 StreamsGroupCommand.java#L559-L620默认即为 dry-runresetOffsets()中boolean dryRun opts.options.has(opts.dryRunOpt) || !opts.options.has(opts.executeOpt)——不显式传--execute就绝不会真正改动偏移量只允许对空组操作仅当组状态为Empty或Dead时才重置输入主题偏移量其他状态直接报错Assignments can only be reset if the group id is inactive, but the current state is state.scope 二选一--all-input-topics或一个/多个--input-topic name--input-topic支持topic:0,1,2形式指定分区子集若使用--from-file则可不指定 scopereset specifier 必须恰好选择一个--to-earliest、--to-latest、--to-current、--to-offset n、--by-duration PnDTnHnMnS、--to-datetime YYYY-MM-DDTHH:mm:SS.sss、--shift-by n正负均可、--from-fileCSV执行时连带删除内部主题--execute模式下工具会删除与该组关联的内部主题可在--execute基础上追加--delete-internal-topic name指定部分或--delete-all-internal-topics全部删除——因为重置偏移后状态存储必须重建重置结果以表格打印GROUP / TOPIC / PARTITION / NEW-OFFSET还可以加--export将待重置偏移量以 CSV 格式输出便于归档或配合--from-file复用CSV 导出/导入逻辑复用CsvUtils见 StreamsGroupCommand.java#L359-L383底层通过alterStreamsGroupOffsets写入新偏移量StreamsGroupCommand.java#L960-L983。删除偏移量以强制重新消费--delete-offsets删除全部或指定输入主题的偏移量让组在下次启动时重新读取数据# 删除全部输入主题的偏移量 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app \ --delete-offsets --all-input-topics # 删除指定主题的偏移量 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app \ --delete-offsets --input-topic input-a --input-topic input-b--delete-offsets一次只支持一个组--group支持多个主题--input-topic同样支持topic:0,1,2分区子集语法。执行成功打印Request succeeded for deleting offsets from group id.随后按分区输出TOPIC / PARTITION / STATUSSuccessful 或 Error 详情。底层通过adminClient.deleteStreamsGroupOffsets实现并对INVALID_GROUP_ID、GROUP_ID_NOT_FOUND、NON_EMPTY_GROUP、GROUP_SUBSCRIBED_TO_TOPIC等顶层错误做了分类提示StreamsGroupCommand.java#L697-L757。删除 Streams group--delete清理删除 broker 侧的 Streams 组元数据偏移量、拓扑、分配并可选择删除内部主题# 删除 Streams group 元数据 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --delete --group my-streams-app # 连带删除全部内部主题谨慎使用 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --delete --group my-streams-app \ --delete-all-internal-topics删除前的安全检查preAdminCallChecks见 StreamsGroupCommand.java#L853-L875组必须存在且为 Streams 组否则报Group id does not exist or is not a streams group.组状态不能是DEAD且必须是EMPTY非空组报Streams group id is not EMPTY.。--delete支持--all-groups作用于所有组。内部主题的识别逻辑从subtopologies中提取repartitionSourceTopics与stateChangelogTopics排除源主题并通过isInferredInternalTopic校验命名是否为可推断的内部主题applicationId-topic-...模式非推断的内部主题不会被连带删除并会打印提示StreamsGroupCommand.java#L908-L958。若 broker 版本过旧不支持工具会提示改用kafka-topics.sh手动处理内部主题。底层删除通过adminClient.deleteStreamsGroups与adminClient.deleteTopics完成。全部选项与标志核心动作选项说明--list列出 Streams groups。可用--state显示/按状态过滤--describe描述由--group选中的组。可组合--state组状态与 epochs、--members成员与分配、--offsets输入与 repartition 主题偏移量/lag、--topologybroker 拓扑描述插件记录的处理拓扑--verbose提供更多细节如适用的 leader epochs--reset-offsets重置输入主题偏移量一次一个组实例应停用。必须且只能选一个 specifier--to-earliest、--to-latest、--to-current、--to-offset n、--by-duration PnDTnHnMnS、--to-datetime YYYY-MM-DDTHH:mm:SS.sss、--shift-by n±、--from-fileCSV。scope--all-input-topics或一个/多个--input-topic name。安全要求必须--dry-run或--execute--execute下可追加--delete-internal-topic name或--delete-all-internal-topics删除内部主题--delete-offsets删除--all-input-topics或指定--input-topic的偏移量--delete删除 Streams group 元数据可追加--delete-all-internal-topics删除全部内部主题通用标志选项说明--group id目标 Streams group即application.id--all-groups作用于所有组允许用于--delete--bootstrap-server host:port要连接的 broker必填--command-config file传给 AdminClient 的属性文件安全、超时等配置--timeout ms部分操作中等待组稳定的时间默认 30000ms--dry-run/--execute偏移量重置操作的预览与执行--help/--version/--verbose用法、版本、详细输出需要留意的是--timeout的默认值 30000ms 定义在 StreamsGroupCommandOptions.java#L146-L150defaultsTo(30000L)它会透传给listGroups等 Admin 调用--verbose的语义随子命令不同而不同状态/成员/偏移量视图分别展示不同列完整说明见 StreamsGroupCommandOptions.java#L77-L80。最佳实践与安全要点先预览再执行偏移量重置前务必用--dry-run验证主题范围与影响确认无误后再--execute。源码保证未显式传--execute时绝无任何写入确保实例停用、组为空--reset-offsets与--delete都要求组处于Empty或Dead状态非空组会被拒绝执行谨慎使用内部主题删除--delete-internal-topic与--delete-all-internal-topics会删除状态存储所依赖的主题repartition 主题与状态变更日志主题。仅在确实希望从输入主题重建状态时才使用工具对非推断的内部主题会拒绝连带删除避免误删用户数据版本兼容性--describe --topology依赖 broker 4.4 的group.streams.topology.description.plugin.class配置旧版本会以UnsupportedVersionException失败内部主题的自动识别/删除也依赖对应 broker 版本能力版本过旧时工具会明确提示改用kafka-topics.sh手工处理充分利用 CSV--export导出待重置偏移量 CSV、--from-file导入 CSV适用于大批量、可审计的偏移量重置场景。测试与验证仓库在 tools/src/test/java/org/apache/kafka/tools/streams/ 下提供了完整的测试覆盖可作为理解工具行为的参考ListStreamsGroupTest、DescribeStreamsGroupTest含拓扑描述各分支、ResetStreamsGroupOffsetTestdry-run/execute 与各类 specifier、DeleteStreamsGroupTest、DeleteStreamsGroupOffsetTest以及StreamsGroupCommandTest对整体参数校验的测试。此外streams/integration-tests 中的TopologyDescriptionPluginIntegrationTest、TopologyDescriptionNoPluginIntegrationTest等用例验证了拓扑描述插件与--topology功能的端到端行为。本文档完整记录了kafka-streams-groups.sh对 KIP-1071 Streams groups 的全部能力。相关文档索引 Documentation Kafka Streams Developer Guide。【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考