• mysql的数据表同步工具 canal的使用


    一  canal的扫盲

    1.1 canal的介绍

    canal是阿里巴巴旗下的一款开源项目,使用java语言进行开发,基于数据库增量日志解析,提供增量数据订阅与消费的功能。是一款很好用的数据库同步工具。目前只支持mysql。

    二  canal的搭建

    2.1 架构流程

     2.2 配置服务器mysql

    canal的原理是基于mysql binlog技术,所以,这里一定要开启mysql的binlog写入的功能。
    1.开启mysql服务:service mysqld  start 或  service mysql start
    2.检测binlog功能是否开启,如果是off,则没有开启,如果是on表示开启
    show variables like 'log_bin';

    3.如果binlog的显示为off,则修改配置文件  my.cnf 进行配置开启

    vi   /etc/my.cnf

    1. # set zhucongfuzhi
    2. server_id = 86               # 设置服务器编号
    3. log_bin = master-bin        # 启用二进制日志,并设置二进制日志文件前缀 
    4. expire_logs_days=7          #自动清理 7 天前的log文件,可根据需要修改
    5. binlog_format=ROW           #选择row模式

    4.重启mysql数据库

    切换到 mysql的 隶属用户:hd-mysql

    [root@localhost local]# su hd-mysql
    [hd-mysql@localhost etc]$ service mysql start
    Starting MySQL. SUCCESS! 



    重启后,再查看binlog的值,为on,则表示已经开启了。

    5.创建远程访问用户,并授权访问

    进入mysql的命令模式:

    1. create user 'canal'@'%'IDENTIFIED BY 'boc123'
    2. grant all on *.* to 'canal'@'%'
    3. flush privileges;

    mysql> create user 'canal'@'%'IDENTIFIED BY 'boc123';
    Query OK, 0 rows affected (0.01 sec)

    mysql> grant all on *.* to 'canal'@'%';
    Query OK, 0 rows affected (0.01 sec)

    mysql> flush privileges;
    Query OK, 0 rows affected (0.02 sec)

    mysql> 
     

     2.3 配置安装canal同步工具

    1.软件包下载地址

    Releases · alibaba/canal · GitHub

    2.上传软件包到服务器 

     3.解压并修改配置文件

    将软件安装到:/usr/local/ 目录下 ,完整路径为 /usr/local/canal  这个目录

    [root@localhost local]# mkdir -p canal
    [root@localhost local]# cd canal/
    [root@localhost canal]# ls
    [root@localhost canal]# pwd
    /usr/local/canal
    [root@localhost canal]# tar -zxvf /root/export/servers/canal.deployer-1.1.6.tar.gz -C .
    bin/startup.bat
    bin/restart.sh
    bin/startup.sh
    bin/stop.sh
    conf/metrics/
    conf/example/
    4.修改配置文件

    vi conf/example/instance.properties

    [root@localhost example]# pwd
    /usr/local/canal/conf/example
    [root@localhost example]# vi instance.properties 
    修改内容如下

    mysql 数据解析关注的表,Perl正则表达式.
    多个正则之间以逗号(,)分隔,转义符需要双斜杠(\\) 
    常见例子:
    1.  所有表:.*   or  .*\\..*
    2.  canal schema下所有表: canal\\..*
    3.  canal下的以canal打头的表:canal\\.canal.*
    4.  canal schema(这里的canal是数据库的名字,test1 为表名)下的一张表:canal.test1
    5.  多个规则组合使用:canal\\..*,mysql.test1,mysql.test2 (逗号分隔)
    注意:此过滤条件只针对row模式的数据有效(ps. mixed/statement因为不解析sql,所以无法准确提取tableName进行过滤) 

    3.进入bin目录下启动

    1.进入到安装目录: /usr/local/canal

    2.启动命令:  sh bin/startup.sh

    [root@localhost canal]# ./bin/startup.sh
    cd to /usr/local/canal/bin for workaround relative path
    LOG CONFIGURATION : /usr/local/canal/bin/../conf/logback.xml
    canal conf : /usr/local/canal/bin/../conf/canal.properties

    截图如下

      2.4  关闭防火墙

     2.5  编写接收java程序

    1.项目结构

     2.pom文件

    1. "1.0" encoding="UTF-8"?>
    2. <project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
    3. xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    4. <modelVersion>4.0.0modelVersion>
    5. <parent>
    6. <groupId>org.springframework.bootgroupId>
    7. <artifactId>spring-boot-starter-parentartifactId>
    8. <version>2.2.1.RELEASEversion>
    9. <relativePath/>
    10. parent>
    11. <groupId>com.ljfgroupId>
    12. <artifactId>canal-demoartifactId>
    13. <version>1.0-SNAPSHOTversion>
    14. <name>canal-demoname>
    15. <url>http://www.example.comurl>
    16. <properties>
    17. <project.build.sourceEncoding>UTF-8project.build.sourceEncoding>
    18. <maven.compiler.source>1.8maven.compiler.source>
    19. <maven.compiler.target>1.8maven.compiler.target>
    20. properties>
    21. <dependencies>
    22. <dependency>
    23. <groupId>org.springframework.bootgroupId>
    24. <artifactId>spring-boot-starter-webartifactId>
    25. dependency>
    26. <dependency>
    27. <groupId>mysqlgroupId>
    28. <artifactId>mysql-connector-javaartifactId>
    29. dependency>
    30. <dependency>
    31. <groupId>commons-dbutilsgroupId>
    32. <artifactId>commons-dbutilsartifactId>
    33. <version>1.7version>
    34. dependency>
    35. <dependency>
    36. <groupId>org.springframework.bootgroupId>
    37. <artifactId>spring-boot-starter-jdbcartifactId>
    38. dependency>
    39. <dependency>
    40. <groupId>com.alibaba.ottergroupId>
    41. <artifactId>canal.clientartifactId>
    42. <version>1.1.0version>
    43. dependency>
    44. dependencies>
    45. <build>
    46. build>
    47. project>

    3.配置文件

    1. # 服务端口
    2. server.port=10010
    3. # 服务名
    4. spring.application.name=canal-client-t14
    5. # 环境设置:dev、test、prod
    6. spring.profiles.active=dev
    7. # mysql数据库连接
    8. spring.datasource.driver-class-name=com.mysql.cj.jdbc.Driver
    9. spring.datasource.url=jdbc:mysql://localhost:3306/xx_db?serverTimezone=GMT%2B8
    10. spring.datasource.username=root
    11. spring.datasource.password=cloudiip

    4.处理类

    1. package com.ljf.canal;
    2. import com.alibaba.otter.canal.client.CanalConnector;
    3. import com.alibaba.otter.canal.client.CanalConnectors;
    4. import com.alibaba.otter.canal.protocol.CanalEntry.*;
    5. import com.alibaba.otter.canal.protocol.Message;
    6. import com.google.protobuf.InvalidProtocolBufferException;
    7. import org.apache.commons.dbutils.DbUtils;
    8. import org.apache.commons.dbutils.QueryRunner;
    9. import org.springframework.stereotype.Component;
    10. import javax.annotation.Resource;
    11. import javax.sql.DataSource;
    12. import java.net.InetSocketAddress;
    13. import java.sql.Connection;
    14. import java.sql.SQLException;
    15. import java.util.Iterator;
    16. import java.util.List;
    17. import java.util.Queue;
    18. import java.util.concurrent.ConcurrentLinkedQueue;
    19. @Component
    20. public class CanalClient {
    21. //sql队列
    22. private Queue<String> SQL_QUEUE = new ConcurrentLinkedQueue<>();
    23. @Resource
    24. private DataSource dataSource;
    25. /**
    26. * canal入库方法
    27. */
    28. public void run() {
    29. CanalConnector connector = CanalConnectors.newSingleConnector(new InetSocketAddress("192.168.152.141",
    30. 11111), "example", "canal", "boc123");
    31. int batchSize = 1000;
    32. try {
    33. connector.connect();
    34. connector.subscribe(".*\\..*");
    35. connector.rollback();
    36. try {
    37. while (true) {
    38. //尝试从master那边拉去数据batchSize条记录,有多少取多少
    39. Message message = connector.getWithoutAck(batchSize);
    40. long batchId = message.getId();
    41. int size = message.getEntries().size();
    42. if (batchId == -1 || size == 0) {
    43. Thread.sleep(1000);
    44. } else {
    45. dataHandle(message.getEntries());
    46. }
    47. connector.ack(batchId);
    48. //当队列里面堆积的sql大于一定数值的时候就模拟执行
    49. if (SQL_QUEUE.size() >= 1) {
    50. executeQueueSql();
    51. }
    52. }
    53. } catch (InterruptedException e) {
    54. e.printStackTrace();
    55. } catch (InvalidProtocolBufferException e) {
    56. e.printStackTrace();
    57. }
    58. } finally {
    59. connector.disconnect();
    60. }
    61. }
    62. /**
    63. * 模拟执行队列里面的sql语句
    64. */
    65. public void executeQueueSql() {
    66. int size = SQL_QUEUE.size();
    67. for (int i = 0; i < size; i++) {
    68. String sql = SQL_QUEUE.poll();
    69. System.out.println("[sql]----> " + sql);
    70. this.execute(sql.toString());
    71. }
    72. }
    73. /**
    74. * 数据处理
    75. *
    76. * @param entrys
    77. */
    78. private void dataHandle(List<Entry> entrys) throws InvalidProtocolBufferException {
    79. for (Entry entry : entrys) {
    80. if (EntryType.ROWDATA == entry.getEntryType()) {
    81. RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
    82. EventType eventType = rowChange.getEventType();
    83. if (eventType == EventType.DELETE) {
    84. saveDeleteSql(entry);
    85. } else if (eventType == EventType.UPDATE) {
    86. saveUpdateSql(entry);
    87. } else if (eventType == EventType.INSERT) {
    88. saveInsertSql(entry);
    89. }
    90. }
    91. }
    92. }
    93. /**
    94. * 保存更新语句
    95. *
    96. * @param entry
    97. */
    98. private void saveUpdateSql(Entry entry) {
    99. try {
    100. RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
    101. List<RowData> rowDatasList = rowChange.getRowDatasList();
    102. for (RowData rowData : rowDatasList) {
    103. List<Column> newColumnList = rowData.getAfterColumnsList();
    104. StringBuffer sql = new StringBuffer("update " + entry.getHeader().getTableName() + " set ");
    105. for (int i = 0; i < newColumnList.size(); i++) {
    106. sql.append(" " + newColumnList.get(i).getName()
    107. + " = '" + newColumnList.get(i).getValue() + "'");
    108. if (i != newColumnList.size() - 1) {
    109. sql.append(",");
    110. }
    111. }
    112. sql.append(" where ");
    113. List<Column> oldColumnList = rowData.getBeforeColumnsList();
    114. for (Column column : oldColumnList) {
    115. if (column.getIsKey()) {
    116. //暂时只支持单一主键
    117. sql.append(column.getName() + "=" + column.getValue());
    118. break;
    119. }
    120. }
    121. SQL_QUEUE.add(sql.toString());
    122. }
    123. } catch (InvalidProtocolBufferException e) {
    124. e.printStackTrace();
    125. }
    126. }
    127. /**
    128. * 保存删除语句
    129. *
    130. * @param entry
    131. */
    132. private void saveDeleteSql(Entry entry) {
    133. try {
    134. RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
    135. List<RowData> rowDatasList = rowChange.getRowDatasList();
    136. for (RowData rowData : rowDatasList) {
    137. List<Column> columnList = rowData.getBeforeColumnsList();
    138. StringBuffer sql = new StringBuffer("delete from " + entry.getHeader().getTableName() + " where ");
    139. for (Column column : columnList) {
    140. if (column.getIsKey()) {
    141. //暂时只支持单一主键
    142. sql.append(column.getName() + "=" + column.getValue());
    143. break;
    144. }
    145. }
    146. SQL_QUEUE.add(sql.toString());
    147. }
    148. } catch (InvalidProtocolBufferException e) {
    149. e.printStackTrace();
    150. }
    151. }
    152. /**
    153. * 保存插入语句
    154. *
    155. * @param entry
    156. */
    157. private void saveInsertSql(Entry entry) {
    158. try {
    159. RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
    160. List<RowData> rowDatasList = rowChange.getRowDatasList();
    161. for (RowData rowData : rowDatasList) {
    162. List<Column> columnList = rowData.getAfterColumnsList();
    163. StringBuffer sql = new StringBuffer("insert into " + entry.getHeader().getTableName() + " (");
    164. for (int i = 0; i < columnList.size(); i++) {
    165. sql.append(columnList.get(i).getName());
    166. if (i != columnList.size() - 1) {
    167. sql.append(",");
    168. }
    169. }
    170. sql.append(") VALUES (");
    171. for (int i = 0; i < columnList.size(); i++) {
    172. sql.append("'" + columnList.get(i).getValue() + "'");
    173. if (i != columnList.size() - 1) {
    174. sql.append(",");
    175. }
    176. }
    177. sql.append(")");
    178. SQL_QUEUE.add(sql.toString());
    179. }
    180. } catch (InvalidProtocolBufferException e) {
    181. e.printStackTrace();
    182. }
    183. }
    184. /**
    185. * 入库
    186. * @param sql
    187. */
    188. public void execute(String sql) {
    189. Connection con = null;
    190. try {
    191. if(null == sql) return;
    192. con = dataSource.getConnection();
    193. QueryRunner qr = new QueryRunner();
    194. int row = qr.execute(con, sql);
    195. System.out.println("update: "+ row);
    196. } catch (SQLException e) {
    197. e.printStackTrace();
    198. } finally {
    199. DbUtils.closeQuietly(con);
    200. }
    201. }
    202. }

    5.启动类

    1. package com.ljf;
    2. import com.ljf.canal.CanalClient;
    3. import org.springframework.boot.CommandLineRunner;
    4. import org.springframework.boot.SpringApplication;
    5. import org.springframework.boot.autoconfigure.SpringBootApplication;
    6. import javax.annotation.Resource;
    7. /**
    8. * Hello world!
    9. *
    10. */
    11. @SpringBootApplication
    12. public class App implements CommandLineRunner
    13. {
    14. @Resource
    15. private CanalClient canalClient;
    16. public static void main( String[] args )
    17. {
    18. System.out.println( "Hello World!" );
    19. SpringApplication.run(App.class, args);
    20. }
    21. @Override
    22. public void run(String... strings) throws Exception {
    23. //项目启动,执行canal客户端监听
    24. canalClient.run();
    25. }
    26. }

    6.启动服务

      2.6 测试验证

    1.在目的端的数据库,新建一个同样名字的数据库,同样名字的数据表。

    如这里源数据库 xx_db, 数据表tb_student;

     目的端:

     2.在源表中新增数据

     3.在目的库中查看

    4.查看console

    总结: 可以看到新增数据已经同步过来了! 

    源代码见:   https://gitee.com/jurf-liu/canal-demo.git

  • 相关阅读:
    Python的文件操作
    从 Google 离职,前Go 语言负责人跳槽小公司
    深度学习发展下的“摩尔困境”,人工智能又将如何破局?
    java计算机毕业设计ssm基于C程序课程的题库在线平台(源码+系统+mysql数据库+Lw文档)
    数据结构第二课 -----线性表之顺序表
    CSDN每日一练 |『坐公交』『盗版解锁密码』『n边形划分』2023-09-17
    神经网络有哪些基本功能,常见的神经网络有哪些
    世界坐标系、相机坐标系和图像坐标系的转换
    难点解释-理解寄主机通过虚拟网络连接到虚拟机的概念
    并查集学习笔记
  • 原文地址:https://blog.csdn.net/u011066470/article/details/126734578