Kafka Streams 如何离线迁移到 group.protocol=streams 并用 kafka-streams-groups.sh 管理 Streams Group

发布时间:2026/9/12 7:08:28
Kafka Streams 如何离线迁移到 group.protocol=streams 并用 kafka-streams-groups.sh 管理 Streams Group
Kafka Streams 如何离线迁移到 group.protocolstreams 并用 kafka-streams-groups.sh 管理 Streams Group【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka如果你的 Kafka Streams 应用目前跑在 classic 协议上想切换到 broker 驱动的 Streams Rebalance ProtocolKIP-1071需要完成两件事在一个维护窗口内把应用离线迁移到group.protocolstreams以及迁移后用bin/kafka-streams-groups.sh查看、核对和管理这个 Streams Group。适用前提broker 和客户端都运行 Apache Kafka 4.2 或更高版本由于 4.2.0 存在离线迁移代码的 broker 端缺陷KAFKA-20254classic 到 streams 的迁移建议在 4.2.1 或更高版本上进行全新创建的 streams group 不受该缺陷影响。迁移完成后broker 侧唯一保留的 group 数据是已提交的 offsetschangelog 和 repartition 等内部 topic 会作为普通 Kafka topic 继续存在。主要参考文档Streams Rebalance Protocol、kafka-streams-groups.sh。迁移前确认 feature 版本已启用Streams Rebalance Protocol 在 Kafka 4.2 之后的新集群上默认启用。对升级到 4.2 的已有集群或者想显式控制该功能的场景先用 feature 描述命令确认当前状态bin/kafka-features.sh --bootstrap-server localhost:9092 describe查看输出中streams.version是否已为FinalizedVersionLevel1。若尚未启用执行升级bin/kafka-features.sh --bootstrap-server localhost:9092 upgrade --feature streams.version1执行后再次运行describe验证变更。回退方式是对应的 downgrade 命令bin/kafka-features.sh --bootstrap-server localhost:9092 downgrade --feature streams.version0以上命令中的localhost:9092是文档示例地址替换为你自己的 bootstrap server。相关 broker 配置如group.coordinator.rebalance.protocols中的streams项、group.streams.session.timeout.ms、group.streams.num.standby.replicas等的完整说明见 broker 配置文档。离线迁移四步操作文档明确说明当前仅支持离线迁移在线迁移应用运行中切换协议在现有版本不可用必须安排维护窗口。完整流程关闭所有应用实例。等待session.timeout.ms过期让 group 变空也可以强制显式 leave group。修改应用配置在 Kafka Streams 应用中设置group.protocolstreams。重启应用实例。反向迁移streams group 转回 classic group是同样流程把配置改回group.protocolclassic。迁移前注意一个客户端配置细节启用 streams 协议后一批客户端配置会被忽略包括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以及session.timeout.ms和heartbeat.interval.ms。后三项在 streams 协议下属于 group 级配置需要改用kafka-configs.sh在 group 维度设置例如bin/kafka-configs.sh --bootstrap-server localhost:9092 \ --alter --entity-type groups --entity-name wordcount \ --add-config streams.num.standby.replicas1其中wordcount是文档示例的 group 名即application.id替换为实际值。group 级可用配置streams.session.timeout.ms、streams.heartbeat.interval.ms、streams.num.standby.replicas、streams.initial.rebalance.delay.ms、streams.assignor.name见 group 配置文档。迁移后验证用 kafka-streams-groups.sh 确认 group 已切换并正常kafka-streams-groups.sh位于bin/下通过--bootstrap-server连接集群安全集群用--command-config传入 AdminClient 属性。--group id指定的就是应用的application.id。先列出集群内所有 Streams group并用--state显示/过滤状态kafka-streams-groups.sh --bootstrap-server localhost:9092 --listStreams group 的状态取值有 Empty、Not Ready、Assigning、Reconciling、Stable、Dead。列出后能看到你的 group 出现说明它已被识别为 streams group 而非 classic consumer group。再深入查看状态和成员# group 状态与 epoch kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --state --verbose # 成员信息当前分配 vs 目标分配以及成员是否仍在使用 classic 协议 kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --members --verbose # 输入 topic 的 offsets 与 lag kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --describe --group my-streams-app --offsetsmy-streams-app是文档示例 group 名替换为实际application.id。验证时重点判断--members输出中迁移成功的成员应显示为 streams 协议成员而不是仍在使用 classic 协议--offsets可以查看处理落后多少如果 group 状态是Not Ready说明 source topic 或内部 topic 缺失或分区配置不满足要求例如 copartition 约束不满足。此时所有心跳仍会正常处理但成员会拿到空分配--describe的 status 会指出问题类型。日常管理操作及其安全边界以下都是文档明确给出的命令。注意--reset-offsets、--delete-offsets、--delete属于变更操作文档要求执行前确认应用实例已停止/不活跃、group 已空且 offset reset 先用--dry-run预览再--execute。重置输入 topic offsets控制重启后的重处理边界。指定符任选其一--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范围用--all-input-topics或一个/多个--input-topic name。文档示例# 先预览把所有输入 topic 重置到指定时间点 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删除 offsets使 group 下次启动时重新消费# 所有输入 topic kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app --delete-offsets --all-input-topics # 指定 topic kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --group my-streams-app \ --delete-offsets --input-topic input-a --input-topic input-b删除 Streams group清理 broker 侧的 offsets、topology、assignments 元数据kafka-streams-groups.sh --bootstrap-server localhost:9092 \ --delete --group my-streams-app如果同时要删除内部 topic追加--delete-internal-topic name或--delete-all-internal-topics。文档特别提醒删除内部 topic 会移除状态所依赖的 topic只有在打算从输入 topic 重建状态时才能这样做。--delete-all-internal-topics是破坏性操作且迁移本身并不删除内部 topic它们作为普通 topic 继续存在除非你确实要清理否则不要顺手加上。限制与已知问题以下限制来自文档规划迁移时需要直接面对不支持在线迁移classic 与 streams 协议之间的迁移只能在所有实例停机后进行。4.2.0 的迁移缺陷KAFKA-20254 是 broker 侧离线迁移代码的缺陷文档不建议在 4.2.0 上执行 classic 到 streams 的迁移修复在 4.2.1 中提供。拓扑更新受限如果拓扑发生重大变化例如新增 source topic、subtopology 数量变化必须创建新的 streams group而不是复用旧 group。不支持正则订阅pattern-based topic 订阅在新协议下不可用。--topology的额外要求--describe --topology需要 broker 运行 Apache Kafka 4.4 或更高版本并配置了group.streams.topology.description.plugin.class在旧版本 broker 上该命令会以UnsupportedVersionException失败。若 broker 未配置 plugin 或应用尚未推送描述工具会打印No topology description is stored for streams group id.并以非零退出码结束。静态成员资格4.4.0 起 streams 协议支持group.instance.id但对没有持久化状态 store 的拓扑Kafka Streams 每次重启会生成新的 process ID导致 broker 重算分配静态成员资格跨重启的收益会被抵消。迁移完成的判定依据就是--list中该 group 以 streams group 出现、--describe --state显示预期状态稳定后为 Stable、--members中不再有 classic 协议成员以及--offsets显示从已提交 offset 继续消费。后续若需要查看 topology 描述细节或 plugin 配置见 Topology Description Plugin配置项全集见 Kafka Streams 配置文档。【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考