news 2026/10/5 13:07:51

canal做mysql的异步传输工具

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
canal做mysql的异步传输工具

文章目录

  • 一、前言
  • 二、思路
  • 三、canal使用前,mysql配置
  • 四、详细介绍Kafka模式
    • (1)启动kafka模式
    • (2)持久化点位(重要)
  • 五、简单介绍TCP模式
    • (1)启动tcp模式
    • (2)spring boot程序编写

一、前言

最近有一个项目,存在着从Mysql数据库同步到oracle数据库。我有不希望在主程序上编写。因此希望用同步工具。后来发现canal可以实现以上目标
canal 分为 canal‑server(抓取 MySQL binlog)、canal‑adapter(可选,同步数据库),但是因为我的数据库表结构不一样,因此本次项目只用到canal-server。

二、思路

canal-server通过抓取binlog实现文件的解析,有两种模式:

  • TCP模式:该模式提供11111端口,用Spring boot编写canal-client,实现自定义处理数据
  • Kafka模式:把数据直接推送到kafka中,然后其他程序用springboot消费kafka
    我选择用第二种模式。这样7天的kafka数据还起到了增量备份的作用。

三、canal使用前,mysql配置

  • (1) mysql必须开启日志,并且日志是row类型,添加以下两行。
[mysqld] binlog-format=ROW # 必须行模式 binlog-row-image=FULL # 必须FULL,update要有before镜像,delete要有完整行
  • (2) 新建一个用户,有replication client权限
CREATEUSERcanal@'%'IDENTIFIEDBY'Canal@123456';GRANTSELECT,REPLICATIONSLAVE,REPLICATIONCLIENTON*.*TOcanal@'%';FLUSHPRIVILEGES;

四、详细介绍Kafka模式

(1)启动kafka模式

dockerrun-d\--namecanal-server-kafka\-p11111:11111\-ecanal.serverMode=kafka\-ekafka.bootstrap.servers=192.168.1.200:9092\-ecanal.mq.topic=canal-bus-topic\-ecanal.mq.partition=0\-ecanal.instance.master.address=192.168.1.100:3306\-ecanal.instance.dbUsername=canal\-ecanal.instance.dbPassword=Canal@123456\-ecanal.instance.filter.regex=bus_db\\.bus_info_.*\-ecanal.instance.gtidon=true\canal/canal-server:v1.1.7
  • 此时 TCP 端口 11111 虽然映射,但是serverMode=kafka,TCP 客户端不能连接消费,全部消息输出 Kafka
  • 这时候去看kafka的topic,发现已经有了canal-bus-topic

(2)持久化点位(重要)

默认 docker 容器位点 meta.dat 存在容器内部,容器删除位点丢失,会重新从头消费 binlog。
生产必须挂载配置目录到宿主机,持久化 instance 位点文件。
canal 容器内部配置目录:/home/admin/canal-server/conf

  • (2.1)把容器中的配置copy出来
# 先启动临时容器,把conf拷贝出来dockerrun--rm--nametemp-canal canal/canal-server:v1.1.7truedockercptemp-canal:/home/admin/canal-server/conf /opt/canal-docker/
  • (2.2) 带挂载启动
dockerrun-d\--namecanal-server-persist\-p11111:11111\-v/opt/canal-docker/conf:/home/admin/canal-server/conf\-ecanal.serverMode=kafka\-ekafka.bootstrap.servers=192.168.1.200:9092\-ecanal.mq.topic=canal-bus-topic\canal/canal-server:v1.1.7

随便到数据库里面操作一下,看看kafka里是不是有数据

五、简单介绍TCP模式

本文的重点是kafka模式。这里稍微带过一下TCP模式

(1)启动tcp模式

dockerrun-d\--namecanal-server\-p11111:11111\-ecanal.serverMode=tcp\-ecanal.instance.master.address=192.168.1.100:3306\-ecanal.instance.dbUsername=canal\-ecanal.instance.dbPassword=Canal@123456\-ecanal.instance.gtidon=true\# 下面这个表示只定义bus_info开头的表。实际应用的时候可以不要-ecanal.instance.filter.regex=bus_db\\.bus_info_.*\canal/canal-server:v1.1.7

(2)spring boot程序编写

pom.xml

<?xml version="1.0" encoding="UTF-8"?><project><dependencies><dependency><groupId>org.springframework.boot</groupId><artifactId>spring‑boot‑starter</artifactId></dependency><!-- canal java客户端 --><dependency><groupId>com.alibaba.otter</groupId><artifactId>canal‑client</artifactId><version>1.1.7</version></dependency><!-- protobuf,canal数据序列化依赖,必须引入 --><dependency><groupId>com.google.protobuf</groupId><artifactId>protobuf‑java</artifactId><version>3.21.9</version></dependency></dependencies></project>

application.yml

canal:server-host:127.0.0.1server-port:11111destination:examplebatch-size:1000# 一次批量拉取多少条

核心代码配置类CanalConfig.java

