Canal — MySQL binlog 监听与数据同步
Canal 概述
Canal(发音为 /kəˈnæl/,意为"运河/管道")是阿里巴巴开源的一款基于 MySQL 数据库增量日志解析的中间件项目,提供增量数据订阅与消费的能力。项目最早诞生于阿里巴巴内部,旨在解决跨机房数据同步场景下的高可用与低延迟问题,后于 2015 年正式开源(GitHub — alibaba/canal)。
核心原理
Canal 通过模拟 MySQL Slave 节点的交互协议,向 MySQL Master 发送 dump 请求,接收并解析 Master 返回的 binlog(二进制日志)数据,将其转化为结构化的数据变更事件,供下游消费。
主要用途
| 场景 | 说明 |
|---|---|
| 缓存同步 | 监听数据库变更,实时更新 Redis / Local Cache,保证缓存与 DB 一致性 |
| ES 索引同步 | 将 MySQL 数据变更实时同步到 Elasticsearch,实现准实时搜索 |
| 异构数据同步 | MySQL → HBase / ClickHouse / MongoDB 等异构存储 |
| 数据迁移 | 基于 binlog 的增量数据迁移,支持不停机迁移 |
| 实时数仓 | 将 MySQL binlog 数据流入 Kafka / Flink,构建实时数仓 ETL 管道 |
| 审计与触发 | 记录所有数据变更操作,或触发业务侧自定义回调逻辑 |
优势与特性
- 非侵入式:无需修改业务代码,通过解析 binlog 获取数据变更
- 低延迟:毫秒级延迟,实时性高
- 高可用:支持基于 ZooKeeper 的 HA 模式
- 顺序保证:同一分区的 binlog 事件严格有序
- 多语言客户端:Java / Python / Go / Node.js / C# 等
MySQL binlog 原理
binlog 概述
binlog(Binary Log)是 MySQL 服务层维护的二进制日志文件,记录了对数据库执行更改的所有操作(不包括 SELECT / SHOW 等查询)。binlog 主要用于主从复制和数据恢复。
binlog 格式
MySQL 支持三种 binlog 格式:
| 格式 | 说明 | 优点 | 缺点 |
|---|---|---|---|
STATEMENT | 记录原始 SQL 语句 | 日志量小 | 依赖于上下文环境(如 NOW()、UUID() 等函数可能导致主从不一致) |
ROW | 记录每一行数据的变更前后映像 | 精确无误,主从绝对一致 | 日志量大(尤其对大字段批量更新) |
MIXED | 混合模式:默认使用 STATEMENT,遇到不确定函数时自动切换为 ROW | 折中方案 | 仍存在部分场景下的不一致风险 |
Canal 要求 MySQL 必须开启 ROW 格式,因为 Canal 需要解析具体的行级变更数据。
-- 查看当前 binlog 格式
SHOW VARIABLES LIKE 'binlog_format';
-- 查看 binlog 是否开启
SHOW VARIABLES LIKE 'log_bin';binlog 文件与位置
binlog 文件由一组二进制文件和一个索引文件组成:
mysql-bin.000001
mysql-bin.000002
mysql-bin.000003
mysql-bin.index每个 binlog 文件包含多个 Event,主要 Event 类型包括:
| Event 类型 | 含义 |
|---|---|
QUERY_EVENT | 执行的 SQL 语句(STATEMENT 格式) |
TABLE_MAP_EVENT | 表结构映射信息(ROW 格式必备) |
WRITE_ROWS_EVENT | 插入操作的行数据 |
UPDATE_ROWS_EVENT | 更新操作的行数据 |
DELETE_ROWS_EVENT | 删除操作的行数据 |
XID_EVENT | 事务提交 |
GTID_LOG_EVENT | GTID 相关事件 |
MySQL Dump 协议
Canal 通过 MySQL 的 COM_BINLOG_DUMP 协议获取 binlog 数据。基本原理如下:
- Canal 伪装为一个 MySQL Slave,向 Master 注册
- 发送
COM_BINLOG_DUMP命令,指定要读取的 binlog 文件名与偏移量 - Master 持续推送 binlog Event 流
- Canal 接收后解析为结构化数据
// 协议关键参数
// 发送 COM_BINLOG_DUMP
// binlog filename: mysql-bin.000001
// binlog position: 4
// server_id: 非 Master 且不与现有 Slave 冲突的 IDGTID(全局事务标识符)
GTID(Global Transaction Identifier)是 MySQL 5.6 引入的特性,为每个事务分配全局唯一的标识符,格式为 UUID:SEQUENCE_NUMBER。
-- 启用 GTID 模式
SET GLOBAL gtid_mode = ON;
SET GLOBAL enforce_gtid_consistency = ON;GTID 的优势:
- 简化主从切换后的定位逻辑(无需记录文件名和位置)
- 支持自动跳过已执行事务
- Canal 支持 GTID 模式下的 binlog 定位
Canal 使用 GTID 定位示例(instance.properties):
# 启用 GTID 模式
canal.instance.gtidon = true
# 指定 GTID 位置(可选)
canal.instance.gtid = d560e8b6-1234-11ec-abcde:1-100Canal 架构
整体架构
Canal 分为 Server 和 Instance 两层:
┌────────────────────────────────────────────┐
│ Canal Server │
│ ┌──────────────────────────────────────┐ │
│ │ Canal Instance (dest) │ │
│ │ ┌──────────┐ ┌─────────┐ ┌───────┐ │ │
│ │ │EventParser│→│EventSink│→│EventStore│ │ │
│ │ └──────────┘ └─────────┘ └───────┘ │ │
│ │ ↓ ↓ │ │
│ │ MetaManager Client连接 │ │
│ └──────────────────────────────────────┘ │
└────────────────────────────────────────────┘核心组件
EventParser(事件解析器)
负责从 MySQL Master 获取 binlog 原始数据并解析为 Canal 内部事件。
职责:
- 维护与 MySQL 的连接,发送 dump 请求
- 接收 binlog Event 流
- 解析
TABLE_MAP_EVENT获取表结构 - 将 binlog Event 转化为
CanalEntry.Entry对象 - 支持断点续传(记录已解析位点)
关键配置(canal.properties):
# 解析线程数
canal.instance.parser.threads = 4
# 批量解析大小
canal.instance.parser.batchSize = 1024
# 过滤 DDL 事件
canal.instance.filter.ddl = false
# 过滤 DML 事件类型(可按需开启)
canal.instance.filter.dml.insert = true
canal.instance.filter.dml.update = true
canal.instance.filter.dml.delete = trueEventSink(事件筛选/调度器)
将 Parser 产出的事件进行过滤、去重、路由,然后写入 EventStore。
职责:
- 支持自定义 Sink 过滤逻辑
- 支持多目的地路由(1 个 Parser → 多个 Store)
- 内置过滤:数据库名、表名、事件类型
EventStore(事件存储)
内存中的环形缓冲区(Ring Buffer),暂存已解析的事件,供 Client 消费。
设计特点:
- 基于 Disruptor(环形无锁队列)实现,高性能
- 可配置存储大小(默认 16 KB 个 entry)
- 支持批量获取和 ACK 机制
# EventStore 配置
canal.instance.memory.buffer.size = 16384
# 缓冲区存储单位(ITEMS 或 MEMSIZE)
canal.instance.memory.buffer.mode = ITEMS
# buffer 内存模式下的内存上限(单位:字节)
canal.instance.memory.buffer.memunit = 1024
# 缓存批次大小
canal.instance.memory.batchSize = 1024Disruptor 是 LMAX 开发的高性能无锁环形队列,Canal 的 EventStore 基于 Disruptor 实现 O(1) 级的读写效率,避免了 GC 压力和锁竞争。
MetaManager(元数据管理器)
记录 Canal Instance 当前消费的 binlog 位点(文件名 + 偏移量 或 GTID),用于:
- 重启后断点续传
- HA 切换时位点交接
支持的元数据存储模式:
| 模式 | 配置值 | 说明 |
|---|---|---|
| 内存 | memory | 默认模式,重启后丢失位点 |
| ZooKeeper | zookeeper | 生产推荐,HA 场景必须 |
| 本地文件 | localFile | 简单持久化,不支持 HA |
| 数据库 | database | 将位点存储在 MySQL 表中 |
| Redis | redis | 将位点存储在 Redis 中 |
# 元数据管理模式
canal.instance.global.mode = zookeeper
# ZooKeeper 集群地址
canal.zkServers = 127.0.0.1:2181,127.0.0.2:2181,127.0.0.3:2181HA 机制(高可用)
Canal 基于 ZooKeeper 实现 Running 模式——同一份 Instance 配置在多个 Canal Server 上部署,通过 ZK 分布式协调确保同一时刻只有一个 Server 提供消费服务。
┌──────────────┐ ┌──────────────┐
│ Canal Server │ │ Canal Server │
│ (Running) │ ←ZK→ │ (Standby) │
│ Instance A │ │ Instance A │
└──────┬───────┘ └──────────────┘
│
▼
┌──────────────────────┐
│ ZooKeeper │
│ /otter/canal/... │
│ └── cluster/node_1 │
│ └── cluster/node_2 │
│ └── destinations/ │
│ └── instance/… │
└──────────────────────┘HA 工作流程:
- 每个 Canal Server 启动时,在 ZK 指定路径创建临时节点
- 成功创建节点的 Server 成为 Running 节点,开始消费
- 其他 Server 监听该节点的删除事件
- Running 节点宕机后,临时节点自动删除,触发 Standby 节点抢占
- 新 Running 节点从 ZK 中读取上次保存的位点,继续消费
# ZK 相关配置
canal.zkServers = 127.0.0.1:2181
# ZK 中 Canal 的根路径
canal.zkRootPath = /otter/canalCanal 部署
环境要求
| 组件 | 版本要求 |
|---|---|
| MySQL | 5.6+(推荐 5.7+ 或 8.0) |
| Canal Server | 1.1.x+ |
| JDK | 1.8+ |
| ZooKeeper(可选) | 3.4+ |
MySQL 配置
开启 binlog
修改 MySQL 配置文件(my.cnf 或 my.ini):
[mysqld]
# 开启 binlog
log_bin = mysql-bin
# binlog 格式必须为 ROW
binlog_format = ROW
# 服务器 ID,不能与 Canal 或其他 Slave 冲突
server_id = 1
# 需要监听的数据库(留空表示全库)
binlog_do_db = your_database
# 或排除某些数据库
binlog_ignore_db = mysql
# binlog 保留天数
expire_logs_days = 7
# binlog 文件大小上限
max_binlog_size = 1G注意:MySQL 8.0 中
expire_logs_days已废弃,改用binlog_expire_logs_seconds:inibinlog_expire_logs_seconds = 604800
创建 Canal 账号
Canal 需要连接 MySQL 并模拟 Slave,需要授予相应权限:
CREATE USER 'canal'@'%' IDENTIFIED BY 'canal_password';
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
FLUSH PRIVILEGES;SELECT:读取表结构信息REPLICATION SLAVE:获取 binlog dump 权限REPLICATION CLIENT:查看 Master 状态(SHOW MASTER STATUS)
Canal Server 部署
下载与安装
# 下载 Canal Server 部署包
wget https://github.com/alibaba/canal/releases/download/canal-1.1.7/canal.deployer-1.1.7.tar.gz
# 解压
mkdir /usr/local/canal-server
tar -zxvf canal.deployer-1.1.7.tar.gz -C /usr/local/canal-server
# 目录结构
cd /usr/local/canal-server
# ├── bin/ # 启动/停止脚本
# ├── conf/ # 配置文件目录
# │ ├── canal.properties # Server 级别配置
# │ ├── logback.xml # 日志配置
# │ └── example/ # instance 目录(默认名 example)
# │ └── instance.properties
# ├── lib/ # 依赖库
# └── logs/ # 日志输出Server 级别配置(canal.properties)
#################################################
# Canal Server 基本配置
#################################################
# Canal Server 唯一标识
canal.id = 1
# Canal Server IP
canal.ip =
# Canal Server 端口(Client 连接端口)
canal.port = 11111
# 可用处理器数量(通常设为 auto)
canal.instances.cors = false
#################################################
# 用户认证
#################################################
canal.user = canal
canal.passwd = E358AAA65E10D8E96E6A2870FEA2C5F2
#################################################
# 实例加载方式
#################################################
# 默认加载 conf/ 目录下所有的 instance
canal.instance.detecting.enable = false
# 支持 spring 或 manager,默认 spring
canal.instance.manager.address =
#################################################
# 流量控制
#################################################
# QPS 限制,0 表示不限制
canal.instance.flowControl = 0
# 是否开启慢日志
canal.instance.flowControl.slow.threshold = 100
#################################################
# 数据传输安全(TCP 模式)
#################################################
canal.instance.transaction.size = 1024
canal.instance.fallbackIntervalInSeconds = 60
#################################################
# Canal Admin(可选)
#################################################
# canal.admin.manager = 127.0.0.1:8089
# canal.admin.port = 11110
# canal.admin.user = admin
# canal.admin.passwd = 4ACFE3202A5FF5CF467898FC58AAB1D615029441Instance 级别配置(instance.properties)
每个 Instance 对应一个要监听的 MySQL 数据源。默认实例名称为 example,可以创建多个 Instance(如 order_db、user_db)。
#################################################
# MySQL 连接配置
#################################################
# Canal 伪装为 Slave 的 server ID(不能与真实 Slave 冲突)
canal.instance.mysql.slaveId = 1234
# 数据源地址
canal.instance.master.address = 127.0.0.1:3306
# 连接认证
canal.instance.dbUsername = canal
canal.instance.dbPassword = canal_password
canal.instance.connectionCharset = UTF-8
# 连接参数
canal.instance.receiveBufferSize = 16384
canal.instance.sendBufferSize = 16384
# 连接超时时间(毫秒)
canal.instance.mysql.connection.timeout = 30
#################################################
# binlog 定位策略
#################################################
# 首次启动时的起始位点
# 取值: LATEST | 指定文件名:偏移量
canal.instance.binlogPosition = LATEST
# 指定 binlog 文件与位置
# canal.instance.master.journal.name = mysql-bin.000001
# canal.instance.master.position = 4
# canal.instance.master.timestamp =
# GTID 模式
canal.instance.gtidon = false
#################################################
# 过滤规则
#################################################
# 数据库名正则过滤(支持正则)
canal.instance.filter.regex = .*\\..*
# 黑名单
canal.instance.filter.black.regex = mysql\\.slave_.*
# 示例:只监听 mydb 下的以 user 开头的表
# canal.instance.filter.regex = mydb\\.user.*
# 表名大小写敏感
canal.instance.filter.table.case.sensitive = false
#################################################
# 事件解析配置
#################################################
# DDL 订阅
canal.instance.filter.ddl = true
# DML 订阅
canal.instance.filter.dml.insert = true
canal.instance.filter.dml.update = true
canal.instance.filter.dml.delete = true
#################################################
# EventStore 配置
#################################################
canal.instance.memory.buffer.size = 16384
canal.instance.memory.buffer.mode = ITEMS
canal.instance.memory.batchSize = 1024启动 Canal Server
# Linux / macOS
cd /usr/local/canal-server
./bin/startup.sh
# 查看日志
tail -f logs/canal/canal.log
tail -f logs/example/example.log
# 停止
./bin/stop.sh
# Windows
# bin\startup.bat验证启动成功:
# 查看 Canal Server 进程
ps -ef | grep canal
# 查看日志输出
# canal.log 中看到如下信息说明启动成功
# "Canal Launcher has been started ..."
# "## start the canal server ..."
# "## the canal server is running now ..."Canal Client
Java Client API
Canal 提供了官方 Java Client,支持 TCP 直连模式连接 Canal Server,获取数据变更事件。
添加依赖
<dependency>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal.client</artifactId>
<version>1.1.7</version>
</dependency>基本使用示例
package com.example.canal;
import com.alibaba.otter.canal.client.CanalConnector;
import com.alibaba.otter.canal.client.CanalConnectors;
import com.alibaba.otter.canal.protocol.CanalEntry;
import com.alibaba.otter.canal.protocol.Message;
import com.google.protobuf.ByteString;
import com.google.protobuf.InvalidProtocolBufferException;
import java.net.InetSocketAddress;
import java.util.List;
public class CanalClientDemo {
private static final int BATCH_SIZE = 1000;
public static void main(String[] args) {
// 1. 创建 Canal 连接器
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), // Canal Server 地址
"example", // Instance 名称
"", // 用户名(无需认证可留空)
"" // 密码
);
int batchId = 0;
try {
// 2. 建立连接
connector.connect();
// 3. 订阅过滤规则("" 表示订阅 Instance 配置的所有表)
connector.subscribe(".*\\..*");
// 可指定具体表:connector.subscribe("mydb\\.user_table");
// 4. 回滚到未 ACK 的位置(首次启动时有必要)
connector.rollback();
System.out.println("Canal Client started, waiting for events...");
while (true) {
// 5. 批量获取消息(阻塞式,超时时间可指定)
Message message = connector.getWithoutAck(BATCH_SIZE, 1000, TimeUnit.MILLISECONDS);
batchId = message.getId();
long size = message.getEntries().size();
if (batchId == -1 || size == 0) {
// 无数据变更,等待
Thread.sleep(100);
continue;
}
// 6. 处理数据变更事件
processEntries(message.getEntries());
// 7. 确认消费完成(ACK)
connector.ack(batchId);
}
} catch (Exception e) {
e.printStackTrace();
// 异常时回滚,下次重新消费
if (batchId > 0) {
connector.rollback(batchId);
}
} finally {
connector.disconnect();
}
}
private static void processEntries(List<CanalEntry.Entry> entries) {
for (CanalEntry.Entry entry : entries) {
// 跳过事务开始/结束事件
if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN
|| entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) {
continue;
}
// 解析行变更数据
CanalEntry.RowChange rowChange;
try {
rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
} catch (InvalidProtocolBufferException e) {
throw new RuntimeException("解析 RowChange 失败", e);
}
CanalEntry.EventType eventType = rowChange.getEventType();
String schemaName = entry.getHeader().getSchemaName();
String tableName = entry.getHeader().getTableName();
System.out.printf("======> Schema: %s, Table: %s, EventType: %s%n",
schemaName, tableName, eventType);
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
if (eventType == CanalEntry.EventType.DELETE) {
// 删除事件:before columns 为删除前的数据
printColumns(rowData.getBeforeColumnsList());
} else if (eventType == CanalEntry.EventType.INSERT) {
// 插入事件:after columns 为插入的数据
printColumns(rowData.getAfterColumnsList());
} else if (eventType == CanalEntry.EventType.UPDATE) {
// 更新事件:before 为更新前,after 为更新后
System.out.println("-----> 更新前数据");
printColumns(rowData.getBeforeColumnsList());
System.out.println("-----> 更新后数据");
printColumns(rowData.getAfterColumnsList());
}
}
}
}
private static void printColumns(List<CanalEntry.Column> columns) {
for (CanalEntry.Column column : columns) {
System.out.printf(" %s = %s (%s, updated=%b)%n",
column.getName(),
column.getValue(),
column.getMysqlType(),
column.getUpdated());
}
}
}批量消费与 ACK 机制
Canal Client 基于 pull 模式 消费数据,提供了完善的 ACK 机制保证至少一次消费语义。
Client Canal Server
│ │
│── getWithoutAck(batchSize) →│ ← 批量拉取事件,不移除
│←── Message(batchId=100) ───│
│ │
│── process(entries) │ ← 业务处理
│ │
│── ack(batchId=100) ───────→│ ← 确认成功,Server 移除
│ │
│ (如果处理失败) │
│── rollback(batchId=100) ──→│ ← 回滚,下次重新推送核心 API 说明:
| 方法 | 说明 |
|---|---|
getWithoutAck(batchSize, timeout, unit) | 批量获取消息,获取后不自动 ACK |
ack(batchId) | 确认指定 batch 已成功处理 |
rollback(batchId) | 回滚指定 batch(下次消费会重新推送) |
subscribe(filter) | 订阅过滤规则 |
connect() / disconnect() | 连接/断开 |
多 Instance 消费
一个 Canal Server 可以运行多个 Instance,Client 分别连接:
// 连接第一个 Instance
CanalConnector orderConnector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), "order_db", "", "");
orderConnector.connect();
orderConnector.subscribe("order\\..*");
// 连接第二个 Instance
CanalConnector userConnector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), "user_db", "", "");
userConnector.connect();
userConnector.subscribe("user\\..*");集群连接
生产环境推荐使用集群模式连接,配合 HA 实现故障自动切换:
// 基于 ZooKeeper 的集群连接
CanalConnector connector = CanalConnectors.newClusterConnector(
"127.0.0.1:2181,127.0.0.2:2181", // ZK 地址
"example", // Instance 名称
"", ""
);Canal Adapter
Canal Adapter 是 Canal 官方提供的开箱即用的数据同步适配器,无需编程即可将 MySQL 数据变更实时同步到目标存储。它支持多种数据源作为接收端,包括 RDB(关系型数据库)、Elasticsearch、HBase、Kafka、RocketMQ 等。
整体架构
┌──────────────────────────────────────────────┐
│ Canal Adapter │
│ │
│ ┌─────────┐ ┌─────────────────────────┐ │
│ │ Canal │ │ OutLoad Adapter │ │
│ │ Server │──→│ ┌─────┐ ┌────┐ ┌────┐ │ │
│ │ (Client) │ │ │RDB │ │ ES │ │MQ │ │ │
│ └─────────┘ │ └─────┘ └────┘ └────┘ │ │
│ └─────────────────────────┘ │
└──────────────────────────────────────────────┘安装 Adapter
# 下载 adapter
wget https://github.com/alibaba/canal/releases/download/canal-1.1.7/canal.adapter-1.1.7.tar.gz
# 解压
mkdir /usr/local/canal-adapter
tar -zxvf canal.adapter-1.1.7.tar.gz -C /usr/local/canal-adapter配置文件结构
conf/
├── application.yml # Adapter 总配置
├── bootstrap.yml # 启动配置
├── es/ # ES 同步配置目录
│ ├── biz_order.yml
│ └── user.yml
├── rdb/ # RDB 同步配置目录
│ └── my_other_db.yml
└── hbase/ # HBase 同步配置目录
└── my_table.yml总配置(application.yml)
server:
port: 8081 # Adapter HTTP 端口
spring:
jackson:
date-format: yyyy-MM-dd HH:mm:ss
time-zone: GMT+8
default-property-inclusion: non_null
# Canal Server 连接配置
canal.conf:
mode: tcp # 连接模式:tcp / kafka / rocketMQ
canalServerHost: 127.0.0.1:11111
# ZK 模式
# zookeeperHosts: 127.0.0.1:2181
# MQ 模式
# mqServers: 127.0.0.1:9092
# flatMessage: true
batchSize: 500
syncBatchSize: 1000
retries: 0
timeout: 30
accessKey:
secretKey:
# 消费方式
consumerProperties:
# canal tcp consumer
canal.tcp.username:
canal.tcp.password:
# 数据源配置(即要读取的 MySQL 数据源)
srcDataSources:
defaultDS:
url: jdbc:mysql://127.0.0.1:3306/mydb?useSSL=false&serverTimezone=UTC
username: canal
password: canal_password
# Adapter 适配器列表
canalAdapters:
- instance: example # 对应 Canal Server 的 Instance 名称
groups:
- groupId: g1
outerAdapters:
- name: logger # 日志输出(调试用)
- name: rdb # 关系型数据库同步
key: mysql1
properties:
jdbc.driverClassName: com.mysql.cj.jdbc.Driver
jdbc.url: jdbc:mysql://127.0.0.1:3306/target_db?useSSL=false
jdbc.username: root
jdbc.password: root
- name: es7 # ES 7.x 同步
key: es1
properties:
es.hosts: http://127.0.0.1:9200
es.batch.size: 200
- name: hbase
key: hbase1
properties:
hbase.zookeeper.quorum: 127.0.0.1:2181RDB 同步配置
映射配置文件
在 conf/rdb/ 目录下创建一个 .yml 文件(文件名随意):
# conf/rdb/user_sync.yml
dataSourceKey: defaultDS # 对应 srcDataSources 中的 key
destination: example # Canal Instance 名称
groupId: g1 # 分组 ID
outerAdapterKey: mysql1 # 对应 outerAdapters 中的 key
# 并发线程数
concurrentThreads: 3
# DDL 同步
ddlSync: true
# 表映射配置
tableMapping:
# 源表 → 目标表的映射
- targetTable: target_db.t_user # 目标表
targetPk: # 目标表主键
id: id
mapColumn: # 字段映射
id: id
user_name: name # 源表 user_name → 目标表 name
email: email
phone: phone
status: status
create_time: create_time
commitBatch: 100 # 批量提交大小
# 插入/更新/删除操作的 SQL 模板(Adpater 自动生成,可自定义)
# etlCondition: "where create_time >= '2024-01-01'"Elasticsearch 同步配置
ES 映射配置文件
在 conf/es/ 目录下创建映射文件:
# conf/es/user_index.yml
dataSourceKey: defaultDS
destination: example
groupId: g1
outerAdapterKey: es1
# ES 索引配置
esMapping:
_index: user_index # ES 索引名称
_type: _doc # ES 类型(7.x 后统一为 _doc)
_id: id # 文档 ID 字段
# 字段映射
upsert: true # 使用 upsert 方式
pk: id # 主键
sql: "select id, user_name as name, email, phone, status, create_time from t_user"
# ETL 条件(全量同步时可指定)
etlCondition: "where create_time >= '2024-01-01'"
# 字段类型覆盖
commitBatch: 200
# 跳过某些字段
# skipFields: []
# 字段映射(如果 ES 字段名与 SQL 列名不同)
# 默认同名映射ES 全量同步
# 触发全量同步(通过 Adapter 的 HTTP API)
POST http://127.0.0.1:8081/etl/es/user_index.yml
# 或全量导入
POST http://127.0.0.1:8081/etl/es/user_index.yml?para=id>0HBase 同步配置
在 conf/hbase/ 目录下创建映射文件:
# conf/hbase/user_hbase.yml
dataSourceKey: defaultDS
destination: example
groupId: g1
outerAdapterKey: hbase1
hbaseMapping:
namespace: default # HBase namespace
tableName: t_user # HBase 表名
rowKey: id # RowKey 字段
# 列族映射
columns:
# 列族:info
info:
- id
- name
- email
- phone
# 批量大小
commitBatch: 100TableMapping 详细说明
tableMapping 是 Adapter 中的核心配置,定义了源表到目标存储的映射关系:
| 配置项 | 说明 | 适用端 |
|---|---|---|
targetTable | 目标表名(RDB)/ 索引名(ES)/ 表名(HBase) | 通用 |
targetPk | 目标表主键映射 | RDB |
mapColumn | 字段映射关系 | RDB |
sql | 查询 SQL,用于全量同步和增量同步获取数据 | ES |
_id | ES 文档 ID 字段 | ES |
rowKey | HBase RowKey 字段 | HBase |
commitBatch | 批量提交大小 | 通用 |
etlCondition | 全量同步过滤条件 | 通用 |
ddlSync | 是否同步 DDL 操作 | RDB |
实战场景
场景一:Redis 缓存同步
通过 Canal Client 监听 MySQL 变更,实时更新 Redis 缓存,保证缓存与数据库的一致性。
架构
MySQL binlog → Canal Server → Canal Client → Redis代码实现
package com.example.canal.cache;
import com.alibaba.otter.canal.client.CanalConnector;
import com.alibaba.otter.canal.client.CanalConnectors;
import com.alibaba.otter.canal.protocol.CanalEntry;
import com.alibaba.otter.canal.protocol.Message;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisPool;
import java.net.InetSocketAddress;
import java.util.List;
import java.util.concurrent.TimeUnit;
public class RedisCacheSync {
private static final String INSTANCE = "example";
private static final String REDIS_KEY_PREFIX = "user:";
private static final int BATCH_SIZE = 500;
private final CanalConnector canalConnector;
private final JedisPool jedisPool;
public RedisCacheSync() {
this.canalConnector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), INSTANCE, "", "");
this.jedisPool = new JedisPool("127.0.0.1", 6379);
}
public void start() {
canalConnector.connect();
canalConnector.subscribe("mydb\\.t_user");
canalConnector.rollback();
try (Jedis jedis = jedisPool.getResource()) {
while (true) {
Message message = canalConnector.getWithoutAck(BATCH_SIZE, 1000, TimeUnit.MILLISECONDS);
long batchId = message.getId();
if (batchId != -1 && message.getEntries().size() > 0) {
for (CanalEntry.Entry entry : message.getEntries()) {
if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) {
continue;
}
CanalEntry.RowChange rowChange =
CanalEntry.RowChange.parseFrom(entry.getStoreValue());
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
switch (rowChange.getEventType()) {
case INSERT:
case UPDATE:
// 写入/更新 Redis
syncToRedis(jedis, rowData.getAfterColumnsList());
break;
case DELETE:
// 删除 Redis 缓存
deleteFromRedis(jedis, rowData.getBeforeColumnsList());
break;
default:
break;
}
}
}
canalConnector.ack(batchId);
} else {
TimeUnit.MILLISECONDS.sleep(200);
}
}
} catch (Exception e) {
e.printStackTrace();
} finally {
canalConnector.disconnect();
}
}
private void syncToRedis(Jedis jedis, List<CanalEntry.Column> columns) {
String id = null;
String name = null;
String email = null;
for (CanalEntry.Column col : columns) {
switch (col.getName()) {
case "id":
id = col.getValue();
break;
case "user_name":
name = col.getValue();
break;
case "email":
email = col.getValue();
break;
}
}
if (id != null) {
String key = REDIS_KEY_PREFIX + id;
jedis.hset(key, "name", name);
jedis.hset(key, "email", email);
// 设置过期时间
jedis.expire(key, 3600);
System.out.println("[Redis 同步] 写入缓存: " + key);
}
}
private void deleteFromRedis(Jedis jedis, List<CanalEntry.Column> columns) {
for (CanalEntry.Column col : columns) {
if ("id".equals(col.getName())) {
String key = REDIS_KEY_PREFIX + col.getValue();
jedis.del(key);
System.out.println("[Redis 同步] 删除缓存: " + key);
break;
}
}
}
}缓存一致性问题
| 策略 | 说明 | 风险 |
|---|---|---|
| 先更新 DB,后删除缓存 | 旁路缓存模式(Cache-Aside),binlog 同步后执行 DEL | 删除失败导致脏数据 |
| 先更新 DB,后更新缓存 | 直写模式 | 并发时可能出现写入顺序错乱 |
| 延迟双删 | 先删缓存 → 更新 DB → 延迟 N 毫秒再删一次 | 实现复杂 |
| binlog 增量驱动 | 利用 Canal 确保 DB 写入成功后才触发缓存更新 | 依赖 Canal 可用性 |
推荐方案:Cache-Aside + Canal 补偿——业务代码维持 Cache-Aside 模式,Canal 作为最终一致性补偿机制。
场景二:Elasticsearch 索引同步
架构
MySQL binlog → Canal Server → Canal Adapter (ES) → Elasticsearch方式一:使用 Canal Adapter(零代码)
- 在
conf/es/下创建索引映射配置 - 配置
sql定义查询逻辑 - 自动监听变更并同步
# conf/es/order_index.yml
dataSourceKey: defaultDS
destination: example
groupId: g1
outerAdapterKey: es7
esMapping:
_index: order_index
_id: order_id
upsert: true
pk: order_id
sql: "SELECT o.id as order_id, o.order_no, o.user_id, u.user_name as buyer_name,
o.total_amount, o.status, o.create_time, o.update_time
FROM t_order o
LEFT JOIN t_user u ON o.user_id = u.id"
commitBatch: 300
# 字段类型映射
# 需要手动在 ES 中创建 mapping方式二:使用 Canal Client 自定义写入
// ES 批量写入示例
private void syncToEs(List<CanalEntry.Entry> entries) {
try (RestHighLevelClient client = esClient()) {
BulkRequest bulkRequest = new BulkRequest();
for (CanalEntry.Entry entry : entries) {
CanalEntry.RowChange rowChange =
CanalEntry.RowChange.parseFrom(entry.getStoreValue());
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
switch (rowChange.getEventType()) {
case INSERT:
case UPDATE: {
Map<String, Object> doc = convertToMap(rowData.getAfterColumnsList());
IndexRequest request = new IndexRequest("order_index")
.id(String.valueOf(doc.get("order_id")))
.source(doc);
bulkRequest.add(request);
break;
}
case DELETE: {
Map<String, Object> before = convertToMap(rowData.getBeforeColumnsList());
DeleteRequest request = new DeleteRequest("order_index")
.id(String.valueOf(before.get("order_id")));
bulkRequest.add(request);
break;
}
}
}
}
if (bulkRequest.numberOfActions() > 0) {
client.bulk(bulkRequest, RequestOptions.DEFAULT);
}
} catch (Exception e) {
// 失败回滚
throw new RuntimeException("ES 同步失败", e);
}
}场景三:MQ 消息写入(Kafka / RocketMQ)
架构
MySQL binlog → Canal Server → Canal MQ Producer → Kafka / RocketMQ → 业务消费者Canal Server 配置 MQ
Canal Server 原生支持将 binlog 事件直接写入 MQ,无需额外部署 Adapter。
# canal.properties —— 开启 MQ 模式
# 选择 MQ 类型
canal.serverMode = kafka
# canal.serverMode = rocketMQ
# Kafka 配置
canal.mq.servers = 127.0.0.1:9092
canal.mq.retries = 0
canal.mq.batchSize = 16384
canal.mq.maxRequestSize = 1048576
canal.mq.lingerMs = 100
canal.mq.bufferMemory = 33554432
# Canal Topic 命名
# 默认 Topic: 实例名 (example)
canal.mq.topic = canal_topic
# 按表名动态 Topic: canal.mq.dynamicTopic = true
# regex 匹配的表对应不同的 Topic 分区
# RocketMQ 配置
# canal.mq.servers = 127.0.0.1:9876
# canal.mq.namespace =
# canal.mq.accessChannel = localInstance 级别 MQ 配置
# instance.properties
# 按表分发到不同 Topic(正则匹配)
canal.mq.dynamicTopic = mydb\\.t_user:user_topic, mydb\\.t_order:order_topic
# 分区数
canal.mq.partitionsNum = 3
# 分区哈希字段(保证同一行的变更进入同一分区)
canal.mq.partitionHash = id
# 消息格式(true: 扁平 JSON / false: protobuf)
canal.mq.flatMessage = true扁平消息格式
当 canal.mq.flatMessage = true 时,MQ 中的消息格式为 JSON:
{
"data": [
{
"id": "1001",
"user_name": "张三",
"email": "zhangsan@example.com",
"status": "1",
"create_time": "2024-01-15 10:30:00"
}
],
"database": "mydb",
"es": 1705293000000,
"id": 5,
"isDdl": false,
"mysqlType": {
"id": "bigint(20)",
"user_name": "varchar(50)",
"email": "varchar(100)",
"status": "tinyint(4)",
"create_time": "datetime"
},
"old": null,
"pkNames": ["id"],
"sql": "",
"sqlType": {
"id": -5,
"user_name": 12,
"email": 12,
"status": -6,
"create_time": 93
},
"table": "t_user",
"ts": 1705293001000,
"type": "UPDATE"
}| 字段 | 说明 |
|---|---|
data | 变更后的数据行(数组,支持批量) |
old | 变更前的数据(UPDATE 时才有) |
database | 数据库名 |
table | 表名 |
type | 事件类型:INSERT / UPDATE / DELETE |
isDdl | 是否 DDL 语句 |
es | 事件产生时间(MySQL 端) |
ts | Canal 处理时间 |
id | Canal 内部事件 ID |
消费者示例(Kafka)
package com.example.canal.consumer;
import com.alibaba.fastjson.JSON;
import com.alibaba.fastjson.JSONObject;
import org.apache.kafka.clients.consumer.*;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.time.Duration;
import java.util.Collections;
import java.util.Properties;
public class CanalKafkaConsumer {
public static void main(String[] args) {
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "127.0.0.1:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "canal-consumer-group");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(Collections.singletonList("canal_topic"));
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
JSONObject message = JSON.parseObject(record.value());
String database = message.getString("database");
String table = message.getString("table");
String type = message.getString("type");
System.out.printf("[Binlog] DB: %s, Table: %s, Type: %s%n",
database, table, type);
if ("INSERT".equals(type) || "UPDATE".equals(type)) {
// 处理数据
for (Object data : message.getJSONArray("data")) {
JSONObject row = (JSONObject) data;
// 业务逻辑处理
System.out.println(" 数据: " + row.toJSONString());
}
}
}
// 手动提交 offset
consumer.commitSync();
}
}
}
}消费者示例(RocketMQ)
package com.example.canal.consumer;
import com.alibaba.fastjson.JSONObject;
import com.alibaba.otter.canal.protocol.FlatMessage;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.common.message.MessageExt;
import java.nio.charset.StandardCharsets;
public class CanalRocketMQConsumer {
public static void main(String[] args) throws Exception {
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("canal-group");
consumer.setNamesrvAddr("127.0.0.1:9876");
consumer.subscribe("canal_topic", "*");
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> {
for (MessageExt msg : msgs) {
String body = new String(msg.getBody(), StandardCharsets.UTF_8);
JSONObject message = JSONObject.parseObject(body);
String table = message.getString("table");
String type = message.getString("type");
System.out.printf("[RocketMQ] %s.%s → %s%n",
message.getString("database"), table, type);
// 业务处理...
// 手动 ACK
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
});
consumer.start();
System.out.println("RocketMQ Consumer 启动成功");
}
}监控与运维
Metrics 指标
Canal Server 内置了 Metrics 指标输出,支持 JMX 和 Prometheus(需开启)。
开启 Prometheus 指标
# canal.properties
# 开启 Metrics
canal.metrics.prometheus.enable = true
# 指标暴露端口
canal.metrics.prometheus.httpServer.port = 9090访问 http://127.0.0.1:9090/metrics 查看指标数据。
核心指标说明
| 指标名 | 类型 | 说明 |
|---|---|---|
canal_instance_received_binlog_bytes | Counter | 已接收的 binlog 字节数 |
canal_instance_parsed_events_total | Counter | 已解析事件总数 |
canal_instance_consumed_events_total | Counter | 已消费事件总数 |
canal_instance_delay | Gauge | 当前延迟时间(毫秒) |
canal_instance_store_size | Gauge | EventStore 缓冲区当前大小 |
canal_instance_store_limit | Gauge | EventStore 缓冲区容量上限 |
canal_instance_latest_position | Gauge | 当前 binlog 位置 |
canal_instance_ack_total | Counter | ACK 次数 |
canal_instance_rollback_total | Counter | 回滚次数 |
Prometheus + Grafana 监控
# prometheus.yml 配置
scrape_configs:
- job_name: 'canal'
static_configs:
- targets: ['127.0.0.1:9090']Grafana 推荐面板:canal 延迟趋势图、QPS 曲线、EventStore 积压量、位点变化图。
日志分析
Canal 的日志目录结构:
logs/
├── canal/ # Server 级日志
│ ├── canal.log # 主日志
│ └── canal_stdout.log # 标准输出日志
└── example/ # Instance 级日志
├── example.log # 实例运行日志
└── example_stdout.log # 实例标准输出日志常见日志及排查
| 日志内容 | 含义 | 处理措施 |
|---|---|---|
something goes wrong with dump, reason: 1236 | binlog dump 异常 | 检查 MySQL 权限、binlog 文件是否存在 |
can't find start position | 找不到启动位点 | 检查 instance.properties 中的位点配置 |
connector disconnected | 与 MySQL 连接断开 | 检查网络、MySQL 连接超时设置 |
Store is full, size=X | EventStore 已满 | 增大 buffer.size 或加快消费速度 |
There is no slave registry | 未成功注册为 Slave | 检查 MySQL 的 server_id 是否冲突 |
TableId not found by tableId | 表结构缓存丢失 | 通常是 DDL 变更引起,重启 Instance |
java.io.IOException: Broken pipe | 写数据到 Client 时连接断开 | Client 端超时或异常断开 |
常用诊断命令
# 查看 MySQL 端 binlog 状态
mysql -h 127.0.0.1 -u canal -p -e "SHOW MASTER STATUS;"
# 查看活跃的 Slave 连接
mysql -h 127.0.0.1 -u canal -p -e "SHOW SLAVE HOSTS;"
# 查看 Canal 进程中的线程状态
jstack -l <canal_pid> | grep -E "parser|sink|store"
# 查看 EventStore 积压
curl -s http://127.0.0.1:9090/metrics | grep canal_instance_store_size
# 实时查看 Canal 日志中最新的事件
tail -f logs/example/example.log | grep "DML\|DDL"性能调优
Canal Server 调优
# 1. 增大 EventStore buffer
canal.instance.memory.buffer.size = 65536
# 2. 增大批处理大小
canal.instance.memory.batchSize = 2048
# 3. 增大 Parser 线程数
canal.instance.parser.threads = 8
# 4. 增大事务处理大小(毫秒)
canal.instance.transaction.size = 2048
# 5. 增大发送缓冲区
canal.instance.sendBufferSize = 65536
canal.instance.receiveBufferSize = 65536
# 6. 设置流量控制阈值
canal.instance.flowControl = 5000MySQL 端调优
[mysqld]
# 增大 binlog 缓存大小,减少磁盘 I/O
binlog_cache_size = 2M
# 增大 binlog 事务缓存
max_binlog_cache_size = 4M
# 减少 fsync 频率(适当容忍崩溃时丢失数据的风险)
sync_binlog = 100
# 增大 binlog 文件大小,减少文件切换
max_binlog_size = 1G
# 使用 ROW 格式的 image 为 full(确保完整行数据)
binlog_row_image = FULLClient 端调优
// 1. 合理设置 batch size(避免单次拉取数据过大)
int batchSize = 2000;
// 建议范围: 500 ~ 5000,根据单行数据量调整
// 2. 调整超时时间(避免频繁空轮询)
Message message = connector.getWithoutAck(batchSize, 2000L, TimeUnit.MILLISECONDS);
// 3. 并行消费(注意顺序性要求)
// 如果不需要严格顺序,可多线程并行处理
executorService.submit(() -> processEntries(message.getEntries()));
// 4. 优化 GC
// JVM 参数: -Xms4g -Xmx4g -XX:+UseG1GC性能瓶颈排查
| 瓶颈点 | 现象 | 排查方法 | 解决方案 |
|---|---|---|---|
| MySQL binlog 写入 | Canal 延迟持续上升 | SHOW MASTER STATUS 看 Position 更新缓慢 | 优化 MySQL 写入性能,增大 sync_binlog |
| 网络带宽 | 接收 binlog 字节数高 | iostat / 网络监控 | 压缩传输(Canal 暂不支持),升级网络 |
| EventStore 满 | Store is full 日志 | Metrics 监控 store_size | 增大 buffer 或加快消费 |
| Client 消费慢 | ACK 间隔长 | 查看 Client 端处理耗时 | 优化业务逻辑,增大消费并发度 |
| GC 停顿 | 延迟曲线出现尖峰 | GC 日志分析 | 调整 JVM 堆大小和 GC 策略 |
常见运维操作
重置消费位点
# 方式一:通过 Canal Admin API(需开启 Admin)
POST http://127.0.0.1:8089/api/v1/instance/reset
{
"destination": "example",
"journalName": "mysql-bin.000010",
"position": 4
}
# 方式二:直接修改 ZK 节点
zkCli.sh -server 127.0.0.1:2181
rmr /otter/canal/destinations/example/cursor
rmr /otter/canal/destinations/example/mark
# 方式三:删除元数据后重新启动 Instance(从 LATEST 开始消费)
# 删除实例配置后重启 Canal Server 会自动重建新增监听表
- 如果使用了
canal.instance.filter.regex = mydb\\..*,新表自动纳入 - 如果是精确匹配,需更新正则并重启 Instance
- Adapter 端需添加对应的映射配置文件
滚动升级 Canal
# 1. 备机标记为 Standby(ZK 模式下自动切换)
# 2. 升级备机 Canal Server
# 3. 重启备机,等待成为 Running 节点
# 4. 升级原主机
# 5. 重启原主机,自动进入 Standby 模式总结
Canal 作为阿里巴巴开源的高性能 binlog 增量订阅组件,已经在大规模生产环境中得到广泛验证。其核心价值在于:
- 解耦数据变更通知——将数据变更以事件流的方式暴露,与业务代码完全解耦
- 异构数据同步的标配——无论是缓存、搜索引擎还是大数据平台,Canal 都是 MySQL 增量同步的事实标准
- 高可用与低延迟——基于 ZK 的 HA 机制和 Disruptor 的高性能存储,保障了企业级的可靠性
典型技术栈组合:
MySQL + Canal + Canal Adapter → Elasticsearch → 搜索服务
MySQL + Canal + Canal Client → Redis → 缓存加速
MySQL + Canal + MQ → Kafka/Flink → 实时数仓
MySQL + Canal + Canal Adapter → HBase → 海量存储