Kafka MirrorMaker2源集群消费组闲置时偏移同步异常问题
Kafka 2.3.0(ZK版)迁移至3.7.1(Kraft版)MM2偏移同步异常问题
背景
将3节点Kafka 2.3.0(Zookeeper版)集群迁移至3节点3.7.1(Kraft版)集群,允许停服但要求无消息丢失。采用MirrorMaker2(MM2)同步主题消息、消费组及偏移量,MM2配置如下:
# specify any number of cluster aliases clusters=source,destination # connection information for each cluster # This is a comma separated host:port pairs for each cluster # for example. "A_host1:9092, A_host2:9092, A_host3:9092" and you can see the exact host name on Ambari > Hosts source.bootstrap.servers=kafka01.rnd.lan:9093,kafka02.rnd.lan:9093,kafka03.rnd.lan:9093 destination.bootstrap.servers=tor1kafka1-clu1.lab.lan:9093,tor1kafka2-clu1.lab.lan:9093,tor1kafka3-clu1.lab.lan:9093 # enable and configure individual replication flows source->destination.enabled=true # regex which defines which topics gets replicated. For eg "foo-.*" source->destination.topics=GEN.LAB.P.TEST.TESTINFO.OLD2NEW offset.max.lag=1 offset-syncs.topic.location=target refresh.topics.enabled=true refresh.topics.interval.seconds=10 checkpoints.topic.replication.factor=3 heartbeats.topic.replication.factor=3 offset-syncs.topic.replication.factor=3 offset.storage.replication.factor=3 status.storage.replication.factor=3 config.storage.replication.factor=3 emit.heartbeats.enabled=true emit.heartbeats.interval.seconds=10 replication.policy.class=org.apache.kafka.connect.mirror.IdentityReplicationPolicy replication.factor=3 refresh.groups.enabled=true refresh.groups.interval.seconds=10 sync.group.offsets.enabled=true sync.group.offsets.interval.seconds=10 emit.checkpoints.enabled=true emit.checkpoints.interval.seconds=10
测试流程
- 在源集群
GEN.LAB.P.TEST.TESTINFO.OLD2NEW主题启动消息生产 - 使用消费组
test-consumer-kafkaupgrade-109在源集群启动消费 - 启动MM2
此时源集群的主题和消费组已在目标集群正确生成,生产/消费活跃时当前偏移、日志末尾偏移同步正常。
异常现象
执行以下操作后:
- 停止源集群的消费者
- 发送少量消息后停止源集群的生产者
预期:源集群生产/消费停止后,MM2应保持源与目标集群的延迟一致。
实际:
- 仅目标集群日志末尾偏移更新,目标集群延迟更大(源集群延迟30,目标集群延迟38)
- 将消费组切换至目标集群后,出现重复消费8条已在源集群消费过的消息
问题原因分析
- 消费组偏移同步触发逻辑:MM2通过
sync.group.offsets.interval.seconds定期同步偏移,当源集群消费组停止后,源集群的消费组偏移不再更新,MM2可能未及时触发最终的偏移同步,导致目标集群偏移滞后。 - Offset Sync Topic机制:配置
offset-syncs.topic.location=target后,偏移同步信息存储在目标集群,但源消费组停止后,MM2的偏移同步任务可能未捕获到偏移的最终状态。 - 版本兼容性:源集群2.3.0与3.7.1版本的MM2之间,消费组元数据格式存在差异,可能导致偏移同步不完整。
解决方案
- 等待最终偏移同步:停止源集群生产消费后,不要立即切换到目标集群,等待至少
sync.group.offsets.interval.seconds(当前为10s)的时间,确保MM2完成最终偏移同步。 - 调整MM2同步配置:
- 临时调小
sync.group.offsets.interval.seconds至5s,加快停止后的偏移同步速度 - 确认
refresh.groups.interval.seconds配置能及时感知消费组状态变化
- 临时调小
- 手动补全偏移:
- 在源集群导出消费组最终偏移:
kafka-consumer-groups.sh --bootstrap-server kafka01.rnd.lan:9093 --describe --group test-consumer-kafkaupgrade-109 - 在目标集群重置消费组偏移到对应值:
kafka-consumer-groups.sh --bootstrap-server tor1kafka1-clu1.lab.lan:9093 --reset-offsets --to-offset <源集群偏移值> --group test-consumer-kafkaupgrade-109 --topic GEN.LAB.P.TEST.TESTINFO.OLD2NEW --execute
- 在源集群导出消费组最终偏移:
- 校验Checkpoint机制:检查目标集群的
mm2-checkpoints.source.internal主题,确认是否存在最新的偏移checkpoint记录;若没有,查看MM2日志排查同步报错。
内容的提问来源于stack exchange,提问作者user27614415
相关产品推荐
相关产品推荐

