Canal这个组件,只要是搞过数据异构、缓存同步或者需要监听MySQL变更的同学,应该都不陌生。我最早接触它是因为当时负责的一个电商项目里,缓存和数据库的一致性总是靠定时任务去刷,延迟高不说,代码还写得特别丑。后来换了Canal做增量订阅,把binlog里的变更事件实时消费掉,整个架构清爽了很多。这篇就把我搭建Canal增量订阅和消费组件的完整过程,包括踩过的坑和优化心得,一次性讲清楚。
1. 为什么需要增量订阅,Canal到底解决了什么问题
1.1 业务系统里的数据同步之痛
先说说咱们遇到的实际场景。一个典型的互联网业务,数据存在MySQL里,但查询压力大的时候,不能所有请求都打到数据库上。于是引入了Redis缓存、Elasticsearch搜索引擎、还有异构的查询库。这时候问题就来了:MySQL里的数据变了,这些下游系统怎么知道?
最常见的做法是业务代码里双写,比如写完MySQL,顺手更新Redis、再调一次ES接口。听起来简单,实际维护起来非常痛苦。一旦某个下游系统暂时不可用,主流程就被拖住了;双写跨多个系统,事务没法保证,最终数据不一致是常态。还有些项目用定时任务全量扫描增量字段,延迟按分钟甚至小时算,数据量大以后效率也很低。
后来我换了思路:不再让业务代码主动通知下游,而是让MySQL把每一次数据变更记录到binlog里,由专门的组件去订阅这个日志流,再推给下游。这就是增量订阅的核心思想,数据库的变更就是我们最可靠的消息源,而Canal就是干这件事的。
1.2 Canal的原理:伪装成一个MySQL从库
Canal是阿里巴巴开源的组件,它的底层原理其实并不神秘。MySQL的主从复制大家应该都知道,从库连上主库,主库把binlog发给从库,从库重放这些日志保持数据一致。Canal做的事情,就是把自己伪装成一个MySQL从库,向主库发送dump协议请求,然后接收并解析binlog。
这里面有几个关键点。第一,binlog必须开启,而且格式必须设置成ROW模式,因为这个模式下binlog里记录的是每一行的字段变更前后值,而不是SQL语句本身,只有拿到行级别的变更数据,下游才能知道具体是哪条记录变了。第二,Canal解析完binlog之后,会把变更数据封装成CanalEntry对象,里面包含库名、表名、事件类型(增删改),以及变更前和变更后的字段值。这样下游消费的时候,根本不用关心MySQL底层的日志格式,拿到的是结构清晰的事件对象。
这里有个生活化的类比:MySQL的binlog就像超市的收银小票,记录了每一笔交易的流水。从库拿到小票照着重放一遍,保持自己货架和主库一致。Canal相当于专门请了一个人盯着这些小票,发现哪位顾客买了什么商品,就打电话通知仓库,让仓库那边同步调整库存。
1.3 增量订阅使用场景的边界
增量订阅这个能力,用对了地方价值很大,但也不是所有同步需求都能套用。我能想到的主要场景有三类:
第一类是缓存一致性,MySQL数据变更后,通过Canal把变更推给业务系统,由业务系统删除或者更新Redis缓存,这比双写优雅得多,而且延迟可以控制在毫秒级。第二类是异构数据同步,比如把MySQL的数据实时同步到Elasticsearch、HBase、ClickHouse等,做个实时搜索引擎或者分析系统,Canal负责把变更事件推过去。第三类是数据驱动的业务流程,比如订单状态变化后,订阅这个事件触发后续的库存扣减、物流通知等,数据变更本身就是业务事件的天然触发源。
不适合直接用Canal的场景也有,比如需要做复杂清洗和聚合的流式计算,这个应该交给Flink这类专业的流处理框架。Canal更合适的定位是数据变更的搬运工,把变更交出去,至于下游拿到之后怎么处理,Canal不管。
2. 搭建Canal Server的前置准备
2.1 MySQL侧必须满足的条件
Canal跑起来之前,MySQL这边先得具备基础条件,否则后面到处是坑。
第一,必须开启binlog。在MySQL配置文件(Linux下通常是/etc/my.cnf)里,要有这么几项:
[mysqld] log-bin=mysql-bin binlog_format=ROW binlog_row_image=FULL server-id=1注意几个细节。binlog_format=ROW是必须的,我之前说过只有ROW模式才有行级别的变更数据。binlog_row_image=FULL也是必须的,它告诉MySQL记录整行所有字段的前后值,而不是只记录变更涉及的字段。如果设置成MINIMAL,Canal解析出来的变更后数据里可能会缺字段,消费端拿到不完整的数据会非常头疼。
另外,server-id不能和Canal的slaveId冲突。Canal默认的slaveId是123456,如果你的MySQL服务没设置server-id,或者设置成了1之类的小数字,有可能会冲突,导致Canal连接被断开。这里建议MySQL的server-id设成一个不常用的值,比如168,Canal那边自定义一个不同的数字。
改完配置需要重启MySQL服务。然后检查一下binlog是否正常生成:
mysql> SHOW VARIABLES LIKE 'log_bin'; mysql> SHOW BINARY LOGS;能看到binlog文件列表,说明开启成功了。
第二,要给Canal创建一个专用的MySQL账号。Canal连接MySQL的时候需要这个账号来获取binlog和元数据信息。不要直接用root,权限给得太宽,而且容易出问题。创建语句是这样的:
CREATE USER 'canal'@'%' IDENTIFIED BY 'canal_pass'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%'; FLUSH PRIVILEGES;三个权限的含义要搞清楚。SELECT权限是Canal在获取表结构信息的时候需要的,因为它要解析binlog里的字段,得先知道表长什么样;REPLICATION SLAVE是让Canal能伪装成从库去请求binlog;REPLICATION CLIENT是让Canal能查询主库的server-id、binlog位置等状态信息。
这里有个版本兼容的坑,MySQL 8.0以上的版本,密码插件默认是caching_sha2_password,而Canal的老版本(1.1.4之前的)用的连接驱动不支持这种认证方式,会报认证失败。解决办法有两个:要么在创建账号的时候显式指定mysql_native_password插件:
CREATE USER 'canal'@'%' IDENTIFIED WITH mysql_native_password BY 'canal_pass';要么升级Canal到1.1.5及以上版本,或者自己替换Canal的MySQL驱动。我个人建议直接升级版本,因为老版本bug实在太多了。
2.2 Canal版本选型和系统环境
Canal的版本选型也是个学问。社区里目前常用的主流版本有1.1.4、1.1.5、1.1.6和1.1.7(现在的)等等。其中1.1.4是个分水岭,之前的版本架构比较旧,之后版本在管理端、适配器这些生态上完善了很多。如果追求稳定,我推荐1.1.5以上的版本,1.1.4在某些场景下会碰到一些已修复的bug。
Canal本身是用Java写的,所以运行环境需要JDK。这里有一个版本对应关系要注意:Canal 1.1.4及之前的版本,基于JDK 7编译,所以JDK 8能跑得挺好;后面的新版本推荐JDK 8或以上,直接装JDK 8是最稳妥的。我的经验是,用OpenJDK 8就足够了,没必要追新。
系统环境方面,Linux服务器先装好Canal再说。解压之后的目录结构是这样的:
canal.deployer/ ├── bin/ │ ├── startup.sh │ └── stop.sh ├── conf/ │ ├── canal.properties │ └── example/ │ └── instance.properties ├── lib/ └── logs/从目录结构就能看出一个概念:Canal Server可以管理多个instance,每个instance对应一组独立的订阅配置。conf/example就是一个默认的instance目录,名字叫example,它对应一份binlog监听配置。你完全可以建多个instance,比如conf/order_database专门监听订单库,conf/user_database专门监听用户库,互不干扰。后面我会讲怎么用。
2.3 确定消费方案:直接客户端还是对接消息队列
在启动Canal之前,得先想清楚下游消费方式,因为这会影响Canal Server的配置。目前主流有两种消费路径。
第一种是Canal Server直接暴露TCP端口,业务系统通过Canal提供的客户端SDK连接这个端口,拉取增量事件。特点是实时性最高,没有中间环节;但也有限制,连接数多了Canal Server压力大,而且如果消费方有多个,要自己管理消费位点。
第二种是Canal Server对接消息队列(RocketMQ或者Kafka),把解析好的binlog事件投递到MQ里,下游系统从MQ订阅消息。这种方式解耦性最强,多个消费者可以独立消费同一个topic而不互相影响,也天然利用MQ的堆积能力。缺点是链路多了一层,还要维护MQ集群。
绝大多数生产场景,我还是建议走MQ方案。原因很简单:数据同步这种需求,下游往往不止一个,比如一方面要更新Redis缓存,另一方面要同步ES索引,可能还要联动一些数据分析链路。如果都直连Canal Server,连接管理是个麻烦事。有了MQ,每个人各自消费各自的topic就行,Canal只负责往MQ里发一次消息。
如果下游只有一两个系统,而且对拓扑简单有要求,直连客户端方案也是一个可行选项。后面我会把两种方式的具体配置和代码都展示出来。
3. 完整落地:搭建一个增量消费链路
3.1 Canal Server的启动配置核心项
先编辑总配置conf/canal.properties。这个文件控制Canal Server的全局行为,我建议重点看这几个:
canal.port = 11111 canal.instance.global.manager.address = ${canal.admin.manager} canal.destinations = example canal.zk.hosts = canal.serverMode = tcp canal.mq.topic = canal_topic如果走kafka或rocketmq方案,canal.serverMode要改成对应的模式。比如:
canal.serverMode = kafka canal.mq.servers = 127.0.0.1:9092 canal.mq.topic = canal_topic canal.mq.partitionsNum = 4 canal.mq.partitionHash = example_db.example_table注意,canal.mq.topic可以设置成固定的topic名,也可以写成按规则动态生成。比如你希望每个表的消息进入独立topic,可以用${db}_${table}这种模板语法。但通常还是统一一个topic,然后在消息里带上库名表名字段,让消费者自己去路由,这样Kafka的分区策略更好控制。
接着编辑具体instance的配置conf/example/instance.properties,这是整个配置里最核心的部分:
canal.instance.master.address = 127.0.0.1:3306 canal.instance.dbUsername = canal canal.instance.dbPassword = canal_pass canal.instance.connectionCharset = UTF-8 canal.instance.tsdb.enable = true canal.instance.filter.regex = test_db\\..*canal.instance.master.address就是MySQL的主库地址,IP和端口都要写对。canal.instance.dbUsername和dbPassword对应我们前面创建的账号。canal.instance.tsdb.enable这个建议设成true,它表示开启表结构存储,Canal会缓存MySQL的表结构信息,因为解析binlog的时候需要知道字段类型才能正确反序列化。如果关掉它,遇到表结构变更(DDL)时Canal可能解析不了变更数据,反而会挂掉。
canal.instance.filter.regex是过滤正则,表示要监听哪些库的哪些表。格式是数据库名\\..*,注意Java正则里反斜杠要转义。比如监听test_db库所有表,就写test_db\\..*。如果想监听多个库,用竖线分隔:db1\\..*|db2\\..*。也有黑名单配置canal.instance.filter.black.regex,格式一样。正则写复杂了容易出问题,我的建议是能简单就简单,宁可多监听一些表也不要用复杂的正则,毕竟binlog解析本身不会消耗太多资源,下游消费时可以再过滤。
配置好之后,启动Canal Server:
sh bin/startup.sh然后看日志:
tail -f logs/example/example.log看到类似start successful的日志,说明Canal已经成功连接上MySQL并开始监听binlog了。这里有两个日志文件要注意:logs/canal/canal.log记录Server的启动过程,logs/example/example.log记录的是具体instance的解析状态。排查问题时先看canal.log,再看example.log,这个顺序可以省很多时间。
3.2 直连客户端消费:走TCP协议
如果MQ方案的条件还不成熟,比如公司暂时没有Kafka或者RocketMQ集群,那直连TCP就是一种能快速落地的方案。Canal官方客户端提供的示例代码很简洁,我这里用Java展示一个最小可用的消费端:
public class CanalClientDemo { public static void main(String[] args) { CanalConnector connector = CanalConnectors.newSingleConnector( "127.0.0.1", 11111, "example", "", "" ); connector.connect(); connector.subscribe("test_db\\..*"); connector.rollback(); while (true) { Message message = connector.getWithoutAck(1000, 1024); long batchId = message.getId(); if (batchId == -1 || message.getEntries().isEmpty()) { try { Thread.sleep(1000); } catch (InterruptedException e) {} continue; } for (CanalEntry.Entry entry : message.getEntries()) { if (entry.getEntryType() == CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); CanalEntry.EventType eventType = rowChange.getEventType(); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (eventType == CanalEntry.EventType.UPDATE) { // 变更前的数据 for (CanalEntry.Column column : rowData.getBeforeColumnsList()) { System.out.println("before: " + column.getName() + "=" + column.getValue()); } // 变更后的数据 for (CanalEntry.Column column : rowData.getAfterColumnsList()) { System.out.println("after: " + column.getName() + "=" + column.getValue()); } } else { for (CanalEntry.Column column : rowData.getAfterColumnsList()) { System.out.println("after: " + column.getName() + "=" + column.getValue()); } } } } } connector.ack(batchId); } } }这里面的核心方法就是getWithoutAck(batchId)和ack(batchId),两者配合构成了一套经典的消费语义。拿到一批消息之后,先不确认,业务处理完了才调用ack告诉Canal这一批已经消费成功了。如果业务处理抛异常,可以不调用ack,那么Canal下次还会再投递这一批,达到重试的效果。这个模型跟Kafka的消费位点提交逻辑是类似的。
实际项目里,如果消费者一端的逻辑比较复杂,比如要落库、要调外部接口,我建议把这个简单的while循环封装成线程池模式,一批消息作为一个任务丢给线程池处理。但要注意,Canal的客户端SDK在同一个连接上是有序的,同一批数据的binlog顺序不能乱。如果开多个线程并发处理,很可能会打乱顺序,带来数据错乱的问题。解决办法是如果必须并发,就要在消息里带上binlog的位置信息(entry的logfileName和logfileOffset),下游自己按这个维度去做顺序控制。这个字段在CanalEntry.Entry的header里面,其实可以取出来用。
3.3 对接Kafka:让多个下游各取所需
生产环境我强烈推荐走Kafka方案。Canal配置成Kafka模式很简单,在canal.properties里改这几项就行:
canal.serverMode = kafka canal.mq.servers = 192.168.1.10:9092 canal.mq.topic = canal_topic canal.mq.partitionsNum = 4 canal.mq.partitionHash = test_db.test_table这里的canal.mq.partitionHash指定哪个库表的binlog进入哪个分区。比如test_db.test_table这串配置,意思是同一张表的变更事件会投递到同一个Kafka分区,从而保证单表内的事件是有序的。这个对很多场景很重要,比如订单表的状态流转,如果消息分散到不同分区,消费者再并行处理,很容易出现后改的消息先处理了、先改的消息反而后处理,最终结果就错了。
Kafka模式的消息体,默认是JSON格式,示例大概长这样:
{ "id": 12345, "database": "test_db", "table": "test_table", "es": 1718185053000, "data": [ {"id": 1, "name": "zhangsan", "age": 20} ], "old": [ {"age": 18} ], "type": "UPDATE", "ts": 1718185053500 }字段含义很清晰:database和table标识是哪张表的数据变更,type是事件类型(INSERT/UPDATE/DELETE),data是变更后的整行数据,old是UPDATE之前被覆盖的旧字段值。这里提醒一下Canal发送MQ消息的时候,默认单个字段的类型信息不会带过去,JSON里全是字符串。如果你下游需要区分数值、时间类型,要么在消费端自己做类型转换,要么借助Canal提供的canal.mq.hashTag或者开启canal.mq.jsonData等配置调整结构。为了简单,我通常都是消费端按照自己掌握的schema去解析,反正data里的字段名是齐全的。
调度消费端的代码,用Kafka的原生consumer即可:
Properties props = new Properties(); props.put("bootstrap.servers", "192.168.1.10:9092"); props.put("group.id", "canal-es-consumer"); props.put("enable.auto.commit", "false"); props.put("auto.offset.reset", "latest"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props); consumer.subscribe(Arrays.asList("canal_topic"));消费端拿到消息之后自己解析JSON,根据database和table字段路由到不同的处理逻辑,比如解析到test_db下的product表发生UPDATE,就取data[0].id,调用ES的update API同步索引。这个过程中我踩过的一个坑是什么?自动提交offset。默认Kafka consumer定期自动提交offset,假设消息从Kafka拉出来、还没处理完进程崩了,恢复之后offset已经提交了,这部分消息就丢了。所以我建议把enable.auto.commit设成false,改成业务逻辑处理完成之后再手动commitSync()。虽然吞吐量会略降,但数据准确性远重要于这点性能。
3.4 消费确认与数据补偿机制
增量订阅方案要保证数据不丢不重,消费确认机制是关键。我梳理一下两种模式下我的做法。
Canal直连模式下,ack机制天然具备重试能力,没确认的消息会重新投递,因此业务处理中要注意幂等。Kafka模式下,Kafka本身不保证只消费一次,业务逻辑里必须做幂等。比如更新ES索引,ES的update操作天然是幂等的,重复执行只会覆盖成相同结果;但如果是累加操作,比如统计计数+1,那重复消费就会算两遍,这时候就要在业务表里记录一个消费位点或者处理过的binlog位置。
我和大家分享一个通用的幂等方案。消息体里不是有es这个时间戳吗?我通常在业务表里加一列last_sync_time,每次处理消息时,先比较消息时间戳和库里已记录的同步时间,如果消息时间戳更晚才处理,更早的直接丢弃。这个方法简单直接,尤其适合有主键且数据包含变更时间的场景。当然,字段级的幂等还有更多高级做法,但刚起步时这个方案够用了。
4. 实战中踩过的坑与消费性能调优
4.1 常见问题速查表
我把这几年监控Canal过程中遇到的典型问题整理成了一个表格,优先级从高到低排列:
| 现象 | 可能原因 | 排查方法与解决方案 |
|---|---|---|
Canal连不上MySQL,报Access denied for user | 授权账号缺失或密码插件不兼容 | 确认账号和权限;8.0以上指定mysql_native_password插件 |
| Canal启动成功但收不到任何增量事件 | 没有开启binlog,或binlog_format不是ROW | 检查MySQL配置,确认log_bin=ON,binlog_format=ROW |
| 消费端拿到的字段值有空值 | binlog_row_image不是FULL | 改成FULL,重启MySQL后重新同步 |
java.io.IOException: Received error packet | server-id冲突 | 调整Canal的slaveId,不要和MySQL的server-id重复 |
| 消费端store长时间收不到消息 | 有别的Canal client没ack | 检查是否有client一直没有调用ack,导致位点卡住 |
Canal频繁重启,日志里有table meta not found | 表结构缓存丢失,且tsdb未开启 | 设置canal.instance.tsdb.enable=true |
| DDL之后Canal解析异常 | 表结构变化,Canal缓存旧结构 | 开启tsdb或者清理Canal的meta信息,重新同步 |
这里单说一个我记忆很深的case。有一次线上MySQL做了小版本升级,结果Canal第二天开始频繁解析失败,打开日志发现binlog里出现了些Canal不认识的新类型标记。当时细节有点记不清了,结论就是Canal版本太老,对MySQL新版本的binlog兼容性有问题。升级Canal到新版本后,问题就消失了。从那以后,MySQL做小版本升级时,我会提前把Canal日志里的解析错误作为升级测试的一个检查项。
4.2 消费位点与数据一致性维护
Canal自己内部会保存位点信息,记录了当前已经消费到哪个binlog文件、哪个偏移量。直连模式下,位点在Canal Server端,由各个连接独立维护;Kafka模式下,位点被搬运到了Kafka的offset上,由消费者自己管理。理解这个差异很重要。
直连模式下,如果消费者处理逻辑异常、一直不ack,Canal Server上的位点就滚动不了。这时候如果消费者重启,Canal会根据上次的位点重新投递一遍,这其实是一种保障。但如果业务逻辑一半成功一半失败,而且没有幂等机制,就会出现重复数据。所以,生产上我的原则是:不管Kafka还是直连模式,消费端一律默认最终一致,通过幂等来兜底。
还有个小问题是misposition的情况。如果MySQL中binlog已经被清理(比如expire_logs_days设置太短),而Canal的位点落后到了已清理的日志位置,Canal启动时会报找不到binlog文件。这时的恢复方式有两个选择。一是重置位点让Canal从当前时刻开始订阅,但会丢一部分数据;二是利用canal.instance.master.gtid做GTID回溯,前提是MySQL开启了GTID模式。重置位点的操作是删除Canal的meta信息:
# 在conf/example/下删除meta.dat rm -f meta.dat # 然后重启Canal sh bin/stop.sh && sh bin/startup.sh清理meta之后注意,Canal会重新写一份位点在当前时刻,之前的数据需要靠全量初始化来补齐。实际生产中,配置expire_logs_days不要小于7天,并且对Canal做巡检,保证位点不落后太多,比出了问题再补救靠谱得多。
4.3 Canal对接多库多表的正则设计
监听多库多表是真实场景下的刚需,正则规则设计不好会导致漏数据。我之前做过的规则示例:
# 监听多个库的所有表 canal.instance.filter.regex = order_db\\..*|user_db\\..*|pay_db\\..* # 排除指定表 canal.instance.filter.black.regex = user_db\\.temp_table|order_db\\.order_log_backup这里需要提醒的是,正则匹配是Java的正则语法,点号要转义成功。如果在配置文件里写order_db..*而不是order_db\\..*,点号会匹配任意字符,能跑通但逻辑是错的。我一开始就踩过这种坑,过滤Regex写错,导致不该监听的表也监听了,消息量大了一倍。
另外,在消费端也建议根据database和table做第二层过滤,虽然Canal会投递所有配置匹配的表的变更,但下游可以只关注自己关心的表,这样即使Canal配置放宽,下游业务也不会被无关消息干扰。
4.4 并发、批量与性能调优经验
Canal的消费吞吐,通常瓶颈不在Canal Server,而在下游消费者。Kafka的消费并发能力取决于分区数。Canal往Kafka投递时,同一个表的变更默认会落到同一分区,这是为了保证顺序,但对应的消费并发就受限于分区数量。如果你的业务对单表内顺序要求不高,可以适当增加分区数,让Canal按主键hash来跨分区投递,这样单个表的消费并发也能提上来。
分区数的配置方法在canal.properties里加canal.mq.partitionsNum和canal.mq.partitionHash。例如:
canal.mq.partitionsNum = 8 canal.mq.partitionHash = test_db.test_table:id意思是test_db.test_table表按id字段哈希平均分到8个分区。这样同一个表的数据也可以被8个消费者并行处理,但需要明确的是,非严格有序了。如果你的场景只是同步ES索引,更新同一个id的文档最终只会覆盖旧值,所以乱序造成的影响不大;但如果业务依赖事件顺序,比如状态机流转,那就必须把所有事件走同一个分区。
批量消费也是提升性能的重要手段。Kafka的poll一次拉取一批消息,别只处理一条就提交。我通常的做法是积攒攒到一定条件再统一提交,比如处理完100条或者攒了500ms再commitSync(),这样做吞吐能提升不少。当然也要注意,commit不及时会造成重复消费,这又回到幂等的重要性上了。
还有一个细节是GC。Canal Server是Java进程,大流量的binlog解析会导致内存中产生大量短暂对象,如果JVM参数不调,频繁FGC会导致吞吐波动。官方默认的启动脚本里,有几个JVM参数可以调:
-Xms4g -Xmx4g -XX:+UseG1GC如果业务量不大,1g到2g堆内存足够;如果在高峰期流量大,我建议把堆内存调到4g以上,并开启G1。Canal的GC优化其实可以单独写一整篇,这里先说一个观察方向,排查时重点看logs/example/example.log里解析TPS和内存使用情况。
4.5 高可用部署思路
单机Canal部署起来很容易,但生产环境中不能容忍单点故障。Canal官方的高可用方案有两种。
一种是借助ZooKeeper做HA,多个Canal Server实例注册到同一个ZooKeeper,同一时间只有一个实例对外提供订阅服务,另一个处于standby状态。active实例挂了,ZooKeeper会自动把standby实例切换成active。这个方案我用过,配置上要指定canal.zk.hosts,每个instance的机器要保证配置一致。
另一种是把Canal的消费逻辑完全依托MQ的集群能力。Canal Server本身可以独立部署多套,每套监听不同的库表,消息投递到同一个Kafka集群。这个架构下,Canal Server本身挂了,影响的只是分配给它的那部分库表的同步,而下游消费端完全不受影响。对于大部分中小团队,我觉得这种按库表拆分的策略比依赖ZooKeeper更实用,拓扑清晰,容易扩展,最多就是Canal Server的配置需要自动化管理起来。
如果团队已经有成熟的注册中心或者容器编排平台,把Canal部署成多个副本,配合脚本做健康检查和自动拉起,也比传统的ZooKeeper方案省心。毕竟Canal本身不是无状态组件,但如果结合MQ,它的状态被转移到MQ的offset里了,Canal这边可以重启并重新订阅,损失的时间窗口是可以接受的。
5. 收尾建议与个人体会
这篇文章写到这里,核心的东西都覆盖到了。Canal不是那种很复杂的组件,它的核心价值在于把MySQL的binlog变成一条可靠的数据管道,让数据下游系统能实时感知上游变化。
我自己的实际体会是,接入Canal最重要的是想清楚消费语义和幂等策略,不要一上来就追求各种高级特性。一步一步来:先本机把binlog和Canal跑通,然后订阅到控制台打印消息,再把消息接入MQ,最后才对接业务逻辑。把这个链路走顺了,后面扩展新场景只是加表加配置的事。
如果再让我重新做一次数据同步方案,我还是会选择Canal,但这次我会把消费端的幂等和异常监控设计得更完善一些,因为在增量同步这个领域,稳定的生产链路比炫技的架构重要得多。希望这篇文章能帮你把Canal这项技术真正用起来,少踩一些我踩过的坑。