乐于分享
好东西不私藏

Canal是如何偷看MySQL日记的?源码深度解析

Canal是如何偷看MySQL日记的?源码深度解析

🔥 MySQL肚里的秘密,Canal全知道 🔥

上期我们提到了Canal,说它是监听MySQL binlog的”终极方案”。

有读者问我:Canal到底是怎么工作的?为什么能实时同步数据?

今天这篇,从原理到源码,彻底讲清楚。


先说说什么是binlog

想象一下,MySQL有一个日记本,每天记录:

「2026-03-22 22:00:00,张三下单了,库存-1」

「2026-03-22 22:01:00,李四注册了,新用户+1」

「2026-03-22 22:02:00,王五修改了密码」

这个日记本,就是 binlog(二进制日志)

binlog里记录了所有对数据库的修改操作:INSERT、UPDATE、DELETE,甚至建表、修改表结构。

🐟 比喻:binlog就像公司的打卡记录本,每一笔操作都被忠实记录,时间、地点、人物、事件,一目了然。连老板都能通过它查出谁在几点几分干了啥。


Canal是干什么的?

Canal,中文意思是”运河”,是阿里开源的数据同步工具

它的核心能力是:偷看MySQL的日记本(binlog),然后把变化同步到其他地方(比如Redis、ES、其他MySQL)。

🐟 比喻:Canal就像派到MySQL公司的卧底间谍,每天蹲在日志室门口,只要有人来记录,他就立刻抄写一份,送到Redis那边去。老板(应用)不用亲自去翻记录,间谍已经整理好送过来了。

Canal的典型应用场景:

数据库缓存一致性:MySQL数据变了,Redis立刻同步更新

搜索索引同步:MySQL数据变了,ES(Elasticsearch)自动更新

异构数据同步:MySQL → Hive、数据仓库、ClickHouse

实时分析:数据库变更实时推送到Kafka,供实时计算使用


Canal的工作流程

Canal的工作可以分为4个步骤:

第一步:假装成MySQL的slave

MySQL有一种复制机制叫主从复制,master(主库)写数据,slave(从库)读数据。

原理是这样的:

主库(Master) ────┬──> 从库1(Slave)
     │              ├──> 从库2(Slave)
     │              └──> Canal(伪装成Slave)
     │
     ▼
  binlog

正常情况下,主库把binlog发给从库。Canal做的事情就是:伪装成一个slave,去连接master。

这样MySQL就会把binlog发送给Canal,就像发给真正的slave一样。

Canal会模拟MySQL的slave协议,发送COM_BINLOG_DUMP命令,告诉主库:”我是从库,把binlog发给我吧!”

第二步:解析binlog

binlog是二进制格式,人眼看不懂。Canal里面有解析器,把它变成我们能看懂的JSON:

原始binlog(二进制):
0001001011010101...

解析后:
{
  "table": "order",
  "type": "INSERT",
  "data": {"id": 1001, "user_id": 666, "amount": 99}
}

Canal支持三种binlog格式:

Statement:记录SQL语句(如 UPDATE order SET amount=99 WHERE id=1)

Row:记录具体的行数据变化(更精确,推荐)

Mixed:混合模式

💡 建议用Row模式,因为Statement模式有时会出问题(比如涉及时间函数、随机数时)。Row模式虽然日志量稍大,但绝对可靠。

第三步:投递到消息队列

Canal解析完不会直接写Redis,而是发到Kafka/RocketMQ

Canal → Kafka → 消费者 → Redis/ES/其他MySQL

为什么要经过MQ?解耦

好处:

削峰填谷:如果Redis扛不住,MQ可以缓冲

系统解耦:Canal不用关心谁来消费,数据变更不用修改业务代码

重复消费:MQ可以重试,消费者可以自行实现幂等

多消费者:一份数据可以同时给Redis、ES、数据仓库使用

🐟 比喻:间谍抄完日记,不是直接送到各部门,而是先放到传达室(MQ),各科室自己来取。这样间谍不会因为某个部门忙而卡住,也不用挨个部门跑腿。

第四步:消费并执行

消费者从MQ拿到数据,执行对应的操作:

// 伪代码:Redis消费者
while(true) {
  msg = kafka.consume();           // 从Kafka取消息
  event = JSON.parse(msg);        // 解析成事件
  
  switch(event.type) {
    case "INSERT":
      redis.set(event.key, event.value);  // 新增
      break;
    case "UPDATE":
      redis.set(event.key, event.value);  // 更新
      break;
    case "DELETE":
      redis.del(event.key);               // 删除
      break;
  }
}

