canal做mysql的异步传输工具

发布时间:2026/10/5 13:08:08
canal做mysql的异步传输工具 文章目录一、前言二、思路三、canal使用前mysql配置四、详细介绍Kafka模式1启动kafka模式2持久化点位重要五、简单介绍TCP模式1启动tcp模式2spring 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-formatROW # 必须行模式 binlog-row-imageFULL # 必须FULLupdate要有before镜像delete要有完整行(2) 新建一个用户有replication client权限CREATEUSERcanal%IDENTIFIEDBYCanal123456;GRANTSELECT,REPLICATIONSLAVE,REPLICATIONCLIENTON*.*TOcanal%;FLUSHPRIVILEGES;四、详细介绍Kafka模式1启动kafka模式dockerrun-d\--namecanal-server-kafka\-p11111:11111\-ecanal.serverModekafka\-ekafka.bootstrap.servers192.168.1.200:9092\-ecanal.mq.topiccanal-bus-topic\-ecanal.mq.partition0\-ecanal.instance.master.address192.168.1.100:3306\-ecanal.instance.dbUsernamecanal\-ecanal.instance.dbPasswordCanal123456\-ecanal.instance.filter.regexbus_db\\.bus_info_.*\-ecanal.instance.gtidontrue\canal/canal-server:v1.1.7此时 TCP 端口 11111 虽然映射但是serverModekafkaTCP 客户端不能连接消费全部消息输出 Kafka这时候去看kafka的topic发现已经有了canal-bus-topic2持久化点位重要默认 docker 容器位点 meta.dat 存在容器内部容器删除位点丢失会重新从头消费 binlog。生产必须挂载配置目录到宿主机持久化 instance 位点文件。canal 容器内部配置目录/home/admin/canal-server/conf2.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.serverModekafka\-ekafka.bootstrap.servers192.168.1.200:9092\-ecanal.mq.topiccanal-bus-topic\canal/canal-server:v1.1.7随便到数据库里面操作一下看看kafka里是不是有数据五、简单介绍TCP模式本文的重点是kafka模式。这里稍微带过一下TCP模式1启动tcp模式dockerrun-d\--namecanal-server\-p11111:11111\-ecanal.serverModetcp\-ecanal.instance.master.address192.168.1.100:3306\-ecanal.instance.dbUsernamecanal\-ecanal.instance.dbPasswordCanal123456\-ecanal.instance.gtidontrue\# 下面这个表示只定义bus_info开头的表。实际应用的时候可以不要-ecanal.instance.filter.regexbus_db\\.bus_info_.*\canal/canal-server:v1.1.72spring boot程序编写pom.xml?xml version1.0 encodingUTF-8?projectdependenciesdependencygroupIdorg.springframework.boot/groupIdartifactIdspring‑boot‑starter/artifactId/dependency!-- canal java客户端 --dependencygroupIdcom.alibaba.otter/groupIdartifactIdcanal‑client/artifactIdversion1.1.7/version/dependency!-- protobufcanal数据序列化依赖必须引入 --dependencygroupIdcom.google.protobuf/groupIdartifactIdprotobuf‑java/artifactIdversion3.21.9/version/dependency/dependencies/projectapplication.ymlcanal:server-host:127.0.0.1server-port:11111destination:examplebatch-size:1000# 一次批量拉取多少条核心代码配置类CanalConfig.javaimportcom.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(destroyMethoddisconnect)publicCanalConnectorcanalConnector(){CanalConnectorconnectorCanalConnectors.newSingleConnector(newInetSocketAddress(host,port),destination,,);connector.connect();// 订阅过滤也可以在canal‑server配置这里客户端再次过滤connector.subscribe(.*\\..*);connector.rollback();returnconnector;}}监听任务 CanalMessageTask.javaimportcom.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;ComponentEnableSchedulingpublicclassCanalMessageTask{AutowiredprivateCanalConnectorcanalConnector;Value(${canal.batch-size})privateintbatchSize;PostConstructpublicvoidstartTask(){// 启动单独线程循环消费不要用Scheduled定时会丢消息newThread(this::processLoop,canal‑consumer‑thread).start();}publicvoidprocessLoop(){while(!Thread.currentThread().isInterrupted()){try{// 拉取一批消息不阻塞MessagemessagecanalConnector.getWithoutAck(batchSize);longbatchIdmessage.getId();intsizemessage.getEntries().size();if(batchId!-1size0){// 解析每一条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(ListCanalEntry.EntryentryList){for(CanalEntry.Entryentry:entryList){// 只处理ROWDATA忽略事务开始、事务结束、DDLif(entry.getEntryType()!CanalEntry.EntryType.ROWDATA){continue;}CanalEntry.RowChangerowChange;try{rowChangeCanalEntry.RowChange.parseFrom(entry.getStoreValue());}catch(Exceptione){thrownewRuntimeException(parse row change error,e);}Stringdatabaseentry.getHeader().getSchemaName();StringtableNameentry.getHeader().getTableName();// DDL语句这里打印canal客户端可以拿到DDLif(rowChange.getIsDdl()){System.out.println(DDL语句:rowChange.getSql());continue;}// 遍历每一行变更for(CanalEntry.RowDatarowData:rowChange.getRowDatasList()){CanalEntry.EventTypeeventTyperowChange.getEventType();System.out.printf(【%s】库:%s 表:%s%n,eventType,database,tableName);if(eventTypeCanalEntry.EventType.INSERT){// insertafter列printColumns(rowData.getAfterColumnsList());}elseif(eventTypeCanalEntry.EventType.UPDATE){// updatebefore旧值after新值System.out.println(---before---);printColumns(rowData.getBeforeColumnsList());System.out.println(---after---);printColumns(rowData.getAfterColumnsList());}elseif(eventTypeCanalEntry.EventType.DELETE){// deletebefore旧值printColumns(rowData.getBeforeColumnsList());}// // 【这里写你的业务逻辑】// 1. 判断database、tableName过滤bus_info_xxxx分表// 2. 取出rowData里面字段做类型转换// 3. 组装实体调用Mapper写入Oracle / 其他逻辑// }}}/** * 打印列名和值 */privatevoidprintColumns(ListCanalEntry.Columncolumns){for(CanalEntry.Columncol:columns){System.out.print(col.getName()col.getValue() );}System.out.println();}}