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

MirrorMaker2仅同步部分消费者组问题排查求助

问题:MirrorMaker2消费者组同步不全

通过Strimzi部署带有SourceConnector和CheckpointConnector的MirrorMaker2,将集成集群(kafka-int)的主题同步至开发集群(kafka-dev)。现有43个消费者组,仅22个被同步至目标集群,日志无报错,组和主题过滤规则配置正确。

关键线索:调试日志显示,未被同步的组在LIST_GROUPS请求中状态为Stable,但随后DESCRIBE_GROUPS请求返回状态为Dead;但通过kafkactl的admin client直接查看时,这些组状态仍为Stable,推测这是MirrorMaker2忽略它们的核心原因。


现有MirrorMaker2配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaMirrorMaker2
metadata:
  name: sync-kafka-from-int-to-dev-85
  namespace: edi-global-dev
spec:
  version: "3.3.1"
  replicas: 1
  connectCluster: "kafka-dev"
  mirrors:
  - sourceCluster: "kafka-int"
    targetCluster: "kafka-dev"
    topicsPattern: ".*"
    groupsPattern: ".*"
    checkpointConnector:
      config:
        checkpoints.topic.replication.factor: 1
        replication.factor: 1 
        sync.group.offsets.enabled: "true"
        replication.policy.class: "io.strimzi.kafka.connect.mirror.IdentityReplicationPolicy"
        refresh.groups.interval.seconds: 30
        refresh.topics.interval.seconds: 30
    sourceConnector:
      config:
        replication.factor: 1 
        offset-syncs.topic.replication.factor: 1
        offset-syncs.topic.location: target
        sync.group.offsets.enabled: "true"
        sync.topic.acls.enabled: "false" 
        replication.policy.class: "io.strimzi.kafka.connect.mirror.IdentityReplicationPolicy"
  clusters:
  - alias: "kafka-int"
    authentication:
      ...
    bootstrapServers: kafka-int.svc.cluster.local:9093
    tls:
      trustedCertificates:
      - certificate: ca.crt
        secretName: "kafka-int-cluster-ca-cert"
  - alias: "kafka-dev"
    authentication:
      ...
    bootstrapServers: kafka-dev.svc.cluster.local:9093
    config:
      config.storage.replication.factor: 1
      offset.storage.replication.factor: 1
      status.storage.replication.factor: 1
    tls:
      ...

关键日志片段

2023-04-27 14:10:46,197 DEBUG [kafka-int->kafka-dev.MirrorCheckpointConnector|worker]
 [AdminClient clientId=adminclient-11] Received LIST_GROUPS response from node 2 for 
request with header RequestHeader(apiKey=LIST_GROUPS, apiVersion=4, clientId=adminclient-
11, correlationId=4): ListGroupsResponseData(throttleTimeMs=0, errorCode=0, groups=[..., 
ListedGroup(groupId='xxx-my-consumer-group', protocolType='consumer', groupState='Stable'), ...])

2023-04-27 14:10:46,211 DEBUG [kafka-int->kafka-dev.MirrorCheckpointConnector|task-0] [AdminClient clientId=adminclient-17] 
Received DESCRIBE_GROUPS response from node 0 for request with header RequestHeader(apiKey=DESCRIBE_GROUPS, apiVersion=5, clientId=adminclient-17, 
correlationId=4): DescribeGroupsResponseData(throttleTimeMs=0, groups=[..., 
DescribedGroup(errorCode=0, groupId='xxx-my-consumer-group', groupState='Dead', 
protocolType='', protocolData='', members=[], authorizedOperations=-2147483648), ...

排查方向与原因分析

  • API版本差异导致状态解析偏差:日志中LIST_GROUPS使用API版本4,DESCRIBE_GROUPS使用API版本5,而kafkactl可能使用了不同的API版本。不同版本的Kafka Admin API对消费者组状态的定义或返回格式存在差异,导致MirrorMaker2收到的Dead状态与实际不符。
  • 跨Broker元数据同步延迟:LIST_GROUPS请求发往node2,DESCRIBE_GROUPS发往node0,kafka-int集群的broker之间可能存在消费者组元数据同步延迟,导致两个节点返回的组状态不一致。
  • 权限限制导致状态返回不准确:DESCRIBE_GROUPS返回的authorizedOperations=-2147483648表示账号无组描述权限,可能导致Broker返回的组状态信息不完整,被MirrorMaker2误判为Dead。需验证MirrorMaker2使用的账号是否拥有DescribeGroups权限。
  • MirrorMaker2的组过滤逻辑:CheckpointConnector的核心逻辑会跳过状态为Dead的消费者组,仅同步Stable、Empty等活跃状态的组。即使实际组状态正常,只要DESCRIBE_GROUPS返回Dead,就会被过滤。
  • 版本兼容性bug:当前使用Kafka 3.3.1版本,该版本的MirrorMaker2可能存在组状态判定的逻辑缺陷,比如对短暂状态波动的处理不当,导致误判组状态。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 17:15:17