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

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应保持源与目标集群的延迟一致。
实际:

  1. 仅目标集群日志末尾偏移更新,目标集群延迟更大(源集群延迟30,目标集群延迟38)
  2. 将消费组切换至目标集群后,出现重复消费8条已在源集群消费过的消息

问题原因分析

  1. 消费组偏移同步触发逻辑:MM2通过sync.group.offsets.interval.seconds定期同步偏移,当源集群消费组停止后,源集群的消费组偏移不再更新,MM2可能未及时触发最终的偏移同步,导致目标集群偏移滞后。
  2. Offset Sync Topic机制:配置offset-syncs.topic.location=target后,偏移同步信息存储在目标集群,但源消费组停止后,MM2的偏移同步任务可能未捕获到偏移的最终状态。
  3. 版本兼容性:源集群2.3.0与3.7.1版本的MM2之间,消费组元数据格式存在差异,可能导致偏移同步不完整。

解决方案

  1. 等待最终偏移同步:停止源集群生产消费后,不要立即切换到目标集群,等待至少sync.group.offsets.interval.seconds(当前为10s)的时间,确保MM2完成最终偏移同步。
  2. 调整MM2同步配置:
    • 临时调小sync.group.offsets.interval.seconds至5s,加快停止后的偏移同步速度
    • 确认refresh.groups.interval.seconds配置能及时感知消费组状态变化
  3. 手动补全偏移:
    • 在源集群导出消费组最终偏移:
      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
      
  4. 校验Checkpoint机制:检查目标集群的mm2-checkpoints.source.internal主题,确认是否存在最新的偏移checkpoint记录;若没有,查看MM2日志排查同步报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 18:44:51