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

Kafka MM2跨集群消息复制出现重复问题求助

Kafka MirrorMaker2(MM2)消息重复复制问题

我有两个Kafka集群和一个Kafka Connect实例,使用MirrorMaker2做消息复制时遇到数据重复问题:每条消息的Value和Timestamp完全一致,均被重复复制。

集群信息

  • 源Kafka集群地址:xx.xx.xx.184:21001,xx.xx.xx.184:21002,xx.xx.xx.184:21003
  • 目标Kafka集群地址:xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006

现象截图

源集群消息:源集群消息
镜像后目标集群消息:镜像后目标集群消息

排查线索

新建了一个Kafka Connect实例并配置新的MM2任务,未复现数据重复问题,推测问题出在旧Connect实例,但未在日志中找到异常信息。

相关MM2配置

MirrorSourceConnector

{
  "name": "1-sc-88884abe403a44b0887107e636c43f4a",
  "config": {
    "connector.class": "org.apache.kafka.connect.mirror.MirrorSourceConnector",
    "replication.factor": "1",
    "offset-syncs.topic.replication.factor": "3",
    "producer.override.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "sync.topic.acls.enabled": "true",
    "consumer.group.id": "test_mm2_04_sc",
    "tasks.max": "1",
    "refresh.topics.enabled": "true",
    "topics": "test-04",
    "sync.topic.acls.interval.seconds": "600",
    "source.cluster.alias": "bsdca",
    "refresh.groups.enabled": "true",
    "sync.topic.configs.interval.seconds": "600",
    "source.cluster.bootstrap.servers": "xx.xx.xx.184:21001,xx.xx.xx.184:21002,xx.xx.xx.184:21003",
    "key.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "target.cluster.alias": "bsdcb",
    "value.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "offset.lag.max": "100",
    "target.cluster.security.protocol": "PLAINTEXT",
    "name": "1-sc-88884abe403a44b0887107e636c43f4a",
    "target.cluster.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "sync.topic.configs.enabled": "true",
    "source.cluster.security.protocol": "PLAINTEXT",
    "refresh.topics.interval.seconds": "20"
  },
  "tasks": [
    {
      "connector": "1-sc-88884abe403a44b0887107e636c43f4a",
      "task": 0
    }
  ],
  "type": "source"
}

MirrorCheckpointConnector

{
  "name": "1-cc-88884abe403a44b0887107e636c43f4a",
  "config": {
    "connector.class": "org.apache.kafka.connect.mirror.MirrorCheckpointConnector",
    "replication.factor": "2",
    "sync.group.offsets.interval.seconds": "20",
    "producer.override.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "consumer.group.id": "test_mm2_04_cc",
    "tasks.max": "1",
    "topics": "test-04",
    "emit.checkpoints.interval.seconds": "60",
    "source.cluster.alias": "bsdca",
    "groups": ".*",
    "source.cluster.bootstrap.servers": "xx.xx.xx.184:21001,xx.xx.xx.184:21002,xx.xx.xx.184:21003",
    "key.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "target.cluster.alias": "bsdcb",
    "value.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "target.cluster.security.protocol": "PLAINTEXT",
    "name": "1-cc-88884abe403a44b0887107e636c43f4a",
    "target.cluster.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "checkpoints.topic.replication.factor": "3",
    "refresh.groups.interval.seconds": "600",
    "source.cluster.security.protocol": "PLAINTEXT",
    "sync.group.offsets.enabled": "true"
  },
  "tasks": [],
  "type": "source"
}

MirrorHeartbeatConnector

{
  "name": "1-hc-88884abe403a44b0887107e636c43f4a",
  "config": {
    "connector.class": "org.apache.kafka.connect.mirror.MirrorHeartbeatConnector",
    "replication.factor": "2",
    "producer.override.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "tasks.max": "1",
    "heartbeats.topic.replication.factor": "3",
    "source.cluster.alias": "bsdca",
    "source.cluster.bootstrap.servers": "xx.xx.xx.184:21001,xx.xx.xx.184:21002,xx.xx.xx.184:21003",
    "key.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "target.cluster.alias": "bsdcb",
    "value.serializer": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "target.cluster.security.protocol": "PLAINTEXT",
    "name": "1-hc-88884abe403a44b0887107e636c43f4a",
    "target.cluster.bootstrap.servers": "xx.xx.xx.184:21004,xx.xx.xx.184:21005,xx.xx.xx.184:21006",
    "emit.heartbeats.interval.seconds": "1",
    "source.cluster.security.protocol": "PLAINTEXT"
  },
  "tasks": [
    {
      "connector": "1-hc-88884abe403a44b0887107e636c43f4a",
      "task": 0
    }
  ],
  "type": "source"
}

请问有没有人遇到过类似问题并解决?


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 22:04:55