源码核心逻辑

Canal的代码主要分为两部分:

📖 1. canal-server(服务端)

负责连接MySQL,解析binlog。核心代码简化:

// 1. 模拟MySQL slave协议,建立连接
public class CanalConnector {
    
    public void connect() {
        // 发送 COM_BINLOG_DUMP 命令
        // 告诉MySQL: 我要订阅binlog,从哪个位置开始
        socket.write(COM_BINLOG_DUMP + position);
    }
    
    // 2. 接收binlog事件(阻塞等待)
    public BinLogEvent fetch() {
        // 读取MySQL发来的binlog数据
        byte[] data = socket.read();  
        
        // 解析binlog事件
        return parseBinlog(data);      
    }
    
    // 3. 解析的具体实现
    private BinLogEvent parseBinlog(byte[] data) {
        // Row模式:解析每行的变化
        // 包括:表名、操作类型、变化前数据、变化后数据
        return new BinLogEvent(tableName, eventType, beforeData, afterData);
    }
}

📖 2. canal-client(客户端)

负责从服务端拉取数据,投递到MQ。核心代码简化:

// 从服务端拉取数据,投递到Kafka
public class CanalMQClient {
    
    private KafkaProducer kafka;
    private CanalConnector canal;
    
    public void start() {
        while(running) {
            // 1. 从Canal服务端拉取数据
            Message message = canal.get(1000);  // 一次最多取1000条
            
            // 2. 遍历每条binlog事件
            for (Entry entry : message.getEntries()) {
                // 3. 转换成JSON
                String json = toJson(entry);
                
                // 4. 发送到Kafka(按表名分区,保证同一表的消息有序)
                kafka.send(entry.getTableName(), json);
            }
            
            // 5. 记录消费位置(支持断点续传)
            savePosition(message.getId());
        }
    }
}

💡 简化理解:Canal源码就做两件事——拉数据(server)→ 发数据(client),中间经过MQ解耦。

整个过程支持断点续传:即使Canal重启,也能从上次的位置继续消费,不会丢数据。


Canal vs 其他方案对比

方案 实时性 复杂度 优点 缺点
定时任务 分钟级 简单,容易实现 有延迟,可能脏读
MQ异步 秒级 ⭐⭐ 解耦,削峰 可能丢消息,需补偿
Canal 毫秒级 ⭐⭐⭐ 零延迟、不侵入业务 需要开启binlog

生产环境使用注意

1. 开启binlog – MySQL配置文件中添加:log-bin=mysql-bin, binlog-format=ROW, server-id=1。注意:开启binlog会有约1%-10%的性能开销,但这是同步的代价,值得。

2. Canal单点问题 – 生产环境建议部署集群,用HA模式(高可用)。官方推荐:ZooKeeper + Canal Server集群,一个Canal Server宕机,ZooKeeper会自动切换到另一个。

3. 消息顺序性 – binlog有顺序,但Kafka不保证全局有序。解决方案:按表名分区(partitionKey),同一张表的消息保证有序。

4. 增量vs全量 – Canal只做增量同步,首次需要全量导入。常见做法:先全量导出一份数据到Redis/ES,然后用Canal增量同步。

5. 数据一致性 – 即使有Canal,也可能出现短暂的不一致(通常毫秒级)。对一致性要求极高的场景(如金融),需要做最终一致性校验。


总结

Canal的原理其实很简单:

1. 伪装成slave,连接MySQL
2. 接收并解析binlog(ROW格式)
3. 发送到Kafka/RocketMQ
4. 消费者订阅,执行具体业务(Redis/ES等)

Canal适合的场景:

• 需要实时同步(毫秒级)

• 不想在业务代码里侵入同步逻辑

• 需要一对多(一份数据同步到多个存储)

🐟 一句话总结:Canal就是MySQL的舔狗——MySQL一发朋友圈(binlog),Canal立刻点赞并转发到各个群(MQ),生怕错过任何一条动态。

下期想看什么?评论区告诉我 👇


觉得有用就点个在看 👇

本站文章均为手工撰写未经允许谢绝转载:夜雨聆风 » Canal是如何偷看MySQL日记的?源码深度解析

猜你喜欢

  • 暂无文章