本文目录导读:

- 核心原理与架构
- 实战案例一:缓存双写一致性(Redis 与 MySQL 同步)
- 实战案例二:异构数据同步(MySQL 到 Elasticsearch)
- 实战案例三:数据变更审计(历史追踪)
- 高频踩坑与调优指南
- 总结建议
Canal 是阿里巴巴开源的一个基于 MySQL 二进制日志(Binlog)的增量数据订阅和消费组件,它的核心原理是伪装成 MySQL 的从库(Slave),向主库(Master)发送 dump 协议,主库则将 Binlog 推送给 Canal,Canal 解析 Binlog 后,将数据变更以特定格式(如 JSON)推送给下游。
下面从基础原理、核心应用场景、实战案例(含代码)、以及踩坑指南四个维度展开,帮助你深入理解 Canal。
核心原理与架构
- 主从复制协议:Canal 模拟 MySQL Slave 的交互协议,向 Master 请求 Binlog。
- Binlog 格式要求:必须设置为
ROW模式(binlog_format=ROW),这样才能获取到每行数据的变更前后镜像。 - 组件角色:
- Canal Server:接收 Binlog,解析为事件(Event)。
- Canal Client:通过 TCP 或 MQ(Kafka/RocketMQ)从 Server 拉取数据,并写入目标存储(如 ES、Redis、HBase)。
实战案例一:缓存双写一致性(Redis 与 MySQL 同步)
业务痛点:高并发下,先更新数据库再删除缓存的策略存在时间窗口,可能导致脏数据,使用 Canal 监听 Binlog 变更,异步更新或删除缓存,可有效避免缓存雪崩和脏读。
场景模拟:
用户表 user 更新了手机号,需要同步更新 Redis 中的用户缓存。
流程图:
应用更新 MySQL -> Canal 监听 Binlog -> Canal Client 消费变更 -> 应用删除/更新 Redis Key。
代码示例(Canal Client 消费核心逻辑):
// 引入依赖(以1.1.5为例,生产建议使用Maven管理)
// com.alibaba.otter:canal.client:1.1.5
public void consumeData() {
// 1. 连接 Canal Server
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), "example", "", "");
try {
connector.connect();
connector.subscribe("dbName.user"); // 订阅特定库表(正则)
connector.rollback(); // 回滚到上次消费位点
while (true) {
// 2. 批量拉取数据(每次拉取100条,超时500ms)
Message message = connector.getWithoutAck(100, 500L);
long batchId = message.getId();
if (batchId == -1 || message.getEntries().isEmpty()) {
Thread.sleep(1000);
continue;
}
// 3. 遍历 Binlog 条目
for (CanalEntry.Entry entry : message.getEntries()) {
// 只处理数据变更类型
if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) continue;
CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
CanalEntry.EventType eventType = rowChange.getEventType();
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
// 4. 获取变更后的数据(BeforeColumns / AfterColumns)
CanalEntry.Column idColumn = rowData.getAfterColumnsList()
.stream()
.filter(col -> col.getName().equals("id"))
.findFirst()
.orElse(null);
if (idColumn != null) {
// 5. 删除 Redis 缓存(或直接更新)
String key = "user:info:" + idColumn.getValue();
redisTemplate.delete(key);
System.out.println("缓存已失效: " + key);
}
}
}
connector.ack(batchId); // 确认消费成功
}
} catch (Exception e) {
e.printStackTrace();
} finally {
connector.disconnect();
}
}
注意点:
- 位点管理:消费失败不能调用
ack,否则数据会丢失。 - 幂等性:Canal 重启后重新投递,删除 Redis 的操作本身是幂等的,不会造成问题。
实战案例二:异构数据同步(MySQL 到 Elasticsearch)
业务痛点:商品搜索需要复杂的全文检索和多字段组合查询,而 MySQL 的 LIKE '%xxx%' 性能极差。
场景模拟:
商品表 product 价格或库存变化,需要实时同步到 ES 的 product 索引中。
架构升级:
Canal Server -> Canal Adapter(官方提供的适配器,支持 ES 索引自动映射)或者 Canal Client + 业务逻辑 -> ES RestHighLevelClient。
方案对比:
- 方案 A(低代码):使用 Canal Adapter,配置
es.yml和mapping文件,直接监听表变化,同步到 ES,优点:不用写代码,缺点:复杂逻辑(如字段拼接、多表关联)难以实现。 - 方案 B(推荐):使用 Canal Client 拉取数据,在业务代码中做 ETL(清洗、转换)后写入 ES。
代码示例(方案 B 核心逻辑):
// 在拉取到 product 表的变动后
if (eventType == CanalEntry.EventType.UPDATE || eventType == CanalEntry.EventType.INSERT) {
// 1. 构建 ES 文档对象
Map<String, Object> doc = new HashMap<>();
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
switch (column.getName()) {
case "product_id":
doc.put("id", column.getValue());
break;
case "product_name":
doc.put("name", column.getValue());
break;
case "price":
doc.put("price", Double.parseDouble(column.getValue()));
break;
// 忽略那些在 ES 中不需要的字段(如 logging_id)
}
}
// 2. 写入 ES
IndexRequest request = new IndexRequest("product_index")
.id(String.valueOf(doc.get("id")))
.source(doc, XContentType.JSON);
try {
esClient.index(request, RequestOptions.DEFAULT);
} catch (IOException e) {
// 失败时可以投递到 MQ 重试,或者记录日志由定时任务补偿
}
} else if (eventType == CanalEntry.EventType.DELETE) {
// 删除 ES 文档
DeleteRequest request = new DeleteRequest("product_index",
rowData.getBeforeColumnsList().get(0).getValue());
esClient.delete(request, RequestOptions.DEFAULT);
}
实战案例三:数据变更审计(历史追踪)
业务痛点:运营误操作修改了用户余额,需要快速定位“谁在什么时候改了什么”。
实现思路: 将 Binlog 变更记录原样(含变更前、变更后数据)写入 HBase 或 ClickHouse(列式存储,适合分析)。
代码处理核心点:
- 解析
RowData的BeforeColumns和AfterColumns。 - 将两个列表的字段名、旧值、新值拼接成 JSON。
- 记录操作时间(
entry.getHeader().getExecuteTime())和操作类型(eventType)。
// 构造审计日志实体
AuditLog log = new AuditLog();
log.setBizType(entry.getHeader().getSchemaName() + "." + entry.getHeader().getTableName());
log.setEventType(eventType.name());
log.setExecuteTime(new Date(entry.getHeader().getExecuteTime()));
// 记录变更字段差异
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
if (column.getUpdated()) { // 判断字段是否被更新
String fieldName = column.getName();
String newValue = column.getValue();
// 查找旧值
String oldValue = rowData.getBeforeColumnsList().stream()
.filter(c -> c.getName().equals(fieldName))
.map(CanalEntry.Column::getValue)
.findFirst().orElse(null);
log.addFieldChange(fieldName, oldValue, newValue);
}
}
auditLogService.save(log);
高频踩坑与调优指南
-
Binlog 模式必须是 ROW:
- 如果数据库是
STATEMENT或MIXED模式,Canal 拿不到具体的数据行变化,无法正常工作。 - 排查方法:执行
SHOW VARIABLES LIKE 'binlog_format';,必须为ROW。
- 如果数据库是
-
MySQL 主从复制用户权限:
- Canal 连接 MySQL 的用户必须有
REPLICATION SLAVE和REPLICATION CLIENT权限。 - 操作:
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
- Canal 连接 MySQL 的用户必须有
-
大事务导致 OOM(内存溢出):
- 如果一个事务更新了 10 万条数据,Canal 会将整个事务的 Binlog 加载到内存解析,如果内存配置不足,Server 会 OOM。
- 调优:减少
canal.instance.memory.buffer.size或者调整 JVM 堆内存-Xms和-Xmx。
-
乱序问题:
- 如果多个线程并发消费并写入目标库,可能会因为网络延迟导致数据乱序(如旧值覆盖新值)。
- 解决方案:对
主键 HASH进行分区/分片,保证同一条主键的数据由同一个线程执行(如 Kafka 按表主键分区)。
-
DDL(数据库结构变更)处理:
- 修改表结构(如新增字段)时,Canal 会发送
DDL事件,此时你需要判断目标端(如 HBase)的表结构是否也需要更新,否则写入会报错,通常需要建立 DDL 的回调机制进行表结构同步或告警。
- 修改表结构(如新增字段)时,Canal 会发送
总结建议
如果只是简单的 MySQL 与 Elasticsearch/HBase 同步,建议直接使用 Canal Adapter(官方插件)或者 Flink CDC(大数据场景,且有状态管理),如果你是 Java 技术栈,且需要复杂的业务逻辑(如缓存更新策略、特殊字段清洗),直接用 Canal Client 异步消费 会更灵活。
希望这些案例能够帮你将 Canal 落地到实际项目中,如果有具体的报错场景,欢迎带着日志进一步探讨。