使用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
相关产品推荐
相关产品推荐

