MirrorMaker2源Broker节点异常后消费组偏移停止同步如何重同步?
MirrorMaker2 指定Topic/Group ID偏移重同步方案
针对你使用的strimzi:0.22.1-kafka-2.7.0版本的MirrorMaker2,在不重建Broker的前提下可以通过清理内部Topic残留脏数据的方式恢复偏移同步,操作步骤如下:
前置操作
首先停止所有运行中的MirrorMaker2同步任务,避免操作过程中产生新的脏数据。
方案1:清理MM2内部偏移同步Topic的残留记录(推荐,可自动恢复后续同步)
MirrorMaker2的消费组偏移同步记录默认存储在目标端名为mm2-offset-syncs.<目标集群别名>.internal的Compact类型内部Topic中,你可以通过写入墓碑消息的方式删除指定Group/Topic的异常残留记录:
- 定位异常记录的键名
该内部Topic的记录键格式为<消费组ID>:<源Topic名称>:<分区号>,可通过以下命令查询目标Group对应的所有异常记录:# 目标端Broker节点执行 ./kafka-console-consumer.sh \ --bootstrap-server <目标端Broker B的接入地址> \ --topic mm2-offset-syncs.<你配置的目标集群别名>.internal \ --property print.key=true \ --from-beginning | grep "<异常的消费组ID>" - 写入墓碑消息删除残留记录
对于需要清理的每一条异常记录,执行以下命令写入对应键的空值消息触发日志压缩删除:# 替换尖括号内容为实际值,分区号按实际查询结果填写,需要清理全部分区就循环执行 echo "<异常消费组ID>:<对应源Topic名称>:<分区号>:" | ./kafka-console-producer.sh \ --bootstrap-server <目标端Broker B的接入地址> \ --topic mm2-offset-syncs.<你配置的目标集群别名>.internal \ --property parse.key=true \ --property key.separator=":" - 加速日志压缩生效
可临时调整内部Topic的压缩参数加快清理速度,清理完成后改回原有配置即可:./kafka-topics.sh \ --bootstrap-server <目标端Broker B的接入地址> \ --alter \ --topic mm2-offset-syncs.<你配置的目标集群别名>.internal \ --config segment.bytes=1048576 \ --config min.cleanable.dirty.ratio=0.01 - 恢复同步
等待1-2分钟日志压缩完成后,重启MirrorMaker2服务,等待1-2个偏移同步周期(你配置的是30秒)后,即可验证偏移同步是否恢复:# 源端Broker A查询原始消费组偏移 ./kafka-consumer-groups.sh --bootstrap-server <源端Broker A地址> --group <消费组ID> --describe # 目标端Broker B查询同步后的消费组偏移,注意MM2同步的消费组默认会加源集群别名前缀 ./kafka-consumer-groups.sh --bootstrap-server <目标端Broker B地址> --group <源集群别名>.<消费组ID> --describe
方案2:手动重置指定消费组偏移(临时救急用)
如果不需要后续自动同步,仅需要恢复当前的偏移对齐,可直接在目标端手动重置消费组偏移到源端的最新偏移值:
./kafka-consumer-groups.sh --bootstrap-server <目标端Broker B地址> \ --group <源集群别名>.<消费组ID> \ --topic <源集群别名>.<源Topic名称> \ --reset-offsets \ --to-offset <源端对应分区的最新偏移值> \ --execute
注意事项
- 操作前建议先导出
mm2-offset-syncs.<目标集群别名>.internal的全量内容做备份,避免误删数据 - 如果有多个异常消费组/Topic,可编写简单脚本批量生成墓碑消息执行,无需手动单条操作
内容的提问来源于stack exchange,提问作者m3nthal
相关产品推荐
相关产品推荐

