设置IdentityReplicationPolicy时如何避免Kafka MM2循环复制?
我们计划基于Kafka MirrorMaker2(MM2)复刻现有MirrorMaker1的主备架构,限制条件如下:
- 应用无法配置为从多主题队列消费,因此不能采用带站点标记的主题复制方案(如
siteA.topic这类命名) - 采用双站点主备模式:siteA为正常主集群,siteB为灾备集群;每个站点部署本地Kafka集群及仅向本地生产的MM2节点,仅启用单向复制避免循环
通过配置replication.policy.class = org.apache.kafka.connect.mirror.IdentityReplicationPolicy实现了无重命名的跨集群主题复制,但切换复制方向(停止siteA→siteB,开启siteB→siteA)时,MM2会将siteA此前同步到siteB的消息再次复制回siteA,导致消息重复。已尝试调整emit.checkpoints.interval.seconds、开启group.offsets.enabled、提高sync.group.offsets.interval.seconds频率,问题仍未解决。
请问是否存在无需重复消息的解决方案,还是只能采用带标记的主题复制?
核心问题根源
使用IdentityReplicationPolicy时,MM2无法区分原始生产到siteA的消息和从siteA复制到siteB的消息。切换复制方向后,MM2会从siteB的主题起始偏移量开始同步,自然会把之前复制过去的消息再次传回siteA,造成重复。
方案1:自定义ReplicationPolicy实现消息过滤(无需修改主题名)
这是最契合你场景的方案:通过自定义ReplicationPolicy,在消息头中添加源集群标识,同时在复制逻辑中过滤掉来自目标集群的消息,避免循环同步。具体步骤:
- 实现
org.apache.kafka.connect.mirror.ReplicationPolicy接口,在复制拦截逻辑中注入并校验源集群标识 - 在消息头中添加源集群ID(如
X-Kafka-Source-Cluster: siteA) - 当MM2从siteB向siteA复制时,检查消息头的源集群标识,如果是siteA则跳过该消息
- 保持主题名称不变,应用无需修改消费配置
示例核心逻辑代码:
public class FilteredIdentityReplicationPolicy extends IdentityReplicationPolicy { private String localClusterId; @Override public void configure(Map<String, ?> props) { super.configure(props); localClusterId = (String) props.get("local.cluster.id"); } // 复制前过滤消息:如果消息来自本地集群则跳过 public boolean shouldReplicate(ConsumerRecord<byte[], byte[]> record) { Header sourceHeader = record.headers().lastHeader("X-Kafka-Source-Cluster"); if (sourceHeader == null) return true; // 无标识的原始消息允许复制 String sourceCluster = new String(sourceHeader.value()); return !localClusterId.equals(sourceCluster); } }
MM2配置中指定自定义类:
replication.policy.class = com.yourcompany.kafka.FilteredIdentityReplicationPolicy local.cluster.id = siteB # siteB的MM2节点配置为siteB,siteA的MM2配置为siteA
方案2:切换时手动指定同步起始偏移量
如果不想开发自定义代码,可通过手动干预避免重复:
- 停止siteA→siteB的MM2任务,等待所有消息同步完成(通过MM2监控指标确认同步进度)
- 用
kafka-consumer-groups.sh或Kafka Admin API,记录siteB上每个主题的最新偏移量 - 启动siteB→siteA的MM2任务时,配置复制任务从记录的最新偏移量开始同步,而非从头消费
注意:该方法需要严格的操作流程,避免遗漏主题或偏移量记录错误。
方案3:优化MM2检查点同步配置
之前的配置调整可能未到位,需确保以下配置精准:
- 设置
emit.checkpoints.interval.seconds=10(足够小的间隔,保证切换前siteB拿到siteA的最新检查点) - 确认
group.offsets.enabled=true,确保消费组偏移量同步到siteB - 切换时等待siteB的检查点完全同步后,再启动反向复制,让MM2从检查点对应的偏移量开始消费
但这种方法依赖MM2的检查点机制,极端场景下仍可能出现少量重复,可靠性不如方案1。
结论
无需被迫切换到带标记的主题复制方案,自定义ReplicationPolicy实现消息过滤是最稳定且符合你应用限制的解决方案。若不想开发自定义类,可尝试手动指定同步起始偏移量的方案,但需严格控制操作流程。
内容的提问来源于stack exchange,提问作者cheronobyl