importcom.alibaba.otter.canal.client.CanalConnector;importcom.alibaba.otter.canal.client.CanalConnectors;importorg.springframework.beans.factory.annotation.Value;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;importjava.net.InetSocketAddress;@ConfigurationpublicclassCanalConfig{@Value("${canal.server-host}")privateStringhost;@Value("${canal.server-port}")privateintport;@Value("${canal.destination}")privateStringdestination;@Bean(destroyMethod="disconnect")publicCanalConnectorcanalConnector(){CanalConnectorconnector=CanalConnectors.newSingleConnector(newInetSocketAddress(host,port),destination,"","");connector.connect();// 订阅过滤,也可以在canal‑server配置,这里客户端再次过滤connector.subscribe(".*\\..*");connector.rollback();returnconnector;}}

监听任务 CanalMessageTask.java

importcom.alibaba.otter.canal.client.CanalConnector;importcom.alibaba.otter.canal.protocol.CanalEntry;importcom.alibaba.otter.canal.protocol.Message;importorg.springframework.beans.factory.annotation.Autowired;importorg.springframework.beans.factory.annotation.Value;importorg.springframework.scheduling.annotation.EnableScheduling;importorg.springframework.scheduling.annotation.Scheduled;importorg.springframework.stereotype.Component;importjavax.annotation.PostConstruct;importjava.util.List;@Component@EnableSchedulingpublicclassCanalMessageTask{@AutowiredprivateCanalConnectorcanalConnector;@Value("${canal.batch-size}")privateintbatchSize;@PostConstructpublicvoidstartTask(){// 启动单独线程循环消费,不要用@Scheduled定时,会丢消息newThread(this::processLoop,"canal‑consumer‑thread").start();}publicvoidprocessLoop(){while(!Thread.currentThread().isInterrupted()){try{// 拉取一批消息,不阻塞Messagemessage=canalConnector.getWithoutAck(batchSize);longbatchId=message.getId();intsize=message.getEntries().size();if(batchId!=-1&&size>0){// 解析每一条binlog entryprocessEntry(message.getEntries());}// ✅消费完成提交ack,推进位点;异常不要ack,下次重启重新消费这批数据canalConnector.ack(batchId);}catch(Exceptione){e.printStackTrace();// 消费异常,回滚,下次重新拉取canalConnector.rollback();try{Thread.sleep(1000);}catch(InterruptedExceptionex){Thread.currentThread().interrupt();}}}}/** * 解析Entry,拿到行变更数据 */privatevoidprocessEntry(List<CanalEntry.Entry>entryList){for(CanalEntry.Entryentry:entryList){// 只处理ROWDATA,忽略事务开始、事务结束、DDLif(entry.getEntryType()!=CanalEntry.EntryType.ROWDATA){continue;}CanalEntry.RowChangerowChange;try{rowChange=CanalEntry.RowChange.parseFrom(entry.getStoreValue());}catch(Exceptione){thrownewRuntimeException("parse row change error",e);}Stringdatabase=entry.getHeader().getSchemaName();StringtableName=entry.getHeader().getTableName();// DDL语句,这里打印,canal客户端可以拿到DDLif(rowChange.getIsDdl()){System.out.println("DDL语句:"+rowChange.getSql());continue;}// 遍历每一行变更for(CanalEntry.RowDatarowData:rowChange.getRowDatasList()){CanalEntry.EventTypeeventType=rowChange.getEventType();System.out.printf("【%s】库:%s 表:%s%n",eventType,database,tableName);if(eventType==CanalEntry.EventType.INSERT){// insert:after列printColumns(rowData.getAfterColumnsList());}elseif(eventType==CanalEntry.EventType.UPDATE){// update:before旧值,after新值System.out.println("---before---");printColumns(rowData.getBeforeColumnsList());System.out.println("---after---");printColumns(rowData.getAfterColumnsList());}elseif(eventType==CanalEntry.EventType.DELETE){// delete:before旧值printColumns(rowData.getBeforeColumnsList());}// =====================// 【这里写你的业务逻辑】// 1. 判断database、tableName过滤(bus_info_xxxx分表)// 2. 取出rowData里面字段,做类型转换// 3. 组装实体,调用Mapper写入Oracle / 其他逻辑// =====================}}}/** * 打印列名和值 */privatevoidprintColumns(List<CanalEntry.Column>columns){for(CanalEntry.Columncol:columns){System.out.print(col.getName()+"="+col.getValue()+" ");}System.out.println();}}
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/5 13:06:37

Wand-Enhancer 完整指南:三步本地解锁 Wand 免费版 2 小时限制

Wand-Enhancer 完整指南&#xff1a;三步本地解锁 Wand 免费版 2 小时限制 【免费下载链接】Wand-Enhancer Advanced UX and interoperability extension for Wand (WeMod) app 项目地址: https://gitcode.com/GitHub_Trending/we/Wand-Enhancer 每天免费使用 Wand 满两…

作者头像 李华
网站建设 2026/10/5 13:02:59

终极指南:快速掌握Ludusavi游戏存档备份的5个核心技巧

终极指南&#xff1a;快速掌握Ludusavi游戏存档备份的5个核心技巧 【免费下载链接】ludusavi Backup tool for PC game saves 项目地址: https://gitcode.com/GitHub_Trending/lu/ludusavi &#x1f3ae; 作为PC游戏玩家&#xff0c;你是否曾经因为系统重装、游戏重装或…

作者头像 李华
网站建设 2026/10/5 12:58:38

高精度模拟加法减法

提示&#xff1a;文章写完后&#xff0c;目录可以自动生成&#xff0c;如何生成可参考右边的帮助文档 高精度模拟加法减法前言一、高精度加法的实现二、高精度减法的实现1.总结1、比如减法实现里面的删除前导零没有考虑到结果全是零的情况2、减法实现里没有实现带有前导零的输入…

作者头像 李华
网站建设 2026/10/5 12:58:15

工业级存储方案:MRAM与FlexNVM EEPROM实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/5 12:57:46

AMD SEV机密计算:虚拟机内存加密原理与KVM/QEMU部署实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华