Canal 是阿里巴巴开源的一款基于 MySQL 数据库增量日志解析的中间件,它通过模拟 MySQL 从库的交互协议,向主库发送 dump 请求,从而持续获取 Binlog 并解析成结构化的变更消息。Spring Boot 应用作为 Canal 的消费端,可以订阅这些消息,将数据库的插入、更新、删除事件实时同步到缓存、搜索引擎、数据仓库或另一个业务库中。这种方案对业务代码侵入小,只依赖 MySQL 的 Binlog 机制,因此在数据一致性要求较高的增量同步场景中被广泛采用。

接入 Canal 的核心并不只在引入依赖和启动客户端,更关键的是理解事件模型与位点确认机制。下面从 Binlog 复制原理开始,逐步展开 Spring Boot 侧的配置、事件解析和故障恢复策略。
一、Canal 监听 Binlog 的工作原理
MySQL 的主从复制是 Canal 实现数据捕获的基础。主库在执行写入操作时,会把变更记录按照事务顺序写入 Binlog 文件,从库通过一个专门的复制线程连接主库,向主库发送 dump 命令。主库接收到 dump 请求后,会持续把新的 Binlog 事件推送给从库,从库再将这些事件写入自己的中继日志并回放,最终完成数据同步。Canal 正是在这个链路中伪装成一个从库,只不过它拿到 Binlog 后并不回放 SQL,而是把二进制事件解析成结构化对象,再交给下游消费者处理。
Canal 本身由两部分组成:服务端和客户端。服务端负责与 MySQL 建立复制连接、接收 Binlog、解析并存储事件,客户端则通过 TCP 或 RocketMQ 等方式从服务端拉取事件。在 Spring Boot 中,我们通常使用官方提供的 canal.client 模块作为客户端,直接与 Canal 服务端交互。连接时会指定 destination,这个名称对应 Canal 服务端配置的一个实例,每个实例可以绑定一个 MySQL 库或多个库的过滤规则。
Binlog 有三种格式:STATEMENT、ROW 和 MIXED。Canal 默认要求 MySQL 使用 ROW 格式,因为 ROW 格式会记录每一行数据变更前后的完整字段值,方便做精确的数据同步。如果 MySQL 还停留在 STATEMENT 格式,只能拿到原始 SQL 文本,解析难度大且容易出现不一致。因此整合 Canal 之前必须确认 MySQL 的 binlog_format 已经设置为 ROW。同时还需要开启 binlog,并确保 Canal 使用的账号具备 REPLICATION SLAVE、REPLICATION CLIENT 以及对应库的 SELECT 权限。
二、Spring Boot 整合 Canal 的依赖与配置
在 Spring Boot 项目中引入 Canal 客户端依赖很简单,只需要在 pom.xml 中添加 canal.client 和 canal.protocol 两个依赖即可。canal.protocol 中包含通过 Protobuf 生成的事件模型类,canal.client 则提供连接器、消息拉取等基础能力。依赖版本建议与服务端保持一致,否则可能出现协议解析不兼容的问题。
<dependency>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal.client</artifactId>
<version>1.1.7</version>
</dependency>
<dependency>
<groupId>com.alibaba.otter</groupId>
<artifactId>canal.protocol</artifactId>
<version>1.1.7</version>
</dependency>
配置参数可以放在 application.yml 中,方便不同环境切换。比较关键的参数有 Canal 服务端地址、destination、批量拉取大小以及订阅的库表过滤规则。过滤规则使用正则表达式,例如订阅所有库所有表可以写成 .*\\..*,如果只想订阅某个库下的 user 表,可以写成 mydb\\.user。注意这里的反斜杠在 YAML 中需要正确保留,Java 字符串里则要写双反斜杠来转义。
canal: server: 127.0.0.1:11111 destination: example batch-size: 1000 subscribe: .*\\..*
在代码中读取这些配置并构建 CanalConnector,可以使用 CanalConnectors.newSingleConnector 工厂方法。它接收服务端地址、destination、用户名和密码,返回一个单机连接器。如果 Canal 服务端部署了集群,也可以切换到 newClusterConnector 使用 ZooKeeper 地址列表。连接建立后需要调用 connect 方法,再调用 subscribe 传入订阅表达式,最后进入一个循环不断拉取消息。
@Configuration
public class CanalConfig {
@Value("${canal.server}")
private String server;
@Value("${canal.destination}")
private String destination;
@Value("${canal.batch-size}")
private int batchSize;
@Value("${canal.subscribe}")
private String subscribe;
@Bean
public CanalConnector canalConnector() {
String[] parts = server.split(":");
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress(parts[0], Integer.parseInt(parts[1])),
destination,
"",
""
);
connector.connect();
connector.subscribe(subscribe);
return connector;
}
}
上面的示例中 server 字段格式为 host:port,在构建连接器时拆分成主机名和端口。批量大小 batchSize 用于控制每次 getWithoutAck 拉取的最大事件条数,如果 Binlog 产生速度很快,可以适当调大以减少网络交互;反之如果处理逻辑较重,可以调小避免单批消息积压过多导致处理超时。
三、实现 Canal 客户端与事件解析
完成连接配置后,最核心的工作是编写消息处理循环。canal.client 提供两种拉取方式:get 和 getWithoutAck。get 方法在返回消息后会自动确认消费位点,适合处理逻辑简单且不关心失败重试的场景;getWithoutAck 则不会自动确认,必须手动调用 ack 或 rollback,这样可以在处理失败时回滚位点,让没有成功的消息重新投递,适合对数据一致性有要求的同步任务。
@Component
public class CanalMessageListener {
private final CanalConnector connector;
private final BinlogEventProcessor processor;
public CanalMessageListener(CanalConnector connector, BinlogEventProcessor processor) {
this.connector = connector;
this.processor = processor;
}
public void start() {
while (true) {
Message message = connector.getWithoutAck(1000);
long batchId = message.getId();
if (batchId == -1) {
continue;
}
List<CanalEntry.Entry> entries = message.getEntries();
try {
if (!entries.isEmpty()) {
processor.process(entries);
}
connector.ack(batchId);
} catch (Exception e) {
connector.rollback(batchId);
// 记录异常日志,等待下一轮重试
}
}
}
}
消息中的核心对象是 Entry,它代表一个 Binlog 事件单元。Entry 可能属于事务开始、事务结束或者具体的行变更。解析时先判断 EntryType,如果是 TRANSACTIONBEGIN 或 TRANSACTIONEND,通常直接跳过;如果是 ROWDATA,则取出 storeValue 字段,使用 RowChange.parseFrom 进行 Protobuf 反序列化。RowChange 中包含事件类型 EventType、数据库名、表名以及变更前后的数据列集合。
public void process(List<CanalEntry.Entry> entries) {
for (CanalEntry.Entry entry : entries) {
if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN ||
entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) {
continue;
}
if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) {
continue;
}
CanalEntry.RowChange rowChange;
try {
rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
} catch (Exception e) {
throw new RuntimeException("解析 Binlog 事件失败", e);
}
CanalEntry.EventType eventType = rowChange.getEventType();
String tableName = entry.getHeader().getTableName();
for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
if (eventType == CanalEntry.EventType.INSERT) {
handleInsert(tableName, rowData.getAfterColumnsList());
} else if (eventType == CanalEntry.EventType.UPDATE) {
handleUpdate(tableName, rowData.getBeforeColumnsList(), rowData.getAfterColumnsList());
} else if (eventType == CanalEntry.EventType.DELETE) {
handleDelete(tableName, rowData.getBeforeColumnsList());
}
}
}
}
RowData 中的 beforeColumns 和 afterColumns 都是 Column 列表,每个 Column 包含列名、列值、是否为主键、是否更新过等元信息。对于 INSERT 事件,只有 afterColumns 有值;对于 DELETE 事件,只有 beforeColumns 有值;UPDATE 事件则两者都有。处理业务时建议先根据表名做路由分发,再把 Column 列表转换成 Map 或者 JSON 对象,方便后续写入目标存储。
private void handleInsert(String table, List<CanalEntry.Column> columns) {
JSONObject data = new JSONObject();
for (CanalEntry.Column column : columns) {
data.put(column.getName(), column.getValue());
}
syncService.syncInsert(table, data);
}
四、数据同步业务落地与位点确认
实际的数据同步往往不是简单地把 Binlog 原样写入另一个库,而是需要根据目标端的特点做适配。比如同步到 Elasticsearch 时,需要把 INSERT 和 UPDATE 统一处理成 upsert 操作,通过主键字段决定是新增文档还是更新文档;同步到 Redis 时,可以只缓存热点字段,并设置合理的过期时间;同步到另一个 MySQL 库时,可能需要将字段值做类型转换后再执行 SQL。业务分发层建议使用策略模式,将不同表的处理逻辑拆成独立的处理器,避免一个类里堆积大量 if else。
位点确认是保证数据不丢的核心。在 getWithoutAck 模式下,只有调用 ack 之后 Canal 服务端才会推进消费位置。如果处理过程中抛异常,需要调用 rollback 把当前批次回滚,这样下一次拉取还能拿到同样的消息。批量确认会带来一个小问题:同一个批次里只要有一条处理失败,整批都会回滚,已经处理成功的消息会被重复消费。因此处理逻辑必须设计成幂等,比如根据主键做 upsert、使用唯一约束去重,或者在目标端先查询再决定插入或更新。
事务边界也需要特别注意。Canal 会把同一个事务里的多条变更打包到同一个 Message 中,但它们可能分散在多个 Entry 里。如果目标端要求事务一致性,比如同步到另一个关系型数据库,最好把同一个 batchId 的消息在目标端也放到一个事务里提交。如果目标端不支持事务,比如 Redis 或 Elasticsearch,则需要考虑部分成功后的补偿机制,例如记录失败的记录到重试表,由后台任务异步修复。
@Service
public class BinlogEventProcessor {
private final Map<String, TableSyncHandler> handlerMap;
public BinlogEventProcessor(List<TableSyncHandler> handlers) {
this.handlerMap = handlers.stream()
.collect(Collectors.toMap(TableSyncHandler::supportTable, h -> h));
}
public void process(List<CanalEntry.Entry> entries) {
// 按表名分组后交给具体 handler 处理
for (CanalEntry.Entry entry : entries) {
// 省略 Entry 解析逻辑
String tableName = entry.getHeader().getTableName();
TableSyncHandler handler = handlerMap.get(tableName);
if (handler != null) {
handler.handle(entry);
}
}
}
}
五、断线重连与常见问题优化
Canal 客户端在运行过程中可能因为网络抖动、服务端重启、MySQL 主从切换等原因断开连接。默认的 CanalConnector 不会自动重连,因此需要在循环外层捕获连接异常,或者实现一个带重试的包装器。比较常见的做法是在 start 方法中判断 connector.checkValid 的返回结果,如果连接失效则重新 connect 并 subscribe。重连后不能简单从最新位点继续,否则会跳过断线期间的变更,应该依赖 Canal 服务端持久化的位点,让服务端从上次确认的位置继续投递。
另一个容易踩坑的地方是订阅表达式的匹配范围。默认情况下 .*\\..* 会订阅实例下所有库所有表,如果 Canal 服务端实例配置了过滤规则,客户端订阅范围不能超出服务端配置的范围。比如服务端只订阅了 order 库,客户端却订阅 .*\\..*,此时可能拉不到任何数据,或者拉取到空消息后进入死循环。排查时可以先在 Canal 服务端日志中确认当前实例的订阅范围,再调整客户端表达式。
性能优化方面,可以从批量大小、解析方式、目标端写入三个方向入手。批量大小建议根据单条消息平均大小和网络时延调优,默认 1000 条通常适合大多数场景。解析层面尽量复用 RowChange 对象,不要在循环里反复创建 Protobuf 解析器。目标端写入如果出现瓶颈,可以采用批量写入、异步队列或线程池分区处理,同时注意保持有序性,避免同一条主键的变更被乱序执行导致最终状态不一致。对于跨机房同步或目标端压力较大的场景,还可以在客户端与业务处理器之间引入本地缓冲队列,削峰填谷,避免瞬时 Binlog 洪峰打垮下游系统。
Spring BootCanalMySQL Binlog修改时间:2026-08-30 18:45:45