Kafka Consumer 如何从 classic 协议在线迁移到 group.protocol=consumer
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
如果你的消费组目前运行在 Kafka 4.0 集群上,还想避免停机切换协议,这篇文章解决的就是这个问题:在不停下消费组的前提下,把KafkaConsumer从classicrebalance 协议切换到group.protocol=consumer指定的新一代 Consumer 协议(KIP-848)。迁移完成后,消费组会由服务端自动从Classic类型转换为Consumer类型,消费不中断。
相关文档:Consumer Rebalance Protocol、Consumer and Share Consumer Configs、Basic Kafka Operations。
前提条件
在线迁移依赖服务端与客户端两侧都具备能力,核对以下三点:
- Broker 为 Apache Kafka 4.0 及以上。4.0 起,新 Consumer 协议在服务端自动启用,由
group.versionfeature flag 控制;升级到 4.0 完成(upgrade finalized)后该协议即生效,无需额外开启。参见 upgrade guide。 - 客户端使用 Apache Kafka 4.0 及以上。4.0 起 Consumer 完全支持新协议,但默认不启用,必须显式设置
group.protocol=consumer。 - 当前消费组使用的 classic 分配策略不内嵌自定义元数据(custom metadata)。这是在线迁移的硬性限制:只有当 classic 组使用不内嵌自定义元数据的 assignor 时,滚动升级才能把组从
Classic转换为Consumer。如果你的组使用了内嵌自定义元数据的分配策略,只能走离线迁移(先停掉所有消费者再整体拉起),见文末说明。
执行在线迁移:滚动发布消费者
在线迁移的操作路径就是**滚动发布(rolling out)**消费者,逐步把实例换成带新配置启动的版本。文档给出的关键机制是:
当第一个使用新 Consumer rebalance 协议的消费者加入组时,该组会从
Classic转换为Consumer,之后 Classic 协议的成员会与新协议成员互操作(interoperated)共同工作。
也就是说,不需要一次性把所有实例切完:先替换一两个实例,组类型即完成转换,其余 classic 实例继续在线消费,直到滚动发布全部完成。
以配置文件方式启动消费者为例,在原有配置中增加一行group.protocol=consumer(group.id替换为你要迁移的现有消费组 ID,bootstrap.servers替换为你的集群地址):
group.id=my-group bootstrap.servers=localhost:9092 group.protocol=consumer启用新协议后,以下客户端配置和 API不再可用,发布前应从配置和代码中移除或停用:
heartbeat.interval.mssession.timeout.mspartition.assignment.strategyenforceRebalance(String)与enforceRebalance()
心跳与会话超时改由服务端配置控制:
group.consumer.heartbeat.interval.msgroup.consumer.session.timeout.ms
分配策略也移到服务端,由group.consumer.assignors指定可用 assignor 列表:默认提供uniform和range两种,uniform是默认值(列表中第一个),客户端可通过group.remote.assignor指定其他已注册的 assignor。原客户端分配策略与服务端 assignor 的对应关系如下:
| 客户端 assignor | 服务端 assignor |
|---|---|
| RangeAssignor | range |
| CooperativeStickyAssignor | uniform |
| StickyAssignor | uniform |
| RoundRobinAssignor | uniform |
如果原来依赖自定义客户端分配策略,注意这属于当前不支持的范围(见“限制与边界”)。
验证迁移结果
迁移完成与否,文档给出的判断标准是组类型发生了转换——第一个consumer协议成员加入后组即从Classic转为Consumer。可以结合以下两种手段观察:
查看消费组状态,使用 basic-kafka-operations.md 中给出的描述命令:
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group my-group注意一个权限差异:如果消费组使用的是 consumer 协议,admin client 需要对组内成员订阅的所有 topic拥有
DESCRIBE权限;classic 协议没有这个要求。如果你的环境里此前没有给 admin client 授权全量 topic 的 DESCRIBE 权限,迁移后这条命令可能报权限错误,需要先补齐 ACL。观察服务端组计数指标。监控文档中定义了按协议区分的组数量指标(见 monitoring.md):
kafka.server:type=group-coordinator-metrics,name=group-count,protocol={consumer|classic|streams}迁移过程中
protocol=consumer计数应包含你的组、protocol=classic对应减少。文档未给出该指标的具体示例数值,以你集群的实际监控输出为准。
如果滚动发布过程中某个实例启动失败或长期停留在 classic,检查其配置是否仍带着上面列出的不可用配置项(如partition.assignment.strategy),以及是否使用了内嵌自定义元数据的分配策略——后者会阻断在线转换。
回滚与降级路径
文档同时说明了反向操作:用相反的过程滚动把消费者换回group.protocol=classic(或不设置该配置),当最后一个使用新 Consumer 协议的消费者离开组时,组会从Consumer转换回Classic。因此在线迁移是可逆的,这可以作为发布出问题时的回滚手段。
需要注意两个不可逆点:
- 一旦新协议被消费组使用,集群只能降级到 3.4.1 或更高版本(见 upgrade guide 中 4.0 升级说明)。
- 协议演进路线(文档预期时间线):4.0 GA;5.0 中
KafkaConsumer将默认使用Consumer协议但仍兼容Classic;6.0 中KafkaConsumer仅支持Consumer协议,broker 端继续兼容Classic。
限制与边界
- 客户端自定义 assignor 不受支持:自定义分配策略不在新协议范围内,文档建议通过 KAFKA-18327 反馈;若你有此类依赖,本方案的在线迁移不适用。
- Rack-aware 分配策略尚未完全支持(工作正在进行中,见 KAFKA-19387)。依赖 rack 感知分配的组应谨慎评估后再迁移。
- 离线替代路径:如果无法在线转换(例如 classic 组使用了内嵌自定义元数据的 assignor),可以在所有消费者停机后以
group.protocol=consumer重新拉起——空组会自动在Classic与Consumer之间转换。代价是消费组必须整体停机。 - 启用新协议后,正则订阅可用
subscribe(SubscriptionPattern)方法,表达式采用 RE2J 格式并在服务端求值;新协议还新增了覆盖改进线程模型的指标,具体名称见文档中引用的 KIP 848 相关 wiki 页面(文档未在本仓库内展开)。
完成迁移并确认组类型为Consumer后,后续版本升级时客户端无需再关心group.protocol的默认值变化——5.0 起它将成为KafkaConsumer的默认协议。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考