文章目录

  • 需求
  • 准备工作
  • 自定义复制策略
  • 编译代码

需求

使用MM2同步集群数据,topic名称不能变,默认的复制策略为:DefaultReplicationPolicy,这个策略会把同步至目标集群的topic都加上一个源集群别名的前缀,比如源集群别名为A,topic为:bi-log,该topic同步到目标集群后会变成:A.bi-log,为啥这么做呢,就是为了避免双向同步的场景出现死循环。

官方也给出了解释:

这是 MirrorMaker 2.0 中的默认行为,以避免在复杂的镜像拓扑中重写数据。 需要在复制流设计和主题管理方面小心自定义此项,以避免数据丢失。 可以通过对“replication.policy.class”使用自定义复制策略类来完成此操作,所以本文主要记录一下自定义复制策略的流程。

准备工作

下载源码

https://kafka.apache.org/downloads

kafka源码是使用Gradle编译的,需要安装Gradle,具体安装操作不赘述了,可以百度。

源码使用IDEA打开后,在connect模块下找到接口:org.apache.kafka.connect.mirror.ReplicationPolicy

自定义复制策略

ReplicationPolicy这个接口主要有几个方法:

  • formatRemoteTopic:重命名topic名称
  • topicSource:根据topic获取source集群别名
  • upstreamTopic:获取topic在source集群中的名称
  • originalTopic:获取topic原始的名称(针对多次同步过程中,被重命名过多次的topic)
  • isInternalTopic:判断是否为内部topic

根据我们的需求,自定义策略需要满足:

  • 不重命名source集群中topic的名称
  • 能返回source集群别名

实现很简单,就是保证topic原封不动即可,完整代码如下:

package org.apache.kafka.connect.mirror;import org.apache.kafka.common.Configurable;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;import java.util.Map;
import java.util.regex.Pattern;/*** Defines remote topics like "us-west.topic1". The separator is customizable and defaults to a period.*/
public class CustomReplicationPolicy implements ReplicationPolicy, Configurable {// In order to work with various metrics stores, we allow custom separators.public static final String SEPARATOR_CONFIG = MirrorClientConfig.REPLICATION_POLICY_SEPARATOR;public static final String SEPARATOR_DEFAULT = ".";private static final Logger log = LoggerFactory.getLogger(CustomReplicationPolicy.class);private String separator = SEPARATOR_DEFAULT;private Pattern separatorPattern = Pattern.compile(Pattern.quote(SEPARATOR_DEFAULT));@Overridepublic void configure(Map<String, ?> props) {if (props.containsKey(SEPARATOR_CONFIG)) {separator = (String) props.get(SEPARATOR_CONFIG);log.info("Using custom remote topic separator: '{}'", separator);separatorPattern = Pattern.compile(Pattern.quote(separator));}}/*** 拼接Topic名(if you need)** @param sourceClusterAlias 源集群标识* @param topic              源Topic名称* @return java.lang.String* @date 2023/03/03 4:28 下午*/@Overridepublic String formatRemoteTopic(String sourceClusterAlias, String topic) {return topic;}/*** 获取源集群标(source.cluster.alias)** @param topic Topic nameMirrorSourceConnector* @return source alias*/@Overridepublic String topicSource(String topic) {// 和source.cluster.alias配置的一致,可通过读取配置,为了方便直接返回return "source";}/*** 截取上游真实Topic名称** @param topic Topic name* @return java.lang.String* @date 2023/03/03 4:22 下午*/@Overridepublic String upstreamTopic(String topic) {return topic;}/*** 获取原始Topic名,没做过加工,直接返回即可** @param topic 源Topic名* @return java.lang.String* @date 2023/03/03 6:42 下午*/@Overridepublic String originalTopic(String topic) {return topic;}
}

还需要修改一个地方:org.apache.kafka.connect.mirror.MirrorSourceConnector#isCycle

这个方法是判断是否出现循环复制,会递归调用,如果不修改会死循环:

原始代码:

修改为:

    // Recurse upstream to detect cycles, i.e. whether this topic is already on the target clusterboolean isCycle(String topic) {String source = replicationPolicy.topicSource(topic);if (source == null) {return false;} else {return source.equals(sourceAndTarget.target());}}

不改的话,后果如下:

编译代码

只需要编译connect模块即可,从Gradle视图中找到对应模块的build方法,修改参数,跳过单元测试(不跳过的话电脑得卡死):

build完成之后,在项目目录下找到对应jar文件,用这两个jar文件替换掉你执行脚本所使用kafka的libs目录下的jar即可:(原jar文件记得备份,以防万一)

libs目录示例:

完成上述操作之后,修改MM2配置文件中的:

replication.policy.class=org.apache.kafka.connect.mirror.CustomReplicationPolicy

还有一种方法是直接将class文件上传到classpath下,这种方式我没试。

再次执行脚本,可以看到同步后的topic已经保持原来的名称了,大功告成!

