三年前接到订单中心拆分需求时,我就面临一个选择:用 Canal 实时从 MySQL 向其它库同步数据,还是继续靠定时脚本硬扛。核心需求是 A 库订单主表变更后,10 秒内要出现在 B 库,并且不碰业务代码。当时团队里有两种主流声音,一是给每张表加 updated_at 做增量轮询,二是直接上消息队列让业务主动发变更事件。这两种方案我都不太满意:轮询有定时扫描的延迟,还需要容忍脏读;业务主动发事件等于把所有写路径都改一遍,风险太大。
后来我注意到 MySQL 主从复制里那个很成熟的 binlog 协议——如果我们自己能消费 binlog,不就等于拿到了数据库每一次变更的完整清单吗?带着这个思路,我跑了小半年 Canal,从单实例同步到 Adapter 落地,中间踩了不少坑,也把原理翻来覆去捋了好几遍。这篇文章就是那段经历的系统整理,适合正在做数据中台、分库分表迁移、缓存更新、搜索引擎索引同步的开发者,尤其是第一次接触 Canal 的人。
我会先把 Canal 的能力边界和原理讲清楚,再给出我在生产环境验证过的 MySQL 端配置、Canal 部署参数、Adapter 同步链路,最后整理高频报错和排查思路。文章里所有配置都来自真实的项目实例,你可以直接照着抄,但前提是真的理解了每项配置在干什么。
1. 先说清楚:Canal 能解决什么问题,不能解决什么问题
1.1 我是在什么场景下盯上 Canal 的
当时那个订单中心拆分项目里,订单表大概每天新增两百万行,业务方希望 A 库写入完成后,B 库能在秒级拿到同一份数据,用来做实时查询和报表。最原始的方案当然是用 DataX 每小时跑一次同步,但业务方明确说不行,因为报表要看分钟级趋势,而且一旦源库做了更新,目标库是空的,整个查询链路就废了。
我在技术选型时列过三个候选:MySQL 原生主从复制、基于应用层的 MQ 消息、Canal 解析 binlog。原生主从复制只能把一个 MySQL 同步到另一个 MySQL,没法把变更同时送给 ClickHouse、Elasticsearch;MQ 消息虽然灵活,但要求业务代码在事务提交后自行发消息,遇到事务回滚会很尴尬——消息发出去了,数据库却回滚了,下游拿到的是假变更。Canal 的定位刚好卡在中间:它不侵入业务,不改表结构,只借助 MySQL 主从复制协议拿到 binlog,然后把变更事件翻译成结构化数据,再以 TCP 或 MQ 的形式推给任意下游。
那次项目之后,我又在两个场景里用过 Canal:一个是把 MySQL 订单表同步到 ClickHouse 做实时宽表,另一个是把商品表同步到 Elasticsearch 做搜索索引。可以说,只要你能接受"最终一致性"这个前提,Canal 能覆盖大多数跨库增量同步需求。
1.2 一次"伪从库"把自己塞进复制链路
Canal 为什么能做到这些?核心是模拟 MySQL 从库。熟悉复制原理的同学都知道,MySQL 主从复制里,从库会先向主库注册自己的 server-id,然后发送COM_BINLOG_DUMP指令,主库就会源源不断地把 binlog 推给从库。Canal 做的事情,就是用 Java 重新实现了这个"从库"协议。
它不需要你给业务表增加字段,也不需要业务方配合调用 API。主库提交事务、写 binlog、每个 binlog 事件被解析成一条或多条行变更,Canel 把这些行变更按事务边界打包成Insert、Update、Delete三类事件,再通过 Tcp 端口、Kafka 或者 RocketMQ 送出去。整个过程相当于把数据库内部的复制能力"外包"了,下游接不接都不影响 MySQL 本身。
理解这一点后,很多配置项就好解释了:为什么 Canal 需要 MySQL 开 binlog?为什么不支持 STATEMENT 格式?为什么必须给一个复制账号?答案都指向同一个事实——Canal 不是靠 SQL 查询去抓增量,而是靠解析二进制日志来还原数据变更。
1.3 它能做和做不到的事
我做选型时会先列一份能力清单,帮你省掉踩坑时间。
Canal 能做的:
- 从 MySQL 实时增量同步到 MySQL、ClickHouse、Elasticsearch、HBase 等目标端;
- 将变更事件投递到 Kafka / RocketMQ,供自研消费者处理;
- 支持单表单库、过滤正则、字段映射和自定义转换;
- 配合 Canal Admin 做集群切换和位点管理。
Canal 做不到的:
- 不能做历史数据全量初始化,全量数据得用 DataX、mysqldump 导一次;
- 不负责目标端表结构自动创建,尤其 ES 索引和 ClickHouse 表,最好预先建好;
- 不做复杂多表关联清洗,如果有这个需求,请把事件投到 MQ 后自研处理;
- 不保证下游重复消费时的幂等,这需要目标端配合唯一键和覆盖写入。
所以 Canal 是一个"变更事件搬运工",不是一个完整 ETL 框架。你在规划同步链路时,先把"增量同步"和"全量初始化"两条路径分开,后面会轻松很多。
2. 同步原理拆解:为什么必须依赖 binlog
2.1 从 MySQL 的三种日志里找到正确答案
MySQL 里至少有三种日志,很多人容易混:redo log 是 InnoDB 存储引擎的物理日志,用于崩溃恢复;undo log 用于事务回滚和 MVCC;binlog 是 Server 层的逻辑日志,记录的是"哪张表、哪个主键、被改成什么",主要服务于主从复制和数据恢复。
Canal 只关心 binlog。原因很简单:redo log 只属于 InnoDB,并且它是物理页层面的变更,其他引擎的表根本没记录;undo log 是内部结构,外部程序很难安全消费。而 binlog 是 MySQL 对外输出的标准接口,有固定的事件格式,主从复制、Canal、CDC 工具全部建立在这条通道上。
所以 MySQL 开启 binlog,不是 Canal 的"建议配置",而是"硬前提"。如果你的数据库从没开过 binlog,那 Canal 无论如何也拿不到增量数据,这属于第一步就决定了路线能不能走通的问题。
2.2 Canal 作为"伪从库"的完整工作过程
我梳理了 Canal 处理一条 INSERT 的全过程,这是理解后续所有参数的基础。
第一步,业务事务提交,InnoDB 将变更写入 binlog。第二步,Canal 启动时先向源 MySQL 发送注册命令,带上自己配置的 server-id。第三步,主库响应后,Canal 发送COM_BINLOG_DUMP指令,指定从哪个 binlog 文件、哪个 offset 开始拉取。第四步,binlog 解析器将RowsEvent解析为InsertRowsEvent、UpdateRowsEvent、DeleteRowsEvent,同时带上表名、数据库名、主键位置以及变更前后行数据。第五步,Canal 根据事务边界,把同事务内的事件打包成一个 batch,写入内存环形缓冲区,再由 dispatcher 推送到客户端或者 MQ。
值得一提是,Canal 拉取 binlog 用的是原生的复制协议,而不是 JDBC 的SELECT,所以你不用担心它会拖垮源库的查询性能。解析过程发生在 Canal 进程内,CPU 消耗主要在解析事件上,生产环境一般需要给 Canal 单独分配机器。
2.3 ROW 格式与 STATEMENT 的差别,以及为什么非它不可
binlog 有STATEMENT、ROW、MIXED三种格式,Canal 正确工作的前提是ROW格式。这个不是偏好,是硬约束。
STATEMENT 格式记录的是原始 SQL,比如UPDATE user SET age = age + 1 WHERE id > 100。Canal 如果只拿到这条 SQL,它根本不知道具体哪些行发生了什么变化,更不可能把变更事件拆成一条条数据记录。就算可以执行 SQL 去反推,也得重新连接源库,代价太高且容易产生更大偏差。
ROW 格式记录的是每一行数据变更前后图像:某一行 UPDATE 前 id=100、age=20,UPDATE 后 id=100、age=21。Canal 拿到这样的原始事件,就能稳定输出字段级别的增量。MIXED 格式虽然在某些情况下会切到 ROW,但为了稳定,我建议直接强制用 ROW。
另外提一句binlog_row_image,我默认设为FULL,这样 UPDATE 事件会带上完整的前像和后像。如果你把它调成MINIMAL,binlog 会只记录修改过的字段以及唯一键字段,日志量更小,但下游做复杂清洗时可能缺少旧值字段。对于同步到普通 MySQL 表这种场景,FULL最省心。
2.4 GTID 模式下的位点策略
现在的 MySQL 8.0 安装时基本默认开启gtid_mode=ON。GTID 的全称是 Global Transaction Identifier,每个事务都有全局唯一编号,主从复制可以通过执行过的 GTID 集合来避免重复应用事务。
Canal 1.1.x 对 GTID 有专门支持。当源 MySQL 开启 GTID 后,建议在 instance.properties 里设置canal.instance.gtidon=true,这样 Canal 的位点不再是简单的(filename, position),而是能感知 GTID 集合,failover 时更不容易错乱。
如果你用的是 MySQL 5.7,并且已经有其他从库在跑,GTID 的配置要和现有复制架构保持一致,不要单独给 Canal 开一套不对称玩法。强一致性环境里,复制链路的所有参与者最好用同一套规则。
3. MySQL 端准备:binlog 开启、账号权限与常见配置错误
3.1 修改 my.cnf 的正确姿势与重启验证
不管是用 RPM、Docker 还是源码安装,MySQL 的 binlog 配置都集中在my.cnf。最核心的几项如下:
[mysqld] log-bin=mysql-bin binlog_format=ROW binlog_row_image=FULL server-id=100 expire_logs_days=7 max_binlog_size=256M # 如果确定要开启 GTID gtid_mode=ON enforce_gtid_consistency=ONserver-id必须保证全局唯一,不能和现有主从集群里的任何节点重复。如果 MySQL 只有单一节点,给个 100 也完全可以。expire_logs_days建议按你的数据保留策略设置,太短会导致 Canal 长时间停机后位点失效,太长会占用磁盘空间。个人建议至少 7 天,能扛得住一次周末维护。
改完配置后重启 MySQL,然后执行:
SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format'; SHOW MASTER STATUS;如果第一行输出ON,格式是ROW,SHOW MASTER STATUS能查到 binlog 文件和 position,说明 MySQL 端已经具备条件。注意,Docker 部署的 MySQL 要用docker exec进入容器去改/etc/mysql/my.cnf,然后重启容器,别在宿主机上修改一个容器根本不读的配置文件。
3.2 创建同步账号:权限最小化
Canal 需要一个专门账号来拉 binlog,不要直接拿 root 给 Canal 用,这是我最想强调的安全习惯。账号权限只需要三个:SELECT、REPLICATION SLAVE、REPLICATION CLIENT。
CREATE USER 'canal'@'%' IDENTIFIED BY 'canalpass'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%'; FLUSH PRIVILEGES;REPLICATION SLAVE是复制协议通信必需的,REPLICATION CLIENT用于查询主库状态和 binlog 坐标,SELECT用于 Canal 内部做位点校验和表结构分析。如果你的源库是 MySQL 8.0,默认认证插件是caching_sha2_password,Canal 1.1.5 之后的版本已经支持,但如果你用的 Canal 版本偏老,或者目标端 JDBC 驱动不兼容,可以在建账号时指定WITH mysql_native_password:
CREATE USER 'canal'@'%' IDENTIFIED WITH mysql_native_password BY 'canalpass';开发环境无所谓,生产环境把这个账号的 Host 尽可能限制成 Canal 服务器的内网 IP,而不是%。同步账号一旦泄露,就能读取所有库的 binlog,敏感数据风险你懂的。
3.3 容易现场翻车的三个配置错误
先说server-id冲突。如果这台 MySQL 已经有从库,从库的 server-id 是 101,你给 Canal 也配置了 101,那么 Canal 连接后,MySQL 会因为重复 server-id 把已存在的从库连接踢掉,然后 Canal 自己也可能拿不到完整 binlog。排障时如果发现原从库频繁Got fatal error 1236,优先检查是不是 Canal 的 server-id 撞车了。
其次是 binlog 清理策略。很多人只设了expire_logs_days=3,结果 Canal 实例因为业务发布停机两天半,启动时直接报找不到 binlog 文件,位点失效。位点失效后的恢复不是小事,稍后我会专门讲恢复链路。
最后是大小写敏感性。MySQL 的lower_case_table_names参数如果和 Canal 解析时依赖的库表名大小写不一致,很容出现目标端找不到表的报错。Linux 下默认是 0,也就是区分大小写,Windows 下默认是 1。如果你的源库表名是大写,目标库是小写,Adapter 映射时又不做处理,同步就会失败。我建议所有环境统一用一个规范:库名和表名一律小写,从源头消灭这类问题。
4. Canal 服务端部署与核心配置逐项说明
4.1 下载、目录结构与启动自检
Canal 官方发布分为 deployer、adapter、admin 和 client 几部分。做单机增量同步只需要 deployer 和 adapter。我用得最顺的版本是 1.1.7,这个版本对 MySQL 8.0 和 GTID 支持都算稳定,下载地址在这里不多说,按官方 release 选择即可。
解压后目录结构如下:
canal.deployer-1.1.7/ ├── bin │ ├── startup.sh │ └── stop.sh ├── conf │ ├── canal.properties │ ├── logback.xml │ └── example │ └── instance.properties ├── lib └── logs启动前先看一眼conf/canal.properties,这里有三组关键信息:监听端口、server 模式、批次大小。默认canal.port=11111是 TCP 端口,客户端从这个端口接收数据;canal.serverMode=tcp表示直接通过 TCP 推给 Adapter,如果走 Kafka 就改成kafka,并填上 broker 地址。
启动时执行bin/startup.sh,然后看日志:
tail -f logs/canal/canal.log tail -f logs/example/example.log看到Canal Launcher has started和start successful这类日志,基本就启动成功了。如果日志里出现canal: no suitable driver或数据库连接失败,请先回头检查 MySQL 配置,服务器起来了不代表已完成同步注册。
4.2 canal.properties 里和同步链路直接相关的配置
canal.properties是整个 Canal 服务的全局配置,我不建议你盲改,只挑生产中真正会动的参数讲。
canal.port是服务端口,默认 11111,客户端和 Adapter 都靠它连接。canal.destinations是 instance 列表,默认example,对应conf/example目录。如果你有多个业务库需要拆开管理,可以新建多个 instance 目录,并用逗号分隔。canal.serverMode决定事件推送到哪,tcp模式下 Adapter 用 TCP 拉取;kafka模式下 Canal 会把事件写入指定的 Kafka topic,适合下游有独立消费组的场景。
canal.instance.global.mode控制位点存储方式,默认memory,意味着 Canal 重启后会从当前 binlog 的最新位置开始拉取,这样可能丢失停机期间的变更。生产环境建议改为zookeeper,并配置canal.zkServers=zk1:2181,zk2:2181,这样实例位点存在 ZK 里,故障切换时才能恢复。
4.3 instance.properties 配置逐行拆解
每个 instance 目录下有一个instance.properties,这才是真正定义"从哪儿拉、拉什么表"的文件。
# 源 MySQL 地址 canal.instance.master.address=127.0.0.1:3306 # 同步账号 canal.instance.dbUsername=canal canal.instance.dbPassword=canalpass canal.instance.connectionCharset=UTF-8 # 是否开启 GTID,源库开启 GTID 时需要设为 true canal.instance.gtidon=false # 解析器缓冲大小,默认 16384 canal.instance.memory.buffer.size=16384 # 库表过滤:正则匹配,默认全部 canal.instance.filter.regex=.*\\..*filter.regex特别值得讲。默认.*\..*表示同步所有库所有表。生产环境我强烈建议精确到库,比如只同步订单库下所有表:
canal.instance.filter.regex=order_db\\..*正则语法遵循 Java 正则,表名之间用竖线分隔。注意点号前要写双反斜杠,因为 Java 字符串里\\.才表示一个普通点号。用错这个配置不会立刻报错,而是同步过去的表比预期少,这个坑我踩过。
4.4 启动自检的观察点:日志和端口
配置完成后,不要急着接 Adapter,先用自带命令验证是否已经拉到 binlog。看instance.log是否持续出现类似batchId : xxx的日志,这些就是解析出来的 binlog 事件。再看canal.log是否出现success register等关键信息。
你还可以用 Canal 官方提供的CanalClient示例代码跑一个 Java 小程序,连上 11111 端口打印事件。这里贴一个简单的连接逻辑,能兜底验证链路通不通:
CanalConnector connector = CanalConnectors.newSingleConnector( new InetSocketAddress("127.0.0.1", 11111), "example", "", null); connector.connect(); connector.subscribe(".*\\..*"); while (true) { Message message = connector.getWithoutAck(100); long batchId = message.getId(); if (batchId == -1) { Thread.sleep(1000); continue; } for (CanalEntry.Entry entry : message.getEntries()) { System.out.println(entry.toString()); } connector.ack(batchId); }看到entry里有行数据,说明 binlog 解析链路完全打通。这时候再往下游接,问题定位范围就小多了。
5. 基于 Canal Adapter 实现跨库同步:MySQL 到 MySQL 的完整链路
5.1 为什么直接用 Adapter,而不是自研消费者
如果你的最终目标就是把 MySQL 表同步到另一个 MySQL 表,Canal Adapter 是最快落地的方式。它相当于一个内置消费者:从 Canal 获取事件,按映射规则拼接目标 SQL,然后写入目标库。
自研消费者的优势是灵活,但代价是要自己管理连接、位点、失败重试、顺序保证。前期调试就能耗掉一到两周,而 Adapter 在单表单库场景下几乎开箱即用。所以我给团队定的原则是:简单镜像同步用 Adapter,复杂清洗或者多下游分发再走上游 MQ。
当然 Adapter 也不是没有坑。最典型的问题是它只做增量事件落库,不会反向创建表,也不会同步源表的历史数据。所以上线前必须先用某一种方式把存量数据导到目标库,然后才能开启 Adapter 增量。
5.2 Adapter 环境准备与 application.yml 配置
Adapter 的包名类似canal.adapter-1.1.7.tar.gz,解压后核心配置是conf/application.yml。先配置数据源和 canal 连接。
server: port: 8081 canal.conf: mode: tcp canalServerHost: 127.0.0.1:11111 batchSize: 500 syncBatchSize: 1000 retries: 0 timeout: 30000 srcDataSources: defaultDS: url: jdbc:mysql://127.0.0.1:3306/order_db?useUnicode=true&serverTimezone=Asia/Shanghai username: canal password: canalpassmode: tcp表示 Adapter 作为 TCP 客户端连接 Canal 服务端。canalServerHost是 Canal 的端口。srcDataSources这个名字容易误导,它其实配的是源 MySQL 连接,Adapter 需要连源库做表结构的一些辅助查询。batchSize是每次从 Canal 拉取的批次行数,syncBatchSize是写入目标端的最大批大小。目标库的配置不放在这个文件里,而是放在具体映射文件中。
启动 Adapter 用bin/startup.sh,默认会读取conf/application.yml和conf/rdb/*.yml。如果日志里出现canal adapter started,说明已经连接成功。
5.3 同步到 MySQL:写一个可复用的表映射
假设要把order_db.user同步到report_db.user_bak,在 Adapter 的conf/rdb目录下新建一个user.yml:
dataSourceKey: defaultDS destination: example groupId: g1 outerAdapterKey: mysql1 concurrent: true dbMapping: database: order_db table: user targetTable: report_db.user_bak targetColumns: id: id name: name age: age create_time: create_time primaryKey: iddataSourceKey要和application.yml里的数据源 key 对应;destination对应 Canal 的 instance 名称;dbMapping.database和table是源库表,targetTable是目标库表。注意这里我写的是report_db.user_bak,Adapter 会自动拼接目标库前缀,前提是目标库已经存在。
primaryKey必须明确指定。Canal 解析出的 UPDATE 事件依赖主键定位目标行,如果目标表没有主键或者映射里没配主键,Adapter 无法生成带WHERE id = ?的更新 SQL,更新就会丢或错。这个配置在 MySQL 目标端是硬性的。
5.4 验证一致性:一条 INSERT 走到底的观察
配置完成后,我习惯先做一个最小验证。
先在源库执行:
INSERT INTO order_db.user (id, name, age, create_time) VALUES (10001, '测试用户', 18, NOW());再查目标库:
SELECT * FROM report_db.user_bak WHERE id = 10001;目标库出现同一行数据,说明 INSERT 链路通了。然后再测试 UPDATE:
UPDATE order_db.user SET age = 20 WHERE id = 10001;等一秒看目标库是否变成 20。如果变了,再测试 DELETE,删除后目标库对应行也要消失。这套链路验证完,说明四种基本操作都覆盖了。
另外你要注意,Adapter 在启动时不会自动加载源表存量数据。所以生产上线顺序应该是:先做一次全量数据迁移(DataX 或 mysqldump),再启动 Adapter,最后再放开业务写入。如果顺序反了,源库一边写一边导全量,目标库会出现主键冲突或数据缺失。
6. 同步到 ClickHouse / Elasticsearch 的差异与落地建议
6.1 不同下游对事件模型的需求不同
同样一条INSERT INTO order_db.user,同步到 MySQL 可以直接执行一条INSERT,但同步到 ClickHouse 时要考虑表引擎的更新语义,同步到 Elasticsearch 时要处理主键和字段映射。这就说明下游不同,Canal Adapter 的适配策略也不同。
Adapter 的conf目录下会有rdb、es7、clickhouse等子目录,每个目录对应一种目标端类型。你用哪个,就在启动时选用哪种映射文件。但要注意,Adapter 不会帮你创建表或索引,它默认只做行级事件映射,所以目标端表结构要提前设计好。
6.2 同步到 ClickHouse 的关键点
ClickHouse 本身是一个分析型数据库,默认的MergeTree引擎不直接支持单行 UPDATE 和 DELETE,只能通过异步合并去重。Canal 同步时,如果直接按 MySQL 的 INSERT 语义往里写,数据会不断追加,多次更新同一主键后会出现重复行。所以生产上我建议目标表使用 ReplacingMergeTree,并显式指定版本列:
CREATE TABLE report_db.user_bak ( id UInt64, name String, age UInt32, create_time DateTime, update_time DateTime ) ENGINE = ReplacingMergeTree(update_time) ORDER BY id;这样每次 Canal 投递一条带update_time的数据,ClickHouse 后台合并时就会按update_time取最新值,从而收敛重复。但要注意,ClickHouse 的 UPDATE 并不会实时生效,合并是滞后的,所以查询时如果需要精确最新值,要加上FINAL关键字,或者依赖足够短的合并周期。
类型映射上,MySQL 的varchar建议写成Nullable(String),datetime用DateTime。尤其要注意时区,源 MySQL 的 datetime 不带时区,ClickHouse 的 DateTime 也有时间歧义,建议同步连接串里统一传serverTimezone=Asia/Shanghai,两边保持一致。
6.3 同步到 Elasticsearch 的 index 映射
ES 的同步和阿里巴巴 Canal 生态搭配很成熟。Adapter 的 es 映射文件里,不是直接写 SQL 字段名,而是写 index 字段映射。比如:
dataSourceKey: defaultDS destination: example outerAdapterKey: es1 esMapping: index: user_index id: id sql: "SELECT id, name, age FROM order_db.user" objFields: name: keyword age: integer这里index对应 ES 的索引名称,id指定文档 ID,sql是 Adapter 从源库取数据时用到的查询语句。objFields定义字段类型映射。
ES 端有个常见坑:Adapter 不会自动创建索引映射,启动同步前需要你先手工创建 index 和 mapping。如果 mapping 是动态生成的,默认会把所有字符串字段识别成 text 加 keyword,字段语义可能出错,排序和聚合时就会出幺蛾子。建议提前设计映射:
{ "mappings": { "properties": { "name": { "type": "keyword" }, "age": { "type": "integer" } } } }6.4 什么时候建议走 MQ 自研消费而不是 Adapter
Adapter 好用,但一旦同步需求超过"单表简单镜像"的范围,它就开始显得笨拙。比如:你需要把 A 表和 B 表 join 后写入宽表;你需要对事件做窗口去重;你需要同时写多个下游并保证最终一致性。这些场景让 Adapter 承担规则引擎,会让配置文件变成一团乱麻。
这时正确的做法是把 Canal 的 serverMode 改成 Kafka,让 Canal 把 binlog 事件投递到 MQ,然后你在业务服务里自由消费。你可以基于 Canal Adapter 内置的某种格式重新解析,也可以直接用 Canal 官方的canal-client从 Kafka 消息里解析CanalEntry.Entry。这个路径其实更能发挥 Canal 的实时能力,是你从"会用"走向"用得深"的必经之路。
Kafka 模式配置也不复杂,在canal.properties里改成:
canal.serverMode=kafka kafka.bootstrap.servers=127.0.0.1:9092 kafka.acks=allCanal 会按 destination 和表名生成 topic,消费者组自己管理 offset。相比 TCP 模式,Kafka 模式天然支持多消费者、重放和持久化,缺点是需要额外维护 Kafka 集群,且位点回放和顺序保证要自己花精力设计。
7. 真实验证过的坑与排错思路
7.1 一次同步中断排查的完整链路
有一次同事找我,说目标库数据已经半小时没有更新,但源库一直有写入。我先看 Canal 的logs/example/example.log,果然有连续报错,内容是连接源库超时。随后查源 MySQL,发现网络确实没问题,但 MySQL 的错误日志里出现大量Aborted connection。
继续往下挖,发现 Canal 配置的server-id和现有从库重复了。MySQL 的行为是:后注册的复制连接会把先注册的踢掉,然后两边互相重连,不断循环。这个问题的隐蔽点在于,Canal 和从库都不觉得自己有问题,只有目标端数据延迟暴露出异常。
排查到这一步就很有趣了:不查 MySQL 错误日志,根本想不到是配置冲突。所以我想给你一个排查口诀:先看 Canal 日志有没有报错,再盯 MySQL 错误日志,再查连接是否被踢,最后才怀疑目标端写入问题。顺序反了,效率会非常低。
7.2 binlog 被清理导致位点失效的恢复
这是我最不愿意遇到、但最终一定会遇到的一类坑。场景是:Canal 停了三天,binlog 保留期只有两天,再启动时 Canal 报错说找不到指定的 binlog 文件或位置。
要恢复,先确认目标库能不能接受一段时间的缺失数据。如果能接受,最简单的方法就是重置位点。对于单机模式,删除logs/example/meta.dat(它保存了 Canal 上次的位点),再重启 Canal,Canal 会从当前 binlog 最新位置开始拉取。但这样做的代价是:停机期间漏掉的变更永远不会补上,目标库和源库会出现一个永久差异窗口。理想做法是重跑全量初始化,把两边对齐后再启 Canal。
如果开启了 Zookeeper 位点存储,我可以提前从 ZK 里把位点删掉,或者在 Canal Admin 里重置 instance。不过这些操作都是"高危动作",一定要先确认业务是否可以接受数据缺失,否则你就是给线上埋雷。
7.3 SSL 连接错误、字符集与时间时区问题
这些大多是"能跑但数据不对"的隐形问题,比直接报错更难发现。
常见的 SSL 错误出现在 Adapter 连接目标 MySQL 时,表现为类似javax.net.ssl.SSLHandshakeException。原因是 MySQL 8.0 默认开启 SSL,而驱动版本和服务器协商不一致。解决方案是 JDBC URL 显式关掉 SSL,并允许公钥检索:
jdbc:mysql://127.0.0.1:3306/report_db?useSSL=false&allowPublicKeyRetrieval=true&serverTimezone=Asia/ShanghaiallowPublicKeyRetrieval=true这个参数很多人会漏掉,尤其在caching_sha2_password认证下,不配置就容易一直认证失败。
字符集问题表现在中文乱码。根因是源库表是utf8mb4,但 Canal 的connectionCharset配成了 UTF-8 或目标库连接串没带characterEncoding=utf8。我现在的做法是:源库连接串统一useUnicode=true&characterEncoding=utf8,目标库连接串同样处理,Canal 的connectionCharset保持UTF-8,这样从源头到落库整条链路都是一套字符集,乱码概率大大降低。
时间字段相差 8 小时是另一类经典问题。如果 JDBC URL 里没配serverTimezone,驱动会用自己的本地时区去解释数据库时间。解决方式和上面一样,统一配serverTimezone=Asia/Shanghai。这个参数写不写,造成的问题时有时无,非常迷惑。
7.4 幂等性与重复消费:为什么下游必须支持覆盖写入
Canal 的下发机制不是严格 end-to-end 一次,它会因为网络重试、客户端重启、位点回退等原因重复投递同一条事件。如果你的目标端只是盲目执行INSERT,重复消费就会导致主键冲突。所以生产环境的目标表必须有唯一主键,并且映射配置里的primaryKey要和表结构一致。
对于 MySQL 目标端,我建议把同步方式定为覆盖写入,即主键冲突时更新、不存在时插入。Canal Adapter 默认对 RDB 目标就是 upsert 语义,所以它默认要求目标表有主键。如果一张源表没有主键,Canal 解析时拿不到定位字段,更新和删除事件基本就是废的。所以同步任何表之前,先检查源表主键,没有主键的表要特别评估。
7.5 大事务带来的延迟尖刺
Canal 本质上是跟着 binlog 走的,事务越大,Canal 在一个 batch 里收到的行数越多,解析和下发的时间就越长。我遇到过一次凌晨结算任务把一张流水表连续更新了 50 万行,Canal 延迟瞬间冲到 20 分钟,目标库的 batch 写入线程直接被打满。
这时候从 Canal 纬度很难做优化,因为事件已经形成了。正确做法是从源头减少大事务:业务侧把大批量更新拆成每批 1000 行的提交,或者把结算更新改成分批处理。如果无法改业务,就得在目标端加大syncBatchSize、调高连接池,并接受这个延迟尖刺。实时同步里面,源头事务大小直接决定同步的平滑程度,这句话建议写进同步设计评审的检查清单。
8. 运维经验:延迟、位点与性能
8.1 延迟从哪儿看:日志时间差与监控指标
Canal 同步做久了,最关心的不是数据对不对,而是慢了多少。我先说怎么在没有监控的情况下看延迟。
看 Canal instance 日志,里面会输出每批事件的写入时间,以及打印当前批次的 binlog 位点。结合SHOW MASTER STATUS查一下源库最新位点,两者文件一样、position 距离大不大,基本能估算落后量。更实用的方式是在目标库造一条带业务时间的数据,和源库比对插入时间差,这个误差最直观。
如果正式接入监控,建议用 Canal 自带或者社区开源的 Prometheus exporter。它暴露的指标里重点看几个:delayTime、pushLatency、executionTime。delayTime是当前处理事件与最新 binlog 之间的时间差,我最关注这个。把它配上告警,比如超过 5 分钟就报警,能第一时间发现同步卡住。
8.2 位点与高可用:别让单机 instance 裸奔
单机 Canal 部署适合开发和测试,生产环境至少要加上 Canal Admin 和 ZooKeeper。Canal 的位点存到 ZK 后,由 Admin 进行 instance 的 failover 管理。一旦 Canal 节点宕机,Admin 会把 instance 漂移到另一台健康节点,并从上一次记录的位点继续拉取。
这里有个细节,instance 的位点记录在 ZK 的特定节点上,如果 ZK 本身出现脑裂或者频繁变更,也会影响 Canal 稳定性。所以 ZK 集群最少三节点,别图省事用一个单节点。
另外,不管用不用 Admin,我每天都会巡检一次meta.dat或者 ZK 路径的位点变化。位点长期不变,要么是源库没写入,要么是同步链路已经断了。这两种状态需要人工区分,不然数据会慢慢产出一个隐性黑洞。
8.3 性能调优:保证顺序与提高吞吐之间的取舍
Canal 的吞吐瓶颈通常不在解析,而在目标端写入。因此调优思路要分两段说。
第一段是 Canal 拉取端。调大canal.instance.memory.buffer.size可以让 Canal 在内存里缓存更多事件,降低频繁拉取的开销,但要注意 JVM 内存。另外,开启canal.instance.parallel=true可以多线程解析 binlog,但也会带来事件顺序上的不确定性。如果你的下游对跨行顺序有强依赖,比如先插入父表再插入子表,我建议不要轻易开并行解析。
第二段是目标端。Canal Adapter 的concurrent参数决定同一个表的事件是否并行写入。concurrent=true能明显提升吞吐,但可能乱序。比如同一行连续两次 UPDATE,并行写就可能出现后更新的值先落库,导致最终数据还是旧值。我的经验是:普通表可以开启并发,但依赖状态变化的表(比如下单状态从 INIT 到 SUCCESS)必须保持串行,即concurrent: false。
目标库连接池大小同样要调。默认连接数往往不够扛住高峰,我把 Adapter 目标库连接池最大连接数调到了 50,batch 大小控制在 500 到 1000 行之间,实测效果比盲目加大单个 batch 要好很多。
8.4 最后再分享一点个人体会
从第一次用 Canal 到现在,我的心态变化挺大。刚上手时总想着把所有表、所有事件一股脑同步过去,觉得越全越好。后来发现很多同步需求其实用不上那么多事件,反而因为字段类型映射、主键缺失、大事务等问题给自己增加维护量。正确做法是先想清楚下游到底需要什么,再决定同步哪几张表、保留哪些字段,以及是否需要更新事件。
我始终给团队定一个规矩:凡是接 Canal 同步的表,必须在映射配置里写清楚主键、目标表、保留字段,并且在验收时做一次随机的源库变更对照。实时同步链路的稳定,不是靠某一次配置完成的,而是靠持续监控和每一个细节都较真。Canal 本身不复杂,复杂的是你愿不愿意把每一个字段的映射、每一次异常恢复都设计到位。
这个思路也可以继续扩展:当你把 Canal 的事件源同时接上 Kafka 和 Flink,可以做更复杂的流式计算;当你把目标端从 MySQL 换成 StarRocks,又能搭出实时数仓的分析链路。但不管下游怎么变,核心还是那几件事:binlog 开好、位点管好、表结构设计好、监控盯好。把这四样做扎实,Canal 会是你数据架构里最省心的一环。