You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Kafka MirrorMaker 2:重命名Topic并保留消费者偏移量可行性问询

批量重命名Kafka Topic并同步偏移量(含Kafka Streams)

可以通过MirrorMaker 2.0(MM2)同时完成批量重命名Topic和同步保留所有消费者偏移量的操作,无需分开执行两次任务。具体实现步骤如下:

1. 配置MM2的Topic批量重命名规则

MM2支持通过正则表达式匹配源Topic,自动替换为带环境前缀的目标Topic,实现批量重命名。在MM2的配置文件(比如mm2.properties)中添加核心配置:

# 源集群与目标集群地址
clusters=source,target
source.bootstrap.servers=production-kafka:9092
target.bootstrap.servers=development-kafka:9092

# 批量重命名规则:匹配production_开头的Topic,替换为development_前缀
topics.regex=production_(.*)
topics.replacement=development_$1

# 开启单向同步(源集群到目标集群)
source->target.enabled=true

该配置会自动将源集群中所有production_前缀的Topic,同步到目标集群并重命名为development_开头的对应Topic。

2. 配置偏移量同步(覆盖标准消费者与Kafka Streams)

MM2默认支持消费者组偏移量同步,但需要确保Kafka Streams的内部消费者组也被纳入同步范围:

# 同步所有消费者组(包括Kafka Streams的内部组)
consumer.group.filters=include:.*

# 配置偏移量同步的基础参数
offset-syncs.topic.replication.factor=1 # 根据目标集群副本数调整
sync.topic.acls.enabled=false # 无需同步ACL时可关闭

Kafka Streams的消费者组格式通常为kafka-streams-<应用ID>-<随机后缀>,通过include:.*可确保这些内部组的偏移量同步到目标集群,并自动映射到重命名后的Topic上。

3. 执行同步并验证

  • 启动MM2同步任务:
bin/kafka-mirror-maker.sh --config mm2.properties
  • 验证目标集群的Topic是否生成:
bin/kafka-topics.sh --bootstrap-server development-kafka:9092 --list
  • 检查消费者组偏移量(含Streams内部组):
bin/kafka-consumer-groups.sh --bootstrap-server development-kafka:9092 --describe --group <你的消费者组名>
bin/kafka-consumer-groups.sh --bootstrap-server development-kafka:9092 --describe --group kafka-streams-<你的应用ID>-*

确认目标集群的偏移量与源集群对应Topic的偏移量一致,即同步成功。

4. 切换业务应用

同步完成后,修改生产者、消费者及Kafka Streams应用的配置,将输入输出Topic替换为重命名后的development_前缀Topic,启动后会自动从同步好的偏移量位置继续处理数据。

注意事项

  • 确保源集群与目标集群的Kafka版本兼容(MM2要求至少Kafka 2.4及以上)。
  • 若需严格数据一致性,同步期间可暂停源集群的写入操作;允许增量同步的话,可保持写入,MM2会持续同步新消息与偏移量。
  • Kafka Streams应用切换时,无需修改application.id,仅调整Topic配置即可,对应内部消费者组的偏移量已同步到目标集群。

内容的提问来源于stack exchange,提问作者Nicholas Lester

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 14:45:14