【Kafka】MM2同步Kafka集群时如何自定义复制策略(ReplicationPolicy)相关推荐

  1. 大数据运维实战第十九课 Kafka 应用场景、集群容量规划、架构设计应用案例

    Kafka 基础与入门 1. Kafka 基本概念 Kafka 官方的定义:是一种高吞吐量的分布式发布/订阅消息系统.这样说起来可能不太好理解,这里简单举个例子:现在是个大数据时代,各种商业.社交.搜 ...

  2. web集群时session同步的3种方法

    web集群时session同步的3种方法 在做了web集群后,你肯定会首先考虑session同步问题,因为通过负载均衡后,同一个IP访问同一个页面会被分配到不同的服务器上,如果session不同步的话 ...

  3. [转载]web集群时利用memcache来同步session

    web集群时session同步的3种方法 在做了web集群后,你肯定会首先考虑session同步问题,因为通过负载均衡后,同一个IP访问同一个页面会被分配到不同的服务器上,如果session不同步的话 ...

  4. java.nio.channels.UnresolvedAddressException: null [运行storm-0.9.4集群时]

    [问题描述] 在运行storm集群时,发现kafkaspout不能消费kafka的数据,查看stormUI,没有发现有什么异常,但是手动消费kafka的数据又是正确的,通过一步一步问题排查,最好定位到 ...

  5. PXC 避免加入集群时发生SST

    环境 现有集群节点: 192.168.99.210:3101 新加入节点: 192.168.99.211:3101 通过xtrabackup备份还原实例,并通过同步方式追数据: 已有节点情况: roo ...

  6. k8s入坑之报错(9)k8s node节点加入到集群时卡住 “[preflight] Running pre-flight checks”...

    参考文档k8s node节点加入到集群时卡住 "[preflight] Running pre-flight checks"报错: k8s node节点加入到集群时卡住 " ...

  7. 搭建pxc集群时需要先安装mysql么_完美起航-高可用MySQL数据库之PXC集群

    高可用MySQL数据库之PXC集群 前言 在上一篇文章介绍了时下流行的几种数据库产品后(公众号发送"NewSQL"查看),有不少小伙伴表示对自动集群的数据库感兴趣,特别是Cockr ...

  8. mysql redis集群 同步_redis集群和redis主从同步的区别

    很多人认为redis集群就是redis主从同步,其实redis集群跟redis主从同步的机制完全不一样. 1.redis集群包含主从同步:假如你配置了6个节点的redis-server做集群,那么使用 ...

  9. redis cluster 设置密码做集群时gem下client.rb文件修改

    redis cluster 设置密码做集群时gem下client.rb文件修改 来源 https://www.cnblogs.com/shihaiming/p/5949772.html redis节点 ...

最新文章

  1. python 3 输入和输出
  2. 文献学习(part13)--A Sober Look at the Unsupervised Learning of Disentangled...
  3. 国内各大平台的推荐算法,看到360的时候笑喷了……
  4. respondsToSelector的相关使用
  5. H5学习之旅-H5的表单(11)
  6. 计算机在制造业中的应用,计算机技术在机械制造中的应用
  7. infopath转换html,Microsoft Tools to Save InfoPath Forms as HTML
  8. Ubuntu Server 16.04 安装 Redis 3.2.0
  9. Oracle11g x64使用Oracle SQL Developer报错:Unable to find a Java Virtual Machine
  10. ZBrush for Mac的插图技巧
  11. 说出来你可能不信,内核这家伙在内存的使用上给自己开了个小灶!
  12. linux安装openssl、swoole等扩展的具体步骤
  13. Linux中fork函数详解(附图解与代码实现)
  14. 计算机在酒店与应用中的展望,对酒店计算机信息管理系统的分析与展望
  15. android 刷新界面布局,Android输入法弹出刷新界面布局导致性卡顿
  16. Quartus-II之D触发器
  17. 林炳文Evankaka原创作品之mybatis的增删改查简单操作
  18. QQ空间内容同步php网站,同步 Sablog 博客日志到 Qzone
  19. 浅谈EV证书的作用及思考
  20. 中国IT风险投资机构

热门文章

  1. 小米iot业务_未来十年,小米公司的 IOT (物联网)业务预计达到 40%-50%
  2. 网络变压器厂家不传之秘:10G网络滤波器不能消除的噪音是哪些?
  3. “源于梦想、止于现实”成为图书作者的困难——《程序员羊皮卷》揭秘系列 (2)
  4. matlab用数据画热力图,Web数据可视化-手把手教你实现热力图
  5. 软件工程开发文档写作教程(05)—可行性研究报告写作规范
  6. echarts 柱状图 圆角 渐变背景 根据高度实现渐变
  7. vuex(vue的组件库)
  8. 搜狐服务架构优化实践
  9. HTML+CSS+JS网页设计期末课程大作业——上海旅游景点(10页)web前端开发技术 web课程设计 网页规划与设计...
  10. Oracle数据库迁移到MySQL