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

