Canal 源码深入
概述
Canal 伪装成 MySQL 从库,通过 binlog dump 协议拉取 binlog,内部由 EventParser、EventStore、EventSink 三个核心组件协作完成解析与投递。本文讲透三组件架构、binlog dump 协议与 HA 机制。
一、Canal 架构总览
MySQL Master
↓(binlog dump 协议)
Canal Server
├── EventParser(拉取+解析)
├── EventStore(环形缓冲)
└── EventSink(过滤+投递)
↓
MQ / 客户端| 组件 | 职责 |
|---|---|
| EventParser | 拉取 binlog、解析事件 |
| EventStore | 存储解析结果(环形队列) |
| EventSink | 过滤、序列化、投递 |
客户端两种模式:
1. Canal Client(拉模式)
2. MQ 模式(推模式:Kafka/RocketMQ)二、binlog dump 协议
2.1 原理
Canal 模拟 MySQL 从库:
1. 向 Master 发送 dump 请求
2. 指定 binlog 文件名 + 位置
3. Master 持续推送 binlog 事件| 协议要点 | 说明 |
|---|---|
| COM_REGISTER_SLAVE | 注册从库 |
| COM_BINLOG_DUMP | 开始 dump |
| 位点 | 文件 + offset |
| 心跳 | 保活与位点推进 |
2.2 binlog 事件类型
| 事件 | 说明 |
|---|---|
| QUERY_EVENT | 语句(DDL 等) |
| TABLE_MAP_EVENT | 表映射 |
| ROWS_EVENT | 行变更(row 格式) |
| XID_EVENT | 事务提交 |
| ROTATE_EVENT | 日志轮转 |
row 格式下:
WRITE_ROWS:插入
UPDATE_ROWS:更新
DELETE_ROWS:删除2.3 位点管理
Canal 记录消费位点:
binlog 文件名 + 位置
重启后从位点继续
(或从最新/指定时间)三、EventParser
3.1 职责
1. 连接 MySQL,dump binlog
2. 解析 binlog 事件
3. 生成 CanalEntry(结构化变更)
4. 写入 EventStore3.2 工作流程
连接管理 → 位点处理 → dump 拉取
→ binlog 解析 → 事务过滤 → 写入 Store| 子模块 | 说明 |
|---|---|
| 连接 | 主库连接池 |
| 位点 | 持久化/恢复 |
| 解析 | binlog → CanalEntry |
| 事务 | 按事务组织变更 |
3.3 CanalEntry
CanalEntry.RowChange:
事件类型(INSERT/UPDATE/DELETE/DDL)
RowData:beforeColumns / afterColumns
事务 id、时间戳JSON 化输出:
{"type":"UPDATE","data":[...],"old":[...],"ts":...}四、EventStore
4.1 职责
缓冲解析结果:
环形队列存储
解耦解析与消费4.2 实现
环形缓冲区:
put(生产者,Parser)
get/ack/rollback(消费者,Sink)| 操作 | 说明 |
|---|---|
| put | 写入事件 |
| get | 读取(不删除) |
| ack | 确认消费 |
| rollback | 回滚重放 |
ack 机制:
消费者确认 → 释放缓冲
失败 → rollback → 重新获取
保证不丢五、EventSink
5.1 职责
1. 从 Store 拉取
2. 过滤(库/表/过滤规则)
3. 序列化(protobuf/JSON)
4. 投递(MQ/客户端)| 功能 | 说明 |
|---|---|
| 过滤 | 黑白名单 |
| 转换 | 数据转换 |
| 投递 | 批量发送 |
5.2 投递模式
| 模式 | 说明 |
|---|---|
| 客户端拉取 | Client 消费 |
| MQ 推送 | Kafka/RocketMQ |
| 自定义 Sink | 扩展 |
六、HA 机制
6.1 问题
单点故障:
Canal 宕机 → 数据中断6.2 ZK 协调
Canal 集群 + ZooKeeper:
竞争同一 instance 的节点
只有一个工作,其他待命流程:
1. 节点注册 ZK(临时节点)
2. 竞争 instance 锁
3. 获得锁 → 成为工作节点
4. 宕机 → 临时节点消失 → 其他接管| 机制 | 说明 |
|---|---|
| 临时节点 | 会话断开即消失 |
| 竞争 | 抢锁 |
| 位点 | 接管后从位点续传 |
| 切换 | 自动 |
6.3 位点持久化
工作节点定期把位点写入 ZK/DB:
接管节点读位点继续
防重复/防丢| 注意 | 说明 |
|---|---|
| 幂等 | 下游去重 |
| 延迟 | 切换期间中断 |
| 一致性 | 位点与数据 |
七、部署与配置
7.1 配置要点
properties
# canal.properties / instance.properties
canal.instance.master.address=127.0.0.1:3306
canal.instance.dbUsername=canal
canal.instance.dbPassword=***
canal.instance.connectionCharset=UTF-8
canal.instance.filter.regex=db\\.table.*| 配置 | 说明 |
|---|---|
| master.address | 主库地址 |
| filter.regex | 过滤规则 |
| 位点模式 | ZK/文件 |
| 投递 | MQ 配置 |
7.2 监控
| 指标 | 说明 |
|---|---|
| 解析延迟 | Parser 到 Sink |
| 积压 | Store 占用 |
| 位点 | 与主库差距 |
| 错误 | 解析/投递失败 |
八、常见问题
| 问题 | 原因与处理 |
|---|---|
| 延迟高 | 检查消费/投递 |
| 数据丢失 | 检查 ack 与位点 |
| 重复 | 下游幂等 |
| HA 不切换 | 检查 ZK |
| binlog 缺失 | 保留期设置 |
常见问题速查
| 问题 | 要点 |
|---|---|
| 三组件职责 | Parser/Store/Sink |
| dump 原理 | 模拟从库 |
| 不丢怎么保证 | ack/rollback + 位点 |
| HA 怎么做 | ZK 竞争 |
| 接管续传 | 位点持久化 |