Kafka Streams Groups 工具(kafka-streams-groups.sh)实战指南:基于 KIP-1071 的 Streams 组管理
2026/9/10 11:36:24 网站建设 项目流程

Kafka Streams Groups 工具(kafka-streams-groups.sh)实战指南:基于 KIP-1071 的 Streams 组管理

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

kafka-streams-groups.sh是 Apache Kafka 提供的 Streams 组管理命令行工具,用于基于 Streams Rebalance Protocol(KIP-1071)的Streams groups的查看与管理:列出/描述组、查看成员与输入主题偏移量(lag)、重置或删除输入主题偏移量、删除组(可连带删除内部主题)。本文将基于当前仓库的文档与源码,完整讲解该工具的每个命令、全部选项、底层实现原理与安全操作最佳实践,帮助你在生产集群上安全、精准地运维 Kafka Streams 应用。

Streams Group 是什么:与经典 Consumer Group 的区别

在 Kafka 引入 KIP-1071(Streams 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.StreamsGroupCommand(StreamsGroupCommand.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 的打印逻辑):

  • 状态视图--state):GROUP / COORDINATOR (ID) / ASSIGNOR / STATE / #MEMBERS--verbose时追加GROUP-EPOCHTARGET-ASSIGNMENT-EPOCH两列(StreamsGroupCommand.java#L409-L429);
  • 成员视图--members):GROUP / MEMBER / PROCESS / CLIENT-ID / ASSIGNMENTS,其中 ASSIGNMENTS 按ACTIVESTANDBYWARMUP三类任务展示,格式如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);
  • 偏移量视图--offsets):GROUP / TOPIC / PARTITION / OFFSET-LAG--verbose时展示CURRENT-OFFSET / LEADER-EPOCH / LOG-END-OFFSET / OFFSET-LAG(StreamsGroupCommand.java#L431-L457)。

偏移量与 lag 的计算逻辑也值得关注:工具会遍历所有成员的activeTasks收集主题分区,通过listOffsets分别获取 earliest 与 latest,再结合listStreamsGroupOffsets获取已提交偏移量;lag = log-end-offset − current-offset,若从未提交则退化为 latest − earliest(StreamsGroupCommand.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_STOREDERROR分别对应上述两条错误消息,且均返回非零退出码(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就绝不会真正改动偏移量;
  • 只允许对空组操作:仅当组状态为EmptyDead时才重置输入主题偏移量;其他状态直接报错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则可不指定 scope;
  • reset 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-file(CSV);
  • 执行时连带删除内部主题--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 / STATUS(Successful 或 Error 详情)。底层通过adminClient.deleteStreamsGroupOffsets实现,并对INVALID_GROUP_IDGROUP_ID_NOT_FOUNDNON_EMPTY_GROUPGROUP_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中提取repartitionSourceTopicsstateChangelogTopics(排除源主题),并通过isInferredInternalTopic校验命名是否为可推断的内部主题(<applicationId>-<topic>-...模式),非推断的内部主题不会被连带删除并会打印提示(StreamsGroupCommand.java#L908-L958)。若 broker 版本过旧不支持,工具会提示改用kafka-topics.sh手动处理内部主题。底层删除通过adminClient.deleteStreamsGroupsadminClient.deleteTopics完成。

全部选项与标志

核心动作

选项说明
--list列出 Streams groups。可用--state显示/按状态过滤
--describe描述由--group选中的组。可组合:--state(组状态与 epochs)、--members(成员与分配)、--offsets(输入与 repartition 主题偏移量/lag)、--topology(broker 拓扑描述插件记录的处理拓扑);--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-file(CSV)。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-L150(defaultsTo(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/ 下提供了完整的测试覆盖,可作为理解工具行为的参考:ListStreamsGroupTestDescribeStreamsGroupTest(含拓扑描述各分支)、ResetStreamsGroupOffsetTest(dry-run/execute 与各类 specifier)、DeleteStreamsGroupTestDeleteStreamsGroupOffsetTest,以及StreamsGroupCommandTest对整体参数校验的测试。此外,streams/integration-tests 中的TopologyDescriptionPluginIntegrationTestTopologyDescriptionNoPluginIntegrationTest等用例验证了拓扑描述插件与--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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询