做了几年数据同步,MySQL主库到从库、到ES、到Redis、到数仓,各种折腾。今天把Canal这套东西掰开揉碎讲清楚。
Canal是阿里巴巴开源的一个中间件,核心原理是把自己伪装成一个MySQL的从库,订阅主库的Binlog日志,然后把增量变更数据解析成结构化事件,再推给下游。你不需要改业务代码、不需要双写、不用碰触发器,它就能把MySQL里的每一次insert、update、delete流转到任何一个你能想到的目标端。
这篇文章适合正在为“MySQL数据实时同步”发愁的Java后端、运维、数据工程师。我会从原理讲到部署,再从配置讲到排障,最后附上两套最常用的下游接入方式,照着做基本能跑通。
1. 为什么需要Canal:主从复制和Binlog的那点事
1.1 MySQL自带主从复制为什么不够用
MySQL原生自带主从复制,主库开Binlog,从库拉日志回放,这个机制本身很成熟。但你要注意一个关键点:原生主从是针对“MySQL到MySQL”的场景设计的,它的目的就是让从库和主库的数据保持一致。一旦你想把数据同步到MySQL之外的存储——比如Elasticsearch、Redis、ClickHouse、Kafka——原生主从就没辙了,因为它根本不知道怎么把这些Binlog事件翻译成其它系统的写入操作。
这时候有人会说,那我业务代码里同步双写不行吗?行,但代价很大。双写意味着业务逻辑侵入、性能损耗、分布式事务问题、失败补偿逻辑,一整套下来比你想象的复杂得多。而且如果你有二十张表都要同步,每个业务方法里都得写一遍,后期想加一个目标端又得全量改一遍。这显然不是合理方案。
Canal的思路是:既然MySQL的增量数据全在Binlog里,为什么不让中间层去解析Binlog,统一对外输出呢?这样业务系统完全无感知,目标端想加就加,想换就换,数据同步这件事从业务代码里彻底剥离出来。
1.2 Canal是怎么“骗”过MySQL的
这里有个很有意思的设计。MySQL主从复制的流程大致是:主库产生Binlog → 从库发送一个请求说“我要同步” → 主库给从库推送日志。Canal做的事情就是伪装成一个“假从库”,用同样的协议跟主库打交道,然后把收到的Binlog解析成自己能理解的事件对象。
具体到内部实现,Canal大致的链路是这样的:
- Canal Server启动后,连接配置好的MySQL主库,发送dump请求,声明自己是一个从库。
- MySQL主库像对待正常从库一样,把Binlog推给Canal。
- Canal接收Binlog原始字节流,解析出每一个事务、每一条DML操作,转换成Canal特有的Entry模型。
- 解析结果放进内存队列,等待客户端消费,或者通过适配器直接投递到目标端。
这个流程中,Canal做对了两件很关键的事情:一是它完整实现了MySQL的从库交互协议,所以MySQL不觉得它是“外人”;二是它把二进制日志解码成了JSON级别的结构化数据,里面包含了库名、表名、操作类型、变更前镜像、变更后镜像,下游拿到就能直接用。
这也顺带回答了一个很多人问过的问题:Canal能监听SQL Server或者Oracle吗?老实说不建议这么想,Canal的设计从编码到协议全是按MySQL的Binlog来的,它默认就是为MySQL/MariaDB准备的。如果你要同步SQL Server,那得看变更数据捕获之类的机制,那是另一套技术栈。
1.3 Binlog三种格式怎么选
MySQL的Binlog有三种格式:STATEMENT、ROW、MIXED。这是个直接影响Canal能不能正常工作的前置条件。
- STATEMENT格式记录的是SQL语句本身,比如
UPDATE users SET name='a' WHERE id=1。这种格式日志量小,但Canal解析起来会痛不欲生,因为你拿到一条SQL,还得自己去算这条SQL影响了哪些行,这等于把主库执行逻辑再模拟一遍,根本做不到。 - ROW格式记录的是每一行变更前后的具体值,比如id=1这一行,之前name是'a',之后name是'b'。这是Canal的最爱,因为有前镜像和后镜像,解析出来就是一个精准的事件。
- MIXED格式是MySQL自己判断,有的操作走STATEMENT,有的走ROW,混合状态。这种也没法保证Canal拿到的一定是ROW事件。
所以配置上有一条铁律:Canal要想稳定解析,Binlog格式必须是ROW,而且建议设置binlog_row_image=FULL,这样才能拿到完整的变更前镜像和变更后镜像。我见过有人没改binlog_row_image,默认值是FULL还好,但如果生产环境曾经调过MINIMAL,Canal拿到的前镜像就不完整,做对比和回滚就会出问题。
注意:改Binlog格式需要重启MySQL才能生效。如果线上库不能随便重启,优先找运维窗口,千万别在业务高峰期直接改。
2. 环境准备和快速部署
2.1 第一步:打开MySQL的Binlog开关
很多库默认是没开Binlog的,尤其是云数据库RDS之外的普通自建库。怎么确认?连上MySQL执行:
SHOW VARIABLES LIKE 'log_bin';如果返回的Value是OFF,那后面所有操作都白搭,Canal连上了也收不到任何数据。开启方式是在MySQL配置文件my.cnf(Windows是my.ini)里加下面几行:
[mysqld] server-id=1 log-bin=mysql-bin binlog_format=ROW binlog_row_image=FULL expire_logs_days=7这里解释一下几个参数:
server-id:必须设置,主从复制环境下每一台机器都需要唯一标识。Canal伪装从库,主库会按server-id来识别不同的复制链路,如果和目标从库的server-id冲突,可能导致同步错乱。log-bin:开启Binlog并指定日志文件前缀,比如mysql-bin,生成的文件就是mysql-bin.000001、mysql-bin.000002这种。expire_logs_days:控制Binlog保留天数。这个要重视,如果Binlog保留太短,Canal中断几天之后可能找不到需要的日志文件,无法续传,只能重新初始化位点。保留7天算是一个常见的起步配置,实际按你的数据量调整。
改完配置重启MySQL,重启完再查一次log_bin确认生效。
2.2 第二步:给Canal创建一个专用的数据库账号
Canal需要连数据库拉Binlog,但它并不需要业务库的读写权限,只需要三个权限:SELECT、REPLICATION SLAVE、REPLICATION CLIENT。为什么是这三个?因为拉Binlog走的是复制协议,需要REPLICATION相关权限;获取位点、判断当前日志位置需要REPLICATION CLIENT;某些场景解析表结构元数据需要SELECT权限。
创建语句很简单:
CREATE USER 'canal'@'%' IDENTIFIED BY 'canal_pass'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%'; FLUSH PRIVILEGES;这里有个小细节:%代表允许任意主机连接,如果你把Canal部署在和MySQL不同的机器上,这个必不可少。如果只写localhost,Canal远程根本连不过来。权限给了*.*是因为Canal同步过程中需要读取Binlog,Binlog里可能涉及多个库,如果只授权某个库,解析时遇到其它库的事件会报权限不足。
提醒:不要直接用root账号跑Canal。单独账号的好处是回收权限容易,排查问题也清晰,而且安全风险小,这个习惯建议保持。
2.3 第三步:部署Canal Server
Canal的部署有两种主流方式,Docker和压缩包直接跑。我两个都用了,给你参考。
Docker方式是最快的,适合本地折腾和快速验证:
docker pull canal/canal-server:latest docker run -d --name canal \ -p 11111:11111 \ -e canal.instance.master.address=你的MySQL地址:3306 \ -e canal.instance.dbUsername=canal \ -e canal.instance.dbPassword=canal_pass \ -e canal.instance.connectionCharset=UTF-8 \ -e canal.instance.tsdb.enable=true \ canal/canal-server:latest端口11111是Canal Server对外提供客户端接入的端口,记住不能跟其它服务冲突。
压缩包方式适合生产环境,因为你会有更多配置自由度,部署逻辑也更透明。
首先到Canal的GitHub Release页面下载对应的发行包,比如canal.deployer-1.1.7.tar.gz,解压之后目录结构是这样的:
canal.deployer-1.1.7/ ├── bin │ ├── startup.sh │ └── stop.sh ├── conf │ ├── canal.properties │ ├── logback.xml │ └── example │ └── instance.properties └── lib启动前需要改两个配置文件,一个是全局的canal.properties,一个是实例级的instance.properties。先改实例级的,因为能不能连上MySQL、监听哪个库,全看它。
conf/example/instance.properties需要改的关键项:
# MySQL主库地址 canal.instance.master.address=127.0.0.1:3306 # Binlog位点配置,先注释让它自动获取最新位点 # canal.instance.master.journal.name= # canal.instance.master.position= # 数据库账号 canal.instance.dbUsername=canal canal.instance.dbPassword=canal_pass # 连接字符集,库表是UTF-8就填UTF-8 canal.instance.connectionCharset=UTF-8 # 要监听的库和表,这里用正则表达式 # 表示监听test库的所有表 canal.instance.filter.regex=test\\..* # 表名大小写不敏感 canal.instance.filter.black.regex=注意canal.instance.filter.regex这个配置的语法是库名\\.表名,要用正则表达式的写法。比如监听所有库的所有表可以写.*\\..*,只监听test库下的user表就是test\\.user。黑名单同理,正则匹配到的表会被过滤掉。
然后启动:
cd canal.deployer-1.1.7 bin/startup.sh检查日志是排查问题的第一道工序:
tail -f logs/example/example.log如果看到类似“binlog position is xxx”这样的日志,说明Canal已经成功连接MySQL并且开始抓Binlog了。
2.4 第四步:写一个最小客户端验证数据
Canal Server跑起来了,但它自己不会“干活”,你得有一个客户端去连接它并把数据拉出来。这里先写一个最简单的Java客户端验证整个链路通不通。
新建Maven工程,引入依赖:
<dependency> <groupId>com.alibaba.otter</groupId> <artifactId>canal.client</artifactId> <version>1.1.7</version> </dependency>然后写代码:
public class QuickStartClient { public static void main(String[] args) { CanalConnector connector = CanalConnectors.newSingleConnector( new InetSocketAddress("127.0.0.1", 11111), "example", // canal instance名称 "canal", // 客户端账号,1.1.4之后可以自定义 "canal" // 客户端密码 ); connector.connect(); connector.subscribe("test\\\\.*"); while (true) { Message message = connector.getWithoutAck(100); // 获取100条数据 long batchId = message.getId(); if (batchId == -1 || message.getEntries().isEmpty()) { continue; } for (CanalEntry.Entry entry : message.getEntries()) { if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) { continue; } CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); CanalEntry.EventType eventType = rowChange.getEventType(); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (eventType == CanalEntry.EventType.INSERT) { System.out.println("INSERT: " + rowData.getAfterColumnsList()); } else if (eventType == CanalEntry.EventType.UPDATE) { System.out.println("UPDATE, before=" + rowData.getBeforeColumnsList() + ", after=" + rowData.getAfterColumnsList()); } else if (eventType == CanalEntry.EventType.DELETE) { System.out.println("DELETE: " + rowData.getBeforeColumnsList()); } } } connector.ack(batchId); // 确认消费,下次从ack的位置继续 } } }这段代码逻辑很简单:连接Canal Server,订阅test库的所有表,循环拉取日志。你在test库的任意一张表里做一次插入或更新,客户端控制台立刻会打印出对应的变更内容。
这里有一个重要的概念:ack机制。客户端调用getWithoutAck拿到一批数据,处理完成后调用ack(batchId)告诉Canal这批数据我收下了,可以更新消费位点。如果客户端没调用ack,Canal会认为这批数据没有被消费,下一次重新连接时还会再推一次。这既是保障机制,也是重复数据的一个来源,后面会细说。
3. 核心配置和参数解析
3.1 canal.properties全局配置:那些影响性能的关键项
只用默认配置跑通一个demo很简单,但要应付生产环境,必须理解canal.properties里的一些关键参数。
第一个是canal.port,Canal Server监听客户端连接的端口,默认11111。这个端口可以理解为“数据出口”,所有客户端都通过它连进来获取日志流。
第二个是内存队列大小。Canal只负责解析Binlog和投递消息,它的内存模型是一个生产者消费者模式:binlog解析线程把事件写入内存队列,客户端拉取线程从队列里取数据。队列大小由canal.instance.memory.queue.buffer.size控制,默认16384个事件。如果队列满了,Canal会阻塞解析线程,也就是“背压”机制,防止内存溢出。这个参数不是越大越好,太大内存占用高,太小吞吐上不去。一般观察客户端消费速度来调。
第三个是canal.instance.memory.batch.mode,控制批量模式,设置为MEMSIZE时按内存大小分批,设置为ITEMSIZE时按条数分批,各自对应一批数据的最大值。这个影响每次客户端拉取时能拿多少数据,调大了可以减少网络往返次数,但单批数据也会变大。
还有canal.zkServer,这是集群模式下注册到ZooKeeper的地址。如果你只部署单机,不用管它;如果要搞Canal的高可用,多个Canal Server实例抢同一个任务,就必须配ZooKeeper。
其实我的建议是:配置最好按官方默认值起步,通过监控消费速度和延迟慢慢调,不要一上来就拉满,因为你不知道你的下游能不能接住。
3.2 instance.properties实例配置:从哪个位点开始抓数据
一个Canal Server可以配置多个instance,每个instance对应一条独立的同步链路。默认的instance名字叫example,你在conf/example/instance.properties里配置的那个就是它。
最重要的概念是位点,也就是Canal当前消费到Binlog的哪个位置了。这个位置由两个信息组成:日志文件名journal.name(比如mysql-bin.000003)和偏移量position(比如234511)。
第一次启动时,如果这两项留空,Canal会从当前最新的Binlog位置开始监听。这意味着:之前已经产生的历史数据Canal是不管的,它只管启动之后的新增变更。
如果你希望从某个历史时间点开始同步,就得先找到那个时间点对应的Binlog文件和位点。可以用MySQL命令查:
SHOW BINARY LOGS; -- 查看所有Binlog文件 SHOW BINLOG EVENTS IN 'mysql-bin.000003' LIMIT 10; -- 查看指定文件的事件然后把这个信息填到canal.instance.master.journal.name和canal.instance.master.position里。这种方式有个前提:你要找的那个时间点,对应的Binlog文件还必须存在于MySQL服务器上。如果已经被purge掉了,那就没法从那个点续传,只能全量初始化。
位点是Canal同步一致性的核心。我在生产中遇到过Canal挂了一晚上,第二天重启后它自动从上次ack的位置继续拉,数据一条没丢。这个“续传”能力就是靠位点实现的,所以务必保证Conf配置里的位点信息别乱改。
3.3 关于重复消费和幂等
实时同步场景下,数据重复是家常便饭。Canal通过ack机制保证了“至少一次”的投递,但也正因为是“至少一次”,极端情况下就可能出现重复。
什么情况下会重复?客户端处理完一批数据、还没调ack的时候,客户端崩溃了。Canal Server不知道这批数据已经被处理过,等客户端重启后再拉,又把同样的数据推一遍。这就是经典的“at least once”语义。
应对重复消费的唯一解药是下游幂等。如果目标是MySQL,写一条REPLACE INTO或ON DUPLICATE KEY UPDATE就可以保证重复插入不会出错;如果目标是Redis,直接SET本来就是幂等的;如果是ES,按文档ID做upsert就行。千万不要假设消息不会重复,按这个假设写代码必踩坑。
4. 数据怎么用:两种主流接入方式
4.1 方式一:Java客户端直连消费
前面的最小客户端已经展示了直连模式。这种方式的好处是灵活,你想把数据处理成什么样子都行;坏处是要自己处理很多事情——断线重连、消费位点管理、状态监控、幂等兜底。
实际生产中用直连模式,我建议把代码结构整理成四层:
- 连接层:负责建立Canal连接,处理重连逻辑。
- 拉取层:循环调用getWithoutAck,把CanalEntry转成业务消息对象。
- 处理层:业务逻辑,比如写入目标存储、发送MQ消息。
- 确认层:处理完成后调用ack,记录位点。
我这里分享一个重连的要点。Canal连接是长连接,网络抖动或者Canal Server重启都会导致客户端断开。客户端必须要有自动重连机制,否则就是静默死亡。推荐用CanalConnectors.newCanalConnector自己管理连接,比直接用newSingleConnector更可控。
CanalConnector connector; while (running) { try { connector = CanalConnectors.newSingleConnector(addr, destination, username, password); connector.connect(); connector.subscribe(filter); while (running) { Message message = connector.getWithoutAck(batchSize); // 处理... connector.ack(message.getId()); } } catch (Exception e) { // 记录错误,延时重连 Thread.sleep(3000); } finally { if (connector != null) { connector.disconnect(); } } }只要这个外层的while循环包住内层,即使Canal Server闪断重启,客户端也会每3秒尝试重连一次,恢复后从上次ack的位置继续,不会丢数据。
4.2 方式二:Canal Adapter同步到目标库
如果不想写客户端代码,Canal官方还提供了一个叫Canal Adapter的组件,它可以直接把增量数据同步到MySQL、Elasticsearch、HBase、Redis等目标端,只需要配置映射关系,不需要写一行Java代码。
Adapter的架构是:Canal Server解析Binlog → Adapter订阅Canal Server的消息 → 根据映射规则写入目标端。
举个例子,canal.adapter-1.1.7.tar.gz解压之后,配置目录长这样:
conf/ ├── application.yml └── rdb └── mytest_user.yml在application.yml里配置Canal Server的连接信息,以及要加载哪些适配器。然后在rdb/mytest_user.yml里配置目标MySQL连接和表映射关系:
dataSourceKey: defaultDS destination: example groupId: g1 outerAdapterKey: mysql1 concurrent: true dbMapping: database: target_db table: user targetDb: target_db targetTable: user targetPk: id: id mapAll: true配置完成后,启动Adapter进程,它就会自动消费Canal Server的数据,把变更同步到目标库的target_db.user表里。这是最简单的“从MySQL到另一个MySQL”的同步方案,不需要写代码,不需要关心位点管理,Adapter自己会处理。
我实际用下来的体验是:Adapter适合快速交付和简单场景,但复杂映射(字段改名、多表关联、数据清洗)还是得自己写客户端。毕竟Adapter的映射规则就那几种,灵活度有限。
4.3 从MySQL到其它库:主流目标端玩法的思路
同步到MySQL只是最基础的一种,实际业务中更多是把数据同步到非关系型存储。我这里简单讲讲几个主流目标端的接入思路。
同步到Elasticsearch:ES的文档模型和关系型数据库差别很大。常见的做法是Canal监听业务表,拿到变更后把数据转换成ES的JSON文档,按业务主键作为_id执行index接口。关联查询怎么办?要么在同步时把关联表的数据扁平化写入同一个文档,要么用ES的父子关系或嵌套文档。数据空洞问题也要注意:如果主表和子表都有更新,同步顺序乱了可能导致ES里的数据短暂不一致,需要靠定时校对来兜底。
同步到Redis:最简单的模式是Canal监听后,直接把变更后的整行数据序列化成JSON写入Redis,key用表名+主键。更新操作直接覆盖,删除操作DEL键。这种模式很爽,因为你的应用层读Redis拿到的数据永远和MySQL里最新状态一致,只要缓存不过期不淘汰。但要注意大key问题,如果一行数据包含大字段(比如text类型存了很长的内容),序列化后可能有几百KB,放Redis里很浪费内存,同时拉大数据也有性能开销。
同步到Kafka:这是最常见的数据管道方案。Canal Server本身支持直接把消息投递到Kafka Topic,配置在canal.properties里,开启canal.mq.flatMessage=true让消息以扁平JSON格式发送。下游系统只需要消费Kafka消息就能获得数据库变更事件,完全解耦。这个方案的好处是:中间有Kafka缓冲,Canal的消费速度和下游的消费速度互不拖累,削峰填谷。
# canal.properties里开启Kafka投递 canal.serverMode=kafka kafka.bootstrap.servers=127.0.0.1:9092 kafka.topic=canal-data注意这里canal.mq.flatMessage有两种格式:true会输出扁平的JSON,类似{"id":1,"name":"张三","age":20};false会输出Canal的原生Entry格式,包含更多元信息。下游如果只是存数据,用扁平更省事;如果要做审计追踪需要变更前后镜像,就得用原生格式。
5. 常见问题与排查实录
5.1 连不上MySQL、认证失败
这个报错我见的次数最多。Canal启动后日志里出现Access denied for user 'canal'...,几乎都是权限没给够。
排查思路:
SHOW GRANTS FOR 'canal'@'%';确认有没有REPLICATION SLAVE和REPLICATION CLIENT权限。还有一个容易被忽略的点:MySQL 8.0默认的认证插件是caching_sha2_password,Canal老版本可能不支持,会报Authentication plugin 'caching_sha2_password' cannot be loaded。解决办法是创建账号时指定用mysql_native_password:
CREATE USER 'canal'@'%' IDENTIFIED WITH mysql_native_password BY 'canal_pass';5.2 连接成功了但收不到数据
这种问题比连不上更费时间。Canal正常连上MySQL,日志也显示在拉Binlog,但客户端就是收不到任何消息。最快的定位方式是在测试库随便执行一条SQL,然后去MySQL端确认Binlog事件确实产生了:
SHOW BINLOG EVENTS IN 'mysql-bin.000003' LIMIT 5;如果Binlog里有记录,那就是Canal的过滤规则把事件过滤掉了。检查canal.instance.filter.regex是不是写对了,是否匹配了正确的库名和表名。特别要注意大小写:Linux上MySQL的库表名是区分大小写的,正则也要区分大小写。
还有一个坑:Canal第一次启动默认从最新位点开始监听,如果你测试之前已经有Binlog事件了,启动之后Canal只会监听到“启动之后”的新事件。所以测试要在Canal启动之后执行,不是启动之前。
5.3 同步延迟持续增大怎么办
同步延迟是最让人头疼的。先分清是哪个环节延迟:是Canal还没从MySQL拉出来,还是拉出来了客户端消费不过来?
一条命令能看出问题:
tail -f logs/example/example.log | grep -i "delay"日志里如果有delay相关的输出,说明Canal解析出来的数据积压在内存队列里,下游消费不及。常见原因有两个:
一个是抓取线程不够。调canal.instance.channelSize和canal.instance.transaction.query.size能缓解,但这只是治标。根本原因是下游处理逻辑太慢,比如每次处理都同步调用远程接口,延迟就会越积越大。应该检查下游消费代码,把耗时的操作异步化,或者加批量处理。
另一个是单次拉取条数太小,网络往返太频繁。把客户端getWithoutAck(batchSize)的batchSize调大,比如从100调到1000,再配合canal.instance.memory.batch.mode调节批量上限,往往有奇效。
5.4 位点丢失导致从头开始
如果Canal Server宕机,重启后出现大量重复数据,很可能是位点没保存住。检查一下conf目录下有没有一个meta.dat文件,这个文件记录subscription(订阅关系)和cursor(位点)。
如果这个文件被删除或者损坏,Canal会像一个刚启动的进程一样从最新位点开始监听,之前没消费的Binlog事件全部丢失。这是一个很典型的误操作:很多人清理文件时把meta.dat当垃圾文件删了,结果重启后数据对不上。
重要:
meta.dat是Canal的命根子,禁止手动编辑。需要重置位点的时候,宁可把整个instance目录删除重建,也不要手工改meta.dat里的文件名和偏移量,格式不对直接起不来。
结尾:一些经验之谈
做Canal同步这几年,我个人的体会是:大部分“数据对不上”的故障,源头都不是Canal本身,而是上游Binlog配置不对、下游接口的幂等没做、位点管理混乱这三个地方。Canal这个组件本身非常成熟和稳定,只要你能确保它可靠地拿增量、可靠地投递消息,剩下的事情都是下游的工程问题。
最后再分享一个小技巧:给Canal做监控时,别只看CPU和内存,重点要看两个指标——消费延迟和处理吞吐量。延迟能从分钟级掉到小时级,说明下游阻塞了;吞吐量如果出现断崖式下降,说明Binlog解析遇到大事务或者网络有问题。这两板斧盯住了,实时同步链路基本不会出大篓子。