Canal 中间件速通

最近在做项目的时候碰到个没见过的东西,叫 Canal,赶紧学学。

概述

我们都知道对于一个系统,数据是非常重要的,数据存储的方式又很多元,MySQL 等关系型数据库、Elastic Search、HBase、Kafka、Redis等等。而 Canal 是阿里的一个开源组件,一个十多岁的老同志了,其作用在于:同步数据库的增量数据到其他的存储应用

Canal 地址 --> https://github.com/alibaba/canal.git

Canal 定义

GitHub 上对 Canal 的介绍就一句话:

canal [kə'næl],译意为水道/管道/沟渠,主要用途是基于 MySQL 数据库增量日志解析,提供增量数据订阅和消费。

解释成大白话就是:Canal 感知到 MySQL 数据变动,然后解析变动的数据,将变动的数据发送到 MQ、同步到其他数据库,并等待进一步业务逻辑处理。
再进一步解释什么叫数据变动呢?其实就是我们对 MySQL 做一些添加操作、删除操作、还有修改操作,这些操作都会引起我们的数据发生变化。
发生变化了之后,candle 它就感知到了,一感知到他把这些数据全部解析出来,放到要放的地方去,进而起到一个数据同步的作用。看官方给的示意图:

Canal.png

工作原理

MySQL 主从复制的原理

在理解 Canal 的工作原理之前,先温习一下八股:MySQL 的主从复制。
都说 MySQL 的日志比他的数据本身更重要。其中有一种日志是 binlog(binary log 即二进制日志文件) ,主要记录了 MySQL 数据库中数据的所有变化,即数据库执行的所有 DDL 和 DML 语句。因此,我们根据主库的 MySQL binlog 日志就能够将主库的数据同步到从库中,主从复制依赖于 MySQL 的 binlog。
该过程的示意图如下:

MySQL主从.png
  • 主库将数据库中数据的变化写入到 binlog
  • 从库连接主库
  • 从库会创建一个 I/O 线程向主库请求更新的 binlog
  • 主库会创建一个 binlog dump 线程来发送 binlog ,从库中的 I/O 线程负责接收
  • 从库的 I/O 线程将接收的 binlog 写入到 relay log 中。
  • 从库的 SQL 线程读取 relay log 同步数据到本地(也就是再执行一遍 SQL )。

Canal 的工作原理

从图中可以看出,Canal 的工作原理就是把自己伪装成MySQL slave(从库),模拟MySQL slave的交互协议向MySQL Mater发送 dump协议,MySQL mater收到canal发送过来的dump请求,开始推送binary log给canal,然后canal解析binary log,再发送到存储目的地,比如MySQL,Kafka,Elastic Search等等。

Canal 历史背景

官网上对 Canal 的历史背景有过一个简要的介绍,早期阿里巴巴在杭州美国美国都部署了一个机房,那这时候就存在着一个跨机房数据同步的业务需求。最初他们采用的方式是业务的触发(trigger),但这种方式操作起来不算方便,所以在 2010年的时候,他逐步使用了一个数据库解析这种方式来取代。那这种方式有个不好的点在于,出现了很多数据库的增量订阅和消费
这也就是 Canal 诞生的背景,为了解决这一问题,大概在 2014年左右,在天猫双11的时候,就推出这个 Canal 组件,用来解决大型的促销活动中,MySQL 数据库的高并发,引起的读与写的一些问题,然后呢他又在阿里的整个集团内部进行推广和应用,并在 2017 年正式开源。
当前的 canal 支持源端 MySQL 版本包括 5.1.x , 5.5.x , 5.6.x , 5.7.x , 8.0.x

Canal 配置与启动

可参考官方的 QuickStart:https://github.com/alibaba/canal/wiki/QuickStart

这里假设已经安装好了 MySQL。

MySQL 准备

首先需要在 MySQL 终创建一个用户,并授权:

sql
复制代码
-- 使用命令登录:mysql -u root -p -- 创建用户 用户名:canal 密码:canal create user 'canal'@'%' identified by 'canal'; -- 授权 *.*表示所有库 grant SELECT, REPLICATION SLAVE, REPLICATION CLIENT on *.* to 'canal'@'%' identified by 'canal';

检查 binlog 是否开启:

shell
复制代码
mysql> show variables like 'log_bin'; +---------------+-------+ | Variable_name | Value | +---------------+-------+ | log_bin | ON | +---------------+-------+ 1 row in set, 1 warning (0.00 sec)

若没有开启,则到配置文件my.cnf或者my.ini中 设置如下信息:

properties
复制代码
[mysqld] # 启用 binlog log-bin=mysql-bin # 设置 server_id,用于标识 MySQL 实例,必须是一个唯一的整数 server-id=1 # 可选:设置 binlog 的格式,可以是 ROW、STATEMENT 或 MIXED binlog_format=ROW # 可选:指定 binlog 的保留天数,0 表示永远保留 expire_logs_days=7 # 可选:最大 binlog 文件大小 max_binlog_size=100M

然后重启 MySQL 服务,再次验证:

sql
复制代码
mysql> show variables like 'log_bin'; +---------------+-------+ | Variable_name | Value | +---------------+-------+ | log_bin | ON | +---------------+-------+ 1 row in set, 1 warning (0.00 sec) mysql> show binary logs; +------------------+-----------+-----------+ | Log_name | File_size | Encrypted | +------------------+-----------+-----------+ | mysql-bin.000001 | 561 | No | +------------------+-----------+-----------+ 1 row in set (0.00 sec) mysql> show master status; +------------------+----------+--------------+------------------+-------------------+ | File | Position | Binlog_Do_DB | Binlog_Ignore_DB | Executed_Gtid_Set | +------------------+----------+--------------+------------------+-------------------+ | mysql-bin.000001 | 561 | | | | +------------------+----------+--------------+------------------+-------------------+ 1 row in set (0.00 sec)

Canal 安装

不建议最新的版本,一般坑多

从官网下载 Canal 的安装包:https://github.com/alibaba/canal/releases

Canal下载.png


对这几个包简单介绍:

  • Canal.deployer
    • 这是核心的 Canal 部署包,包含 Canal Server,负责监听 MySQL binlog 并将数据解析后推送到目标系统。大多数情况下,都需要这个包来启动 Canal。
  • Canal.adapter
    • 这个组件是为了将 Canal 收集到的数据适配到不同的存储或消息队列中,比如 Kafka、Elasticsearch 等。如果需要将数据同步到特定的系统,则需要这个适配器。
  • Canal.admin
    • 这是 Canal 的管理工具,提供了图形界面和 API,方便对 Canal 集群进行管理和监控。如果有多个 Canal 实例或者需要更方便的管理方式,则可以考虑安装这个组件。
  • Canal.example
    • 是一个示例配置包,通常包含了一些典型的 Canal 配置示例和模板文件,帮助快速了解如何配置和使用 Canal。这个包的用途是为用户提供参考,包含一些常见场景的配置文件,比如 Canal 如何连接到 MySQL 主库、如何配置适配器将数据同步到 Kafka 或 Elasticsearch 等。你可以通过这些示例文件,快速搭建自己的 Canal 环境,或者根据示例调整自己的配置。

解压后,可以看到如下的文件夹:

Canal目录.png


进入到 conf 目录后,可看到如下的文件结构:

Canal配置目录.png
`canal.properties`中对 canal 的 Server 进行了一个整体的配置,其中有一段配置如下:
properties
复制代码
################################################# ######### destinations ############# ################################################# canal.destinations = example # conf root dir canal.conf.dir = ../conf # auto scan instance dir add/remove and start/stop instance canal.auto.scan = true canal.auto.scan.interval = 5 canal.instance.tsdb.spring.xml = classpath:spring/tsdb/h2-tsdb.xml #canal.instance.tsdb.spring.xml = classpath:spring/tsdb/mysql-tsdb.xml

这是用于指定 Canal 的destinations的,说简单点就是关于要监听哪个源的配置,上述的配置文件指定了监听 example 这个路径,从而会再具体去找 example 目录下的 instance.properties 配置文件,后者对应了具体的 MySQL 数据源。那么,如果是要监听 A 和 B 两个 MySQL 源,配置信息则应该写为
canal.destinations = A, B
,之后创建对应的 A、B 两个目录,在其中编写各自的instance.properties

修改完上述配置后,进入到conf.example修改配置文件,只用改下面的几个配置,其他的保持默认即可:

properties
复制代码
# binlog日志名称,与上一步配置一致 canal.instance.master.journal.name=mysql-bin.000001 # 与主库连接时起始的binlog偏移量 canal.instance.master.position=154 # MySQL数据解析表的黑名单,多个表用 , 隔开 canal.instance.filter.black.regex= # 数据库的用户名和密码,也可以设置常用的、root等 canal.instance.dbUsername=canal canal.instance.dbPassword=canal

Canal 启动

启动脚本位于bin目录下,如果是 Windows 环境,则需要简单修改启动脚本 start.bat,在下面这行中:

shell
复制代码
set CANAL_OPTS= -DappName=otter-canal -Dlogback.configurationFile="%logback_configurationFile%" -Dcanal.conf="%canal_conf%"

去掉 -Dlogback.configurationFile="%logback_configurationFile%"
然后双击启动脚本即可。

Java 客户端操作

Canal 支持多种语言的客户端操作:
image.png
此处以 Java 语言的操作为例。参考:https://github.com/alibaba/canal/wiki/ClientExample

依赖引入

从 Maven 仓库中拉取依赖:

xml
复制代码
<!-- https://mvnrepository.com/artifact/com.alibaba.otter/canal.client --> <dependency> <groupId>com.alibaba.otter</groupId> <artifactId>canal.client</artifactId> <version>1.1.4</version> </dependency>

编写一个 Canal 客户端代码

接下来在一个 Spring Boot 的应用中编写一个简单的客户端代码,懒人选择直接 copy 官方示例代码😂

不一定非得是 Spring Boot 应用,普通的 Java 工程也可

java
复制代码
package com.liangshou.springbootdemo.canal; import java.net.InetSocketAddress; import java.util.List; import com.alibaba.otter.canal.client.CanalConnectors; import com.alibaba.otter.canal.client.CanalConnector; import com.alibaba.otter.canal.common.utils.AddressUtils; import com.alibaba.otter.canal.protocol.Message; import com.alibaba.otter.canal.protocol.CanalEntry.Column; import com.alibaba.otter.canal.protocol.CanalEntry.Entry; import com.alibaba.otter.canal.protocol.CanalEntry.EntryType; import com.alibaba.otter.canal.protocol.CanalEntry.EventType; import com.alibaba.otter.canal.protocol.CanalEntry.RowChange; import com.alibaba.otter.canal.protocol.CanalEntry.RowData; import org.springframework.beans.factory.InitializingBean; import org.springframework.stereotype.Component; /** * @author X-L-S */ @Component public class CanalClient implements InitializingBean { private final static int BATCH_SIZE = 1000; @Override public void afterPropertiesSet() throws Exception { // 创建链接 CanalConnector connector = CanalConnectors.newSingleConnector( new InetSocketAddress(AddressUtils.getHostAddress(), 11111), "example", "", ""); try { //打开连接 connector.connect(); //订阅数据库表,全部表 connector.subscribe(".*\\..*"); //回滚到未进行ack的地方,下次fetch的时候,可以从最后一个没有ack的地方开始拿 connector.rollback(); while (true) { // 获取指定数量的数据 Message message = connector.getWithoutAck(BATCH_SIZE); //获取批量ID long batchId = message.getId(); //获取批量的数量 int size = message.getEntries().size(); //如果没有数据 if (batchId == -1 || size == 0) { try { //线程休眠2秒 Thread.sleep(2000); } catch (InterruptedException e) { e.printStackTrace(); } } else { //如果有数据,处理数据 printEntry(message.getEntries()); } //进行 batch id 的确认。确认之后,小于等于此 batchId 的 Message 都会被确认。 connector.ack(batchId); } } catch (Exception e) { e.printStackTrace(); } finally { connector.disconnect(); } } /** * 打印canal server解析binlog获得的实体类信息 */ private static void printEntry(List<Entry> entrys) { for (Entry entry : entrys) { if (entry.getEntryType() == EntryType.TRANSACTIONBEGIN || entry.getEntryType() == EntryType.TRANSACTIONEND) { //开启/关闭事务的实体类型,跳过 continue; } //RowChange对象,包含了一行数据变化的所有特征 //比如isDdl 是否是ddl变更操作 sql 具体的ddl sql beforeColumns afterColumns 变更前后的数据字段等等 RowChange rowChage; try { rowChage = RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e); } //获取操作类型:insert/update/delete类型 EventType eventType = rowChage.getEventType(); //打印Header信息 System.out.println(String.format("================》; binlog[%s:%s] , name[%s,%s] , eventType : %s", entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType)); //判断是否是DDL语句 if (rowChage.getIsDdl()) { System.out.println("================》;isDdl: true,sql:" + rowChage.getSql()); } //获取RowChange对象里的每一行数据,打印出来 for (RowData rowData : rowChage.getRowDatasList()) { //如果是删除语句 if (eventType == EventType.DELETE) { printColumn(rowData.getBeforeColumnsList()); //如果是新增语句 } else if (eventType == EventType.INSERT) { printColumn(rowData.getAfterColumnsList()); //如果是更新的语句 } else { //变更前的数据 System.out.println("------->; before"); printColumn(rowData.getBeforeColumnsList()); //变更后的数据 System.out.println("------->; after"); printColumn(rowData.getAfterColumnsList()); } } } } private static void printColumn(List<Column> columns) { for (Column column : columns) { System.out.println(column.getName() + " : " + column.getValue() + " update=" + column.getUpdated()); } } }

代码的功能比较粗暴,在控制台打印出 Canal 的信息。
之后编写一段建表的 SQL 语句:

sql
复制代码
-- 创建库 create database if not exists canal_demo; -- 切换库 use canal_demo; -- 用户表 create table if not exists user ( id bigint auto_increment comment 'id' primary key, userAccount varchar(256) not null comment '账号', userPassword varchar(512) not null comment '密码', unionId varchar(256) null comment '微信开放平台id', mpOpenId varchar(256) null comment '公众号openId', userName varchar(256) null comment '用户昵称', userAvatar varchar(1024) null comment '用户头像', userProfile varchar(512) null comment '用户简介', userRole varchar(256) default 'user' not null comment '用户角色:user/admin/ban', createTime datetime default CURRENT_TIMESTAMP not null comment '创建时间', updateTime datetime default CURRENT_TIMESTAMP not null on update CURRENT_TIMESTAMP comment '更新时间', isDelete tinyint default 0 not null comment '是否删除', index idx_unionId (unionId) ) comment '用户' collate = utf8mb4_unicode_ci;

简单结果演示

确保 MySQL 服务启动、Canal Server 启动、Java 应用启动,执行对数据库的操作,比如下面的插入一条数据:

sql
复制代码
INSERT INTO user (userAccount, userPassword, unionId, mpOpenId, userName, userAvatar, userProfile, userRole) VALUES ('testAccount', 'testPassword', 'testUnionId', 'testMpOpenId', 'testUserName', 'http://example.com/avatar.jpg', 'This is a test user.', 'user');

即可在控制台中看到打印出来的 Canal 收集的信息:

shell
复制代码
================》; binlog[mysql-bin.000002:406] , name[canal_demo,user] , eventType : DELETE id : 1 update=false userAccount : testAccount update=false userPassword : testPassword update=false unionId : testUnionId update=false mpOpenId : testMpOpenId update=false userName : testUserName update=false userAvatar : http://example.com/avatar.jpg update=false userProfile : This is a test user. update=false userRole : user update=false createTime : 2024-08-31 17:05:58 update=false updateTime : 2024-08-31 17:05:58 update=false isDelete : 0 update=false ================》; binlog[mysql-bin.000002:878] , name[canal_demo,user] , eventType : INSERT id : 2 update=true userAccount : testAccount update=true userPassword : testPassword update=true unionId : testUnionId update=true mpOpenId : testMpOpenId update=true userName : testUserName update=true userAvatar : http://example.com/avatar.jpg update=true userProfile : This is a test user. update=true userRole : user update=true createTime : 2024-08-31 19:22:27 update=true updateTime : 2024-08-31 19:22:27 update=true isDelete : 0 update=true

未完待续……

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP