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

使用MirrorMaker 2迁移EC2 Kafka至Strimzi时部分消费组未同步求助

Kafka MirrorMaker 2消费组同步缺失问题

问题背景

我们正在将EC2上部署的Kafka集群通过MirrorMaker 2迁移至Strimzi环境,当前Topic及事件同步正常,但消费组同步存在缺失:源端约4000个活跃消费组,仅同步成功2104个。我们可以接受同步延迟,但无法容忍消费组完全缺失。

环境版本

  • Kafka版本:3.7.1
  • Strimzi版本:0.43

MirrorMaker 2配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaMirrorMaker2
metadata:
 name: mirror-maker-2
 annotations:
   strimzi.io/use-prometheus: "true"
spec:
 replicas: 40
 version: 3.7.1  
 connectCluster: "target-a"
 metricsConfig:
   type: jmxPrometheusExporter
   valueFrom:
     configMapKeyRef:
       name: mirrormaker-metrics
       key: metrics-config.yml
 clusters:
   - alias: "source-a"
     bootstrapServers: kafka-source.kafka:9094
     authentication:
       type: plain
       username: NDSNQ6GM89NAB3MG
       passwordSecret:
        secretName: my-ec2-secret
        password: plain.password
     config:
       security.protocol: SASL_PLAINTEXT
       sasl.mechanism: PLAIN

   - alias: "target-a"
     bootstrapServers: kafka-target.kafka:9094
     authentication:
       type: scram-sha-512
       username: my-user
       passwordSecret:
        secretName: my-kafka-user-secret-new
        password: plain.password
     config:
       rebalance.timeout.ms: 30000
       connections.max.idle.ms: 600000
       offset.flush.interval.ms: 30000
       producer.batch.size: 655360
       producer.acks: all
       producer.linger.ms: 10
       producer.buffer.memory: 1073741824
       producer.compression.type: none
       security.protocol: SASL_PLAINTEXT 
       sasl.mechanism: SCRAM-SHA-512
       producer.max.request.size: 10485760
       offset_flush_timeout: 250000
       offset.storage.replication.factor: -1
       status.storage.replication.factor: -1
       config.storage.replication.factor: -1
       status.storage.topic: mm2-status.source-a.internal
       config.storage.topic: mm2-configs.source-a.internal

 mirrors:
   - sourceCluster: "source-a"
     targetCluster: "target-a"
     sourceConnector:
       tasksMax: 500
       autoRestart:
        enabled: true
       config:
         producer.override.batch.size: 655360
         producer.override.linger.ms: 100
         producer.request.timeout.ms: 60000
         admin.timeout.ms: 60000
         connector.client.config.override.policy: All
         exactly.once.source.support: disabled
         consumer.isolation.level: read_committed
         replication.policy.class: "org.apache.kafka.connect.mirror.IdentityReplicationPolicy"
         errors.retry.timeout: 500
         errors.retry.delay.max.ms: 100
         errors.log.enable: true
         offset-syncs.topic.location: "target"
         offset.lag.max: 100 
         refresh.topics.enabled: true
         sync.topic.acls.enabled: true
         sync.topic.configs.enabled: true
         replication.factor: 3
         offset-syncs.topic.replication.factor: 3 
         refresh.topics.interval.seconds: 60  
         refresh.topic.list: true
         consumer.poll.timeout.ms: 10000
         consumer.max.poll.records: 2000
         consumer.fetch.max.wait.ms: 100
         consumer.fetch.max.bytes: 104857600 
         consumer.max.poll.interval.ms: 600000
         consumer.heartbeat.interval.ms: 6000
         consumer.session.timeout.ms: 60000
         consumer.fetch.min.bytes: 1048576
         consumer.max.partition.fetch.bytes: 52428800 
         offset.syncs.enable: true
         consumer.request.timeout.ms: 60000
         offset.syncs.commit.interval.ms: 60000
         emit.offset.syncs.enabled: true
         transaction.timeout.ms: 900000     
         max.poll.interval.ms: 1800000      
         enable.idempotence: true           
         consumer.auto.offset.reset: "latest"
         
     checkpointConnector:
       tasksMax: 100
       autoRestart:
        enabled: true
       config:
         admin.timeout.ms: 120000         
         replication.policy.class: "org.apache.kafka.connect.mirror.IdentityReplicationPolicy"
         emit.checkpoints.enabled: true
         emit.checkpoints.interval.seconds: 5
         sync.group.offsets.enabled: true
         sync.group.offsets.interval.seconds: 5
         offset-syncs.topic.location: "target"  
         refresh.groups.interval.seconds: 10
         consumer.poll.timeout.ms: 200
         checkpoints.topic.replication.factor: 3
         consumer.max.poll.records: 500
         consumer.session.timeout.ms: 30000 
         offset-syncs.topic.replication.factor: 3  
         groups: .*
         
     heartbeatConnector: 
      autoRestart:
        enabled: true
      config:
        heartbeats.topic.replication.factor: 3
        emit.heartbeats.enabled: true
        
     topicsPattern: ".*"
     groupsPattern: ".*"
     
     
 resources:
   requests:
     memory: "12Gi"
     cpu: "2"
   limits:
     memory: "12Gi"
     cpu: "2"
 template:
   pod:
    metadata:
    securityContext:
       runAsGroup: 65534
       runAsNonRoot: true
       runAsUser: 65534
   connectContainer:
     securityContext:
       readOnlyRootFilesystem: true
       allowPrivilegeEscalation: false
     env:
       - name: KAFKA_HEAP_OPTS
         value: "-Xms6g -Xmx6g"
       - name: GC_LOG_ENABLED
         value: "true"
       - name: KAFKA_JVM_PERFORMANCE_OPTS
         value: "-XX:+UseG1GC -XX:MaxGCPauseMillis=100 -XX:InitiatingHeapOccupancyPercent=35 -XX:+ExplicitGCInvokesConcurrent -XX:G1HeapRegionSize=16M"
       - name: STRIMZI_JAVA_SYSTEM_PROPERTIES
         value: "-Dcom.sun.management.jmxremote=true -Dcom.sun.management.jmxremote.authenticate=false -Dcom.sun.management.jmxremote.ssl=false -Dkafka.logs.dir=/opt/kafka -Dlog4j.configuration=file:/opt/kafka/custom-config/log4j.properties"

 readinessProbe:
   initialDelaySeconds: 180
   timeoutSeconds: 15
 livenessProbe:
   initialDelaySeconds: 180
   timeoutSeconds: 15

排查与解决建议

1. 检查源端账号权限

确认MirrorMaker 2使用的源端账号NDSNQ6GM89NAB3MG是否具备DescribeGroups和DescribeGroupOffsets权限,无权限访问的消费组会直接跳过同步。

2. 优化Checkpoint Connector配置

  • 延长admin.timeout.ms至300000ms:大量消费组场景下,源端Kafka返回信息需要更长时间,避免超时中断拉取。
  • 调大consumer.poll.timeout.ms至5000ms:过短的超时可能导致消费组信息未完全拉取。
  • 提升consumer.max.poll.records至2000:减少轮询次数,提高消费组信息拉取效率。

3. 调整任务负载

当前Checkpoint Connector的tasksMax为100,面对4000个消费组可适当提升至200(需结合集群资源情况评估),分散同步压力。

4. 日志定位问题

查看MirrorMaker 2 Pod日志,搜索Group相关关键词,确认是否有消费组同步失败的报错(如权限不足、名称匹配异常等)。

5. 对比消费组列表

在源端执行kafka-consumer-groups.sh --list --bootstrap-server <source-bootstrap>,列出所有消费组后与目标端同步结果对比,排查缺失组是否为临时/过期组,或存在特殊字符导致正则匹配失效。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 15:54:54