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;
- 组状态、组 epoch、目标分配 epoch(配合
- 重置输入主题偏移量,用精确的规格(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-EPOCH与TARGET-ASSIGNMENT-EPOCH两列(StreamsGroupCommand.java#L409-L429); - 成员视图(
--members):GROUP / 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); - 偏移量视图(
--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_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-run:
resetOffsets()中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则可不指定 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_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)、--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/ 下提供了完整的测试覆盖,可作为理解工具行为的参考:ListStreamsGroupTest、DescribeStreamsGroupTest(含拓扑描述各分支)、ResetStreamsGroupOffsetTest(dry-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),仅供参考