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

