Article

日志采集 Canal

更新于:2026-07-13

第1章:Canal 概述与核心原理

1.1 什么是 Canal

概念名称说明注意事项
Canal阿里巴巴开源的基于 MySQL 数据库增量日志(Binlog)解析的组件,用于提供增量数据订阅和消费。其名称”Canal”意为”水道”,寓意数据的流动通道。Canal 本身不存储数据,仅解析并转发 Binlog 事件,需配合消费者(Client 或 Adapter)进行数据落地或处理。
开源项目项目托管于 GitHub(https://github.com/alibaba/canal),采用 Java 编写,支持多种数据同步场景。社区活跃,版本迭代稳定,但部分高级功能(如 Canal Admin)需自行部署与维护。
增量同步专注于捕获数据库的变更(INSERT、UPDATE、DELETE),而非全量数据导出。依赖 MySQL 的 Binlog 功能,必须确保数据库开启 Binlog 并配置正确格式(推荐 ROW 模式)。

1.2 Canal 的应用场景

应用场景说明注意事项
数据异构将 MySQL 中的数据实时同步到其他存储系统,如 Elasticsearch(用于搜索)、HBase、Redis(缓存)等。目标系统需具备接收变更的能力,建议使用 Canal Adapter 简化配置。
缓存更新当数据库发生变更时,通过 Canal 通知缓存系统(如 Redis)更新或删除对应缓存,避免脏读。需保证缓存更新的实时性与一致性,可结合消息队列削峰。
数据库审计记录所有对数据库的修改操作,用于安全审计、合规检查或操作追溯。建议记录操作类型、时间、原始值与新值,并持久化存储。
微服务间数据同步在微服务架构中,服务间通过数据库解耦,Canal 可实现跨服务的数据事件通知。避免循环依赖,建议通过 Topic 或 Tag 进行消息路由。
数仓实时入仓将业务库的变更数据实时接入数据仓库或数据湖,支持实时分析。需考虑数据格式转换、字段映射、延迟容忍度等问题。

1.3 MySQL Binlog 基础回顾

概念名称说明注意事项
Binlog(Binary Log)MySQL 的二进制日志,记录所有对数据库的写操作(DDL、DML),是主从复制的基础。必须在 MySQL 配置文件中启用 log-bin 参数,否则 Canal 无法工作。
Binlog 格式:STATEMENT记录 SQL 语句原文。不推荐用于 Canal,因无法精确还原行级变更,且可能引发主从不一致。
Binlog 格式:ROW记录每一行数据的变更前(before)和变更后(after)的完整内容。Canal 必须使用 ROW 模式,以准确解析数据变更。
Binlog 格式:MIXED混合模式,由 MySQL 自动选择使用 STATEMENT 或 ROW。不推荐,可能导致 Canal 解析不稳定。建议显式设置为 ROW。
Binlog 文件与索引Binlog 以文件序列形式存储(如 mysql-bin.000001),并通过索引文件(.index)管理。Canal 会记录消费位点(position),重启后可从中断处继续读取。
server-id每个 MySQL 实例必须拥有唯一 server-id,用于主从复制标识。单机测试时也需配置,否则无法开启 Binlog。

1.4 Canal 的工作原理(基于 Binlog 的解析与订阅)

核心机制说明注意事项
模拟 MySQL SlaveCanal Server 启动后,伪装成 MySQL 的从库(Slave),向主库发起 dump 协议请求,获取 Binlog 流。需在 MySQL 创建具有 REPLICATION SLAVEREPLICATION CLIENT 权限的用户。
Binlog Dump 协议MySQL 主库通过 dump 线程将 Binlog 推送给 Canal,Canal 以流式方式接收。网络中断或 Canal 停止时,MySQL 会保留 Binlog 直至被消费,避免数据丢失。
日志解析(Parse)Canal 接收原始 Binlog 字节流,使用内置解析器将其转换为结构化的 Entry、RowData 等对象。解析过程不依赖 MySQL 实例,性能开销主要在 Canal Server 端。
数据存储(内存 Queue)解析后的数据暂存于内存队列(如 MemoryEventStore),等待 Client 拉取。队列大小可配置,过小可能导致阻塞,过大可能引发 OOM。
客户端订阅(Client Pull)Canal Client 通过 TCP 或 gRPC 协议连接 Server,发送订阅请求,并持续拉取数据。支持单 Client、多 Client 广播、集群模式(通过 ZooKeeper 协调)。
位点管理(Position)Canal 记录当前已解析和已消费的 Binlog 位置(包括 filename、position、timestamp、serverId)。支持自动提交或手动 ack,确保故障恢复后不丢数据也不重复消费。

1.5 Canal 的核心组件架构

组件名称说明注意事项
Canal Server核心服务进程,负责连接 MySQL、解析 Binlog、存储事件、响应 Client 请求。支持多 instance(实例)管理。每个 instance 对应一个 MySQL 数据源,配置独立,资源隔离。
Canal Client数据消费者,通过 API 连接 Server,订阅数据并处理。可部署在业务系统中。建议使用官方 Client-Example 或集成到 Spring 项目中。
ZooKeeper分布式协调服务,用于 Server 高可用(HA)、Client 负载均衡、位点存储等。生产环境推荐使用 ZooKeeper 集群,避免单点故障。
Canal AdminWeb 管理后台,用于集中管理多个 Canal Server 的配置、实例、状态等。非必需,但可大幅提升运维效率,需单独部署。
Canal Adapter数据同步适配器,支持将 Canal 数据写入外部系统(如 ES、HBase、Kafka、RDBMS)。基于 Spring Boot 构建,支持 REST API 和 YAML 配置。
Instance逻辑实例,代表一个数据源(MySQL 节点)的订阅任务,包含数据源配置、过滤规则、解析位点等。支持动态加载、卸载,可通过 Canal Admin 管理。
MetaManager管理 instance 的元信息,如当前位点(position)、运行状态等。可存储在内存、ZooKeeper 或本地文件中,推荐 ZooKeeper 以支持 HA。

第2章:Canal 环境准备与部署

2.1 MySQL 环境配置(开启 Binlog、用户权限)

配置项 / 命令语法 / 示例用途注意事项
log-binlog-bin = mysql-bin启用二进制日志,文件前缀为 mysql-bin必须开启,Canal 依赖 Binlog 获取数据变更
binlog-formatbinlog-format = ROW设置 Binlog 格式为 ROW 模式必须为 ROW,STATEMENT 或 MIXED 不支持精确行级解析
server-idserver-id = 1设置唯一服务器 ID,用于主从复制标识每个 MySQL 实例必须唯一,否则无法开启 Binlog
expire-logs-daysexpire-logs-days = 7Binlog 文件保留天数根据 Canal 消费速度设置,避免日志过早清理导致位点失效
binlog-row-imagebinlog-row-image = FULL记录行变更的完整前像和后像推荐设置为 FULL,确保 UPDATE 操作能获取 before 和 after 数据
创建 Canal 用户CREATE USER 'canal'@'%' IDENTIFIED BY 'canal';
GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
FLUSH PRIVILEGES;
授予 Canal 连接所需的最小权限REPLICATION SLAVE(用于 dump binlog)和 REPLICATION CLIENT(查看日志)是必需权限
验证 Binlog 是否开启SHOW VARIABLES LIKE 'log_bin';
SHOW VARIABLES LIKE 'binlog_format';
确认 Binlog 已启用且格式正确返回值应为 ONROW

2.2 ZooKeeper 搭建与配置(集群模式)

配置项 / 命令语法 / 示例用途注意事项
tickTimetickTime=2000ZooKeeper 心跳间隔(毫秒)通常设为 2000,不宜过大或过小
dataDirdataDir=/var/lib/zookeeper数据存储目录,包含 myid 和快照目录需存在且可写,每个节点路径不同
clientPortclientPort=2181客户端连接端口Canal Server 和 Client 将连接此端口
initLimitinitLimit=10Follower 初始连接 Leader 的超时时间(以 tickTime 为单位)一般为 1015,表示 2030 秒
syncLimitsyncLimit=5Follower 与 Leader 同步数据的超时时间一般为 5~10
集群节点配置server.1=zoo1:2888:3888
server.2=zoo2:2888:3888
server.3=zoo3:2888:3888
定义集群中各节点的通信地址2888 用于 Follower 连接 Leader,3888 用于 Leader 选举
创建 myid 文件在 dataDir 下创建文件 myid,内容为节点编号(1, 2, 3)标识当前节点在集群中的 ID每个节点的 myid 必须唯一且在 server 列表中定义
启动 ZooKeeperbin/zkServer.sh start启动单个节点建议使用 start-foreground 查看日志
验证集群状态echo stat | nc zoo1 2181查看节点角色(Leader/Follower)和连接数确保至少一个 Leader 和多个 Follower 正常运行

2.3 Canal Server 安装与配置

配置项语法 / 示例用途注意事项
下载地址https://github.com/alibaba/canal/releases获取最新版本的 canal.deployer推荐使用稳定版本,如 v1.1.7
解压安装tar -zxvf canal.deployer-$version.tar.gz -C /opt/canal解压到指定目录建议统一部署路径
canal.portcanal.port = 11111Canal Server TCP 服务端口多实例部署时需区分端口
canal.zkServerscanal.zkServers = 192.168.1.10:2181,192.168.1.11:2181指定 ZooKeeper 集群地址若使用 HA 模式,此项必填
canal.instance.mysql.hostcanal.instance.mysql.host = 192.168.1.100MySQL 主库地址支持域名或 IP
canal.instance.mysql.portcanal.instance.mysql.port = 3306MySQL 端口默认 3306
canal.instance.connectionCharsetcanal.instance.connectionCharset = UTF-8数据库连接字符集建议设为 UTF-8,避免乱码
canal.instance.master.journal.namecanal.instance.master.journal.name = mysql-bin.000001指定初始解析位点首次启动可不填,从最新位置开始;恢复时需手动指定
canal.instance.master.positioncanal.instance.master.position = 1234
canal.instance.dbUsernamecanal.instance.dbUsername = canal连接 MySQL 的用户名密码必须具有 REPLICATION 权限
canal.instance.dbPasswordcanal.instance.dbPassword = canal
canal.instance.filter.regexcanal.instance.filter.regex = db1\\.table1,db2\\..*表级过滤正则表达式支持库名、表名匹配,. 需转义为 \\.
启动命令sh bin/startup.sh启动 Canal Server日志位于 logs/canal/canal.loglogs/example/example.log
停止命令sh bin/stop.sh安全停止服务避免直接 kill

2.4 Canal Admin 管理后台部署(可选)

配置项 / 步骤语法 / 示例用途注意事项
下载 canal-admin从 release 包中获取 canal.admin 模块独立的 Web 管理后台需与 canal-deployer 配合使用
数据库准备创建数据库 canal_manager,执行 conf/canal_manager.sql存储配置元数据支持 MySQL 5.7+
canal.admin.managercanal.admin.manager = 192.168.1.200:8089指定 Admin 服务地址部署后填写实际 IP 和端口
canal.admin.portcanal.admin.port = 8089Admin 服务端口默认 8089
canal.admin.usercanal.admin.user = admin登录账号密码(明文存储)建议首次登录后修改
canal.admin.passwdcanal.admin.passwd = admin
启动 Adminsh bin/startup.sh local本地模式启动 Admin生产环境建议外接数据库
访问地址http://<ip>:8089Web 管理界面默认账号密码:admin/admin
关联 Deployer在 Admin 界面中添加 Server 节点 IP 和端口实现集中管理 instance 配置需确保网络互通

2.5 启动与验证 Canal 服务

操作 / 命令语法 / 示例用途注意事项
启动 Canal Serversh bin/startup.sh启动服务进程检查 logs/canal/canal.log 是否出现 “Startup successfully”
查看 instance 日志tail -f logs/example/example.log观察具体 instance 是否连接 MySQL 成功成功标志:connect to mysql success
检查端口监听netstat -anp | grep 11111确认 Canal 服务端口已监听若未监听,检查 canal.port 配置
使用 telnet 测试连接telnet <canal-ip> 11111验证网络可达性成功连接表示服务正常
查看 ZooKeeper 节点echo dump | nc <zk-ip> 2181 | grep canal检查 Canal 是否注册到 ZooKeeperHA 模式下应看到 /otter/canal 节点
修改 MySQL 数据INSERT INTO test_table(id, name) VALUES(1, 'test');触发 Binlog 写入确保表在过滤规则范围内
查看日志输出example.log 中查找 new entry 或 DML 记录验证 Canal 是否成功解析变更成功示例:DML , tablename=test_table , sql=INSERT INTO ...
停止服务sh bin/stop.sh安全关闭避免 abrupt shutdown 导致位点丢失

提示: 本章为部署实操章节,建议按顺序完成 MySQL → ZooKeeper → Canal Server → Admin → 验证 的完整流程。下一章将进入 Canal Client 编程接口,涉及具体 API 使用。

第3章:Canal Client 通信模型与 API 入门

3.1 Canal Client 通信协议(TCP / gRPC)

通信协议说明适用场景注意事项
TCP 协议Canal 最早支持的通信方式,基于自定义二进制协议,轻量高效。单机部署、简单集成、低延迟场景需手动管理连接与重试,不支持服务发现
gRPC 协议基于 HTTP/2 的高性能 RPC 框架,支持双向流、服务发现、TLS 加密。微服务架构、Kubernetes 部署、高安全要求场景需引入 gRPC 依赖,配置稍复杂,Canal Server 需启用 gRPC 端口(默认 11112)
协议选择建议TCP:简单、稳定、兼容性好
gRPC:现代化、可扩展、支持流控
根据架构复杂度选择生产环境推荐 gRPC,尤其是多实例集群场景

3.2 单机模式连接 Canal Server

方法名语法用途代码示例注意事项
CanalConnector.newSingleConnectorCanalConnector.newSingleConnector(String address, int port, String destination, String username, String password)创建单机连接器CanalConnector connector = CanalConnectors.newSingleConnector("192.168.1.100", 11111, "example", "canal", "canal");destination 对应 instance 名称(如 example)
connector.connect()void connect()建立与 Canal Server 的连接connector.connect();必须先连接才能订阅
connector.checkValid()void checkValid()检查连接是否有效connector.checkValid();可用于断线重连检测
connector.disconnect()void disconnect()断开连接connector.disconnect();建议在 finally 块中调用

3.3 集群模式连接(通过 ZooKeeper)

方法名语法用途代码示例注意事项
CanalConnectors.newClusterConnectorCanalConnectors.newClusterConnector(String zkServers, String destination, String username, String password)创建集群连接器,通过 ZooKeeper 发现 ServerCanalConnector connector = CanalConnectors.newClusterConnector("192.168.1.10:2181", "example", "canal", "canal");zkServers 为 ZooKeeper 集群地址
connector.connect()void connect()连接并自动选择可用 Serverconnector.connect();若当前 Server 宕机,会自动切换
connector.subscribe()void subscribe(String filter)订阅数据,可指定表过滤规则connector.subscribe("test\\.user,product\\..*");filter 语法同 instance.properties
connector.unsubscribe()void unsubscribe()取消订阅connector.unsubscribe();一般不需手动调用

3.4 订阅数据的基本流程(connect、subscribe、get、ack)

方法名语法用途代码示例注意事项
connect()connector.connect()建立连接见上节必须第一步调用
subscribe()connector.subscribe(filter)发送订阅请求connector.subscribe("mydb\\.mytable");可为空字符串订阅所有表
getWithoutAck()Message getWithoutAck(int batchSize)
Message getWithoutAck(int batchSize, long timeout, TimeUnit unit)
拉取一批未确认的消息Message msg = connector.getWithoutAck(100);
Message msg = connector.getWithoutAck(100, 3, TimeUnit.SECONDS);
推荐使用带超时版本,避免阻塞
ack()void ack(long batchId)确认消费成功,提交位点connector.ack(msg.getId());必须在处理完成后调用,否则重复消费
rollback()void rollback()
void rollback(long batchId)
回滚位点,重新消费connector.rollback();
connector.rollback(msg.getId());
用于消费失败时重试
get()Message get(int batchSize)拉取并自动 ack(不推荐)Message msg = connector.get(100);容易丢数据,建议用 getWithoutAck + ack 组合

3.5 解析 Entry 与 RowData 数据结构

结构 / 字段说明示例代码片段注意事项
EntryBinlog 事件的最小单位,包含 header 和 entryTypefor (Entry entry : msg.getEntries()) {
  if (entry.getEntryType() == EntryType.ROWDATA) { ... }
}
只处理 ROWDATA 类型
Header包含日志文件名、位置、时间戳、serverId、schemaName、tableNameString logfileName = entry.getHeader().getLogfileName();
String tableName = entry.getHeader().getTableName();
用于定位变更来源
RowChange封装行变更信息,包含事件类型(INSERT/UPDATE/DELETE)和行数据列表RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
EventType eventType = rowChange.getEventType();
必须反序列化 storeValue
RowData单行数据,包含 beforeColumns 和 afterColumns 列表List<Column> afterColumns = rowData.getAfterColumnsList();INSERT 只有 after,DELETE 只有 before
Column列信息,包含 name、value、isUpdate、isNullfor (Column col : afterColumns) {
  System.out.println(col.getName() + ": " + col.getValue());
}
isUpdate 表示该列是否被修改(UPDATE 时有效)

第4章:Canal 配置详解

4.1 instance.properties 配置详解

配置项示例值用途注意事项
canal.instance.mysql.host192.168.1.100MySQL 主库 IP支持域名
canal.instance.mysql.port3306MySQL 端口
canal.instance.master.journal.namemysql-bin.000001起始 Binlog 文件名首次可不填
canal.instance.master.position1234起始 Binlog 位置
canal.instance.master.timestamp1672531200000起始时间戳(毫秒)优先级高于 position
canal.instance.master.serverId1MySQL server-id
canal.instance.dbUsernamecanal连接用户名
canal.instance.dbPasswordcanal连接密码明文存储,注意权限
canal.instance.connectionCharsetUTF-8字符集推荐 UTF-8
canal.instance.filter.regexdb1\\.table1,db2\\..*白名单过滤支持正则
canal.instance.filter.black.regextest\\.temp.*黑名单过滤优先级高于 white
canal.instance.tsdb.enabletrue是否启用时间序列数据库用于解决 DDL 问题

4.2 canal.properties 全局配置说明

配置项示例值用途注意事项
canal.id1当前 Canal 实例 ID集群中唯一
canal.zkServers192.168.1.10:2181ZooKeeper 地址HA 模式必填
canal.port11111TCP 端口
canal.metrics.pull.port11112gRPC 端口
canal.admin.managerhttp://192.168.1.200:8089Admin 地址启用 Admin 时填写
canal.admin.port11110Admin 端口
canal.admin.useradminAdmin 用户名
canal.admin.passwdadminAdmin 密码(明文)建议修改
canal.instance.global.modespring全局 instance 模式
canal.instance.global.lazyfalse是否延迟加载 instance

4.3 过滤规则配置(filter、blacklist)

规则类型语法示例说明注意事项
单表db1.table1只同步该表
所有表db1\\..*db1 库下所有表. 需转义为 \\.
多表db1\.table1,db2\.table2用逗号分隔
通配符\\.\\test$匹配表名以 test 结尾支持正则表达式
黑名单canal.instance.filter.black.regex = test\\.temp.*排除匹配表优先级高于白名单
注意事项区分大小写、特殊字符需转义、建议测试环境验证规则

4.4 HA 与负载均衡配置

配置项示例值用途注意事项
canal.zkServerszk1:2181,zk2:2181指定 ZooKeeper 集群多个 Server 注册至此
canal.instance.global.spring.xmlclasspath:spring/default-instance.xmlinstance 加载方式
canal.serverModetcpkafkarocketmq输出模式HA 通常配合消息队列
cluster 模式连接CanalConnectors.newClusterConnector(...)Client 自动发现可用 Server避免单点故障
故障转移自动当前 Server 宕机,Client 切换到其他 Server依赖 ZooKeeper 选举机制

4.5 日志与监控配置

配置项文件路径用途注意事项
canal.log.levelINFO设置日志级别可设为 DEBUG 调试
canal.log.dirlogs/日志目录
主日志logs/canal/canal.logCanal Server 启动与运行日志查看是否启动成功
Instance 日志logs/example/example.log每个 instance 的解析日志查看 Binlog 解析情况
错误日志logs/canal/canal.pid + .err异常堆栈
JMX 监控com.alibaba.canal:type=destinationJMX MBean 暴露可集成 Prometheus
canal.metrics.enabletrue是否启用指标收集
canal.metrics.pull.port11112指标拉取端口

第5章:Canal 数据消费模式与高级特性

5.1 FlatMessage 模式介绍与使用

概念 / 配置说明代码示例注意事项
FlatMessage将复杂的 Entry 结构扁平化为 JSON 格式,便于传输和消费,常用于 Kafka/RocketMQ 消息体{ "database": "test", "table": "user", "type": "INSERT", "ts": 1672531200000, "data": [ { "id": "1", "name": "Alice" } ] }无需解析 Protocol Buffer,适合跨语言消费
启用 FlatMessage在 instance 配置中设置 canal.instance.tsdb.enable = true(非必须);输出到 MQ 时自动转换canal.mq.flatMessage = true(在 canal.properties 中)仅在 MQ 模式下生效
数据结构字段data: 变更后数据(INSERT/UPDATE),old: 变更前数据(UPDATE/DELETE),type: DML 类型,ts: 时间戳见上例UPDATE 操作同时包含 data 和 old
使用场景同步到 Kafka/RocketMQ、跨语言消费者(Python/Go)、简化数据解析避免在 Java Client 中重复解析 PB
注意事项不支持 DDL 事件的扁平化需确保消费者能处理 JSON 格式

5.2 批量拉取与流式拉取对比

拉取模式说明代码示例优点缺点适用场景
批量拉取一次获取多条 Entry(如 100 条),适合高吞吐Message msg = connector.getWithoutAck(100);吞吐高,网络开销小延迟较高,内存占用大数据同步、批量处理
流式拉取逐条处理,每次只取一条或少量数据while (true) {
  Message msg = connector.getWithoutAck(1, 3, TimeUnit.SECONDS);
  // 处理单条
  connector.ack(msg.getId());
}
实时性高,内存友好吞吐较低,频繁调用实时缓存更新、审计日志
带超时拉取推荐方式,避免无限阻塞getWithoutAck(100, 3, TimeUnit.SECONDS)平衡吞吐与延迟需合理设置超时通用场景
注意事项批量过大可能导致 OOM,流式需注意 ack 频率建议批量大小 50500,超时 310 秒

5.3 位点(Position)管理与手动提交

方法 / 配置语法用途代码示例注意事项
getWithoutAck()Message getWithoutAck(int batchSize, long timeout, TimeUnit unit)获取一批未确认消息,返回 batchIdMessage msg = connector.getWithoutAck(100, 3, TimeUnit.SECONDS);
long batchId = msg.getId();
batchId 用于后续 ack 或 rollback
ack(long batchId)void ack(long batchId)提交位点,表示该批次已成功处理connector.ack(batchId);必须在业务处理成功后调用
rollback()void rollback()回滚到上一次 ack 位置connector.rollback();用于消费失败重试
rollback(long batchId)void rollback(long batchId)回滚到指定 batchIdconnector.rollback(batchId);仅限未 ack 的 batchId
手动位点管理不依赖自动提交,由业务控制结合 try-catch 使用try {
  process(msg);
  connector.ack(msg.getId());
} catch (Exception e) {
  connector.rollback();
}
避免数据丢失或重复消费
注意事项不 ack 会导致重复消费,错误 rollback 可能跳过数据建议记录 batchId 用于追踪

5.4 多表订阅与动态过滤

方法 / 配置语法用途代码示例注意事项
subscribe(String filter)void subscribe(String filter)订阅时指定过滤规则connector.subscribe("db1\\.user,db2\\.order.*");支持正则表达式
动态切换过滤可多次调用 subscribe()修改订阅规则connector.subscribe("new_db\\..*");会覆盖之前规则
白名单过滤canal.instance.filter.regexinstance 级白名单db1\\.t1,db2\\..*
黑名单过滤canal.instance.filter.black.regex排除某些表test\\.temp.*优先级高于白名单
通配符说明\\. 表示点,.* 表示任意字符正则匹配^user_\\d+$ 匹配 user_1, user_2注意转义
注意事项过滤在 Server 端生效,修改后需重新 connect/subscribe建议先在测试环境验证正则

5.5 数据转换与自定义处理器

方法 / 类说明代码示例注意事项
自定义 EventSink继承 AbstractEventSink,重写 onEventpublic class CustomSink extends AbstractEventSink<RowData> {
  @Override
  protected boolean onEvent(RowData row) {
    // 转换逻辑
    return true;
  }
}
高级用法,需修改 instance 配置
在 Client 中转换拉取后进行业务处理List<Column> cols = rowData.getAfterColumnsList();
Map<String, Object> map = new HashMap<>();
for (Column col : cols) {
  map.put(col.getName(), col.getValue());
}
推荐方式,灵活
转换为 POJO映射到 Java 对象User user = new User();
user.setId(Long.valueOf(getValue("id")));
注意类型转换异常
输出到外部系统如 Redis、HTTP、文件jedis.set("user:" + id, JSON.toJSONString(map));注意异常处理与重试
注意事项避免在 onEvent 中阻塞,转换失败应记录日志并 rollback建议异步处理耗时操作

第6章:Canal 与主流中间件集成

6.1 Canal + Kafka 集成配置与使用

配置项示例值用途注意事项
canal.serverModekafka设置输出模式为 Kafka必须配置
canal.mq.servers192.168.1.10:9092Kafka broker 地址
canal.mq.retries3发送失败重试次数
canal.mq.batchSize16384批量发送大小
canal.mq.topiccanal-topic默认 Topic
canal.mq.dynamicTopicdb1.user:topic1,db2.order:topic2动态 Topic 路由支持正则
canal.mq.flatMessagetrue使用 FlatMessage 格式推荐开启
Kafka 消费者使用 KafkaConsumer<String, String>消费 JSON 消息需解析 FlatMessage 结构
注意事项确保 Kafka 可写,Topic 需提前创建或开启自动创建建议监控 Kafka 延迟

6.2 Canal + RocketMQ 集成实践

配置项示例值用途注意事项
canal.serverModerocketmq输出到 RocketMQ
canal.mq.servers192.168.1.10:9876NameServer 地址
canal.mq.producerGroupcanal-producer生产者组名
canal.mq.topiccanal-rocket-topic默认 Topic
canal.mq.dynamicTopicdb1\\..*:topic_db1动态 Topic
canal.mq.accessChannellocalremote访问模式
canal.mq.flatMessagetrue使用 JSON 格式
RocketMQ 消费使用 DefaultMQPushConsumer订阅 Topic需处理消息顺序性
注意事项依赖 canal.mq.rocketmq.core 模块,确保 NameServer 可达RocketMQ 需开启 ACL 时配置密钥

6.3 Canal Adapter 数据同步到 Elasticsearch

配置文件路径说明示例内容注意事项
application.ymlconf/application.yml启用 ES 模式canalAdapters:
  - instance: example
  groups:
  - groupId: g1
  outerAdapters:
  - name: es
  hosts: 192.168.1.200:9200
  properties:
  mode: rest
hosts 为 ES 地址
es/*.ymlconf/es/test_user.yml映射配置dataSourceKey: defaultDS
destination: example
groupId: g1
esMapping:
  _index: user
  _type: _doc
  _id: _id
  sql: "select id as _id, name, email from user"
SQL 必须包含主键映射
同步操作INSERT → ES Index,UPDATE → ES Update,DELETE → ES Delete自动转换需确保 ES 字段类型匹配
注意事项ES 版本兼容性(6.x/7.x/8.x),高频更新建议使用 bulk支持嵌套对象、数组映射

6.4 Canal Adapter 同步 MySQL 到其他数据库

配置文件说明示例内容注意事项
application.yml配置数据源spring:
  datasource:
  canal:
  url: jdbc:mysql://192.168.1.100:3306/target_db?useUnicode=true
  username: root
  password: 123456
支持多数据源
rdb/*.yml映射规则dataSourceKey: canal
destination: example
groupId: g1
outerAdapterKey: mysql1
concurrent: true
dbMapping:
  database: target_db
  table: user_copy
  mapAll: true
mapAll: true 表示字段全映射
支持数据库MySQL、Oracle、PostgreSQL、SQL Server通过 JDBC 连接需放入对应 JDBC 驱动到 lib/
同步类型INSERT、UPDATE、DELETE自动转换为对应 SQL注意外键约束
注意事项目标表需预先创建,字段类型需兼容,支持字段映射 mapping不支持 DDL 同步

6.5 自定义 Sink 输出(如 HTTP、Redis)

输出类型实现方式示例代码片段注意事项
HTTP使用 OkHttpClient 或 RestTemplateRequest request = new Request.Builder()
  .url("http://api.example.com/webhook")
  .post(RequestBody.create(json, MediaType.get("application/json")))
  .build();
client.newCall(request).execute();
注意超时与重试机制
Redis使用 Jedis 或 LettuceJedis jedis = new Jedis("192.168.1.50", 6379);
jedis.set("canal:data:" + id, json);
jedis.close();
建议使用连接池
File写入本地文件Files.write(Paths.get("output.log"), (json + "\n").getBytes(), StandardOpenOption.APPEND);注意磁盘空间
自定义 Sink 类继承 AbstractCanalClientTest 或实现 SinkFunction在 process 方法中添加输出逻辑适用于 Flink/Spark 流处理
注意事项异常处理(网络失败、序列化错误),异步化避免阻塞 Canal建议引入消息队列做缓冲

第7章:Canal 运维与监控

7.1 常见异常与排查方法

异常现象可能原因排查方法解决方案注意事项
connect to mysql failedMySQL 网络不通、用户权限不足、Binlog 未开启telnet host 3306、检查 canal.instance.dbUsername 权限、SHOW VARIABLES LIKE 'log_bin';开启 Binlog,授权 REPLICATION 权限确保 server-id 唯一
command : -1客户端与 Server 版本不兼容检查客户端与 Server 的 Canal 版本是否一致升级为相同版本建议使用 v1.1.7+ 稳定版
ack error, batchId not exist提交了已 ack 或不存在的 batchId检查代码是否重复 ack 或 rollback 到已提交位点避免在 finally 中错误 rollback记录 batchId 用于追踪
parse row data failed表结构变更(DDL)导致列不匹配查看 example.log 是否有 DDL 记录,检查 tsdb 是否启用启用 canal.instance.tsdb.enable=true 或重启 instance建议开启 TSDB 避免解析失败
ZooKeeper session expiredZK 会话超时,网络抖动检查 ZK 集群状态:echo stat | nc zk-host 2181调大 canal.zookeeper.sessionTimeout,优化网络建议 ZK 集群部署在独立节点

7.2 位点漂移问题与解决方案

问题描述原因分析解决方案注意事项
位点回退(重复消费)未调用 ack、错误调用 rollback()、MetaManager 存储异常确保处理成功后调用 ack(batchId)、避免在异常处理中无条件 rollback、使用 ZooKeeper 存储位点建议在业务事务提交后 ack
位点跳跃(数据丢失)手动修改了 meta.dat、指定错误的 journal.name 和 position不要手动修改位点文件、通过 Admin 或日志确定正确位点恢复前先备份 meta.dat
位点不持久化使用内存 MetaManager,Server 重启后丢失配置 canal.instance.global.metaManager = zookeeper生产环境必须使用 ZK 或 file
多 Client 位点冲突多个 Client 订阅同一 destination使用集群模式(ClusterConnector),由 ZK 协调单 destination 应只有一个活跃 Consumer
建议实践使用 getWithoutAck + ack 模式,记录 batchId 与处理状态位点管理是数据一致性的关键

7.3 性能调优建议(批大小、线程数等)

调优项配置项推荐值说明注意事项
批量拉取大小getWithoutAck(batchSize)50 ~ 500提高吞吐,减少网络调用过大会增加内存压力
拉取超时getWithoutAck(..., timeout, unit)3 ~ 10 秒避免无限阻塞太短可能导致空轮询
解析线程数canal.instance.parser.parallelCPU 核数 ~ 2×核数并行解析 Binlog开启需设置 parallel = true
MQ 发送批次canal.mq.batchSize8192 ~ 32768提高 Kafka/RocketMQ 吞吐需平衡延迟与吞吐
内存队列容量canal.instance.memory.bufferSize16384 ~ 32768缓冲 Binlog 事件过小可能导致阻塞
网络参数canal.tcp.soTimeout60000Socket 超时(毫秒)防止连接挂起
建议监控 GC 与内存使用,压测调优参数调优需结合业务场景与硬件资源

7.4 监控指标(JMX、Prometheus)

指标类别JMX MBean 名称Prometheus 路径关键指标说明
实例状态destination=example/metricsrunning、dumpRunning、parseRunning查看 instance 是否运行
解析延迟ParseCountercanal_parse_delay{destination="example"}delay(毫秒)Binlog 时间与解析时间差,核心延迟指标
消费位点MetaCountercanal_mq_position{topic="..."}position、filename当前消费的 Binlog 位置
消息吞吐SinkCountercanal_sink_qps{destination="..."}QPS、TPS每秒处理的消息数
连接状态ClientIdentitycanal_client_count{destination="..."}客户端连接数监控消费者活跃度
JVM 指标java.lang内置堆内存、GC、线程数基础资源监控
配置启用canal.metrics.pull.port=11112访问 http://ip:11112/metrics 获取需开启 canal.metrics.enable=true

7.5 升级与回滚策略

操作步骤注意事项
升级前准备备份 conf/logs/,记录当前位点(meta.dat 或 ZK),停止 Canal Server建议在低峰期操作
升级步骤替换 canal.deployer 包,合并对新版本的配置变更,启动新版本 Server,验证日志与数据同步检查版本兼容性(如 v1.1.4 → v1.1.7)
回滚步骤停止新版本,恢复旧版本包和配置,确认位点文件未被覆盖,启动旧版本若位点已更新,可能需手动回退
灰度发布先升级一个节点,观察稳定后再升级其他节点适用于集群环境
注意事项升级前测试兼容性,避免跨大版本直接升级(如 1.0 → 1.1),关注官方 Release Notes建议使用 Canal Admin 管理升级

第8章:Canal 源码解析(进阶)

8.1 Canal Server 启动流程分析

阶段核心类 / 方法说明注意事项
1. Main 入口CanalLauncher.main()启动类,解析命令行参数
2. 配置加载CanalController.initGlobalConfig()加载 canal.properties设置全局参数
3. Instance 初始化CanalController.initInstance()创建 CanalInstance,加载 instance.properties支持多 instance
4. 组件启动CanalInstanceWithSpring.start()依次启动:EventPublisher → EventParser → EventSink顺序不能错
5. 网络服务启动CanalServerWithNetty.start()启动 TCP/gRPC 服务,监听客户端连接默认端口 11111/11112
6. 运行中EventParser 持续 dump Binlog → EventSink 处理 → EventStore 存储数据流管道故障时自动重试

8.2 Binlog 解析核心类(EventParser)

类名说明核心方法注意事项
AbstractEventParser抽象解析器,定义模板流程start()stop()run()
MYSQLIncrementalParser增量 Binlog 解析器parseRowsEvent()处理 ROW 模式事件
MysqlConnection封装 MySQL 连接与 dump 协议dump(long position)模拟 Slave 发起 dump
LogEventConvert将 Binlog Event 转换为 EntryparseRowsEvent(...)核心转换逻辑
EventStore内存存储解析后的 Entryput(Entry)get(...)默认实现为 MemoryEventStoreWithBuffer
工作流程连接 MySQL → dump Binlog → 解析为 Entry → 存入 EventStore → 通知 EventSink解析失败会重试

8.3 Entry 与 Protocol 结构解析

结构说明字段示例注意事项
Entry顶层结构,包含 header 和 storeValueentryType: ROWDATA, storeValue: bytesstoreValue 为 PB 序列化数据
Header元信息logfileName, position, timestamp, schemaName, tableName, eventType用于定位和路由
RowChange行变更封装eventType: INSERT, rowDatasList需反序列化 storeValue
RowData单行数据beforeColumns, afterColumnsINSERT 无 before,DELETE 无 after
Column列信息name, value, isUpdate, isNullisUpdate 仅 UPDATE 有效
Protocol Buffer使用 Google Protobuf 定义 .proto 文件entry.proto, logproxy.proto编译生成 Java 类

8.4 Client-Server 通信协议实现

协议核心类说明注意事项
TCP 协议CanalServerWithNettySimpleCanalConnector基于 Netty 实现 TCP 长连接自定义二进制协议
gRPC 协议GrpcCanalServerGrpcCanalConnection基于 gRPC-Netty 实现支持流式传输
消息类型ClientAuth、Get、Ack、Rollback定义在 canal.proto使用 Protobuf 编码
认证流程client → auth → server传输 destination、username、password失败则断开连接
数据拉取GetRequest → GetResponse返回 Message(含 batchId 和 entries)支持带超时拉取
心跳机制PingRequest客户端定期发送保持连接活跃

8.5 Canal Adapter 工作机制

模块核心类说明注意事项
CanalAdapterApplication主启动类加载 Spring 上下文
OuterAdapter输出适配器接口ESAdapter、RdbAdapter、LoggerAdapter可扩展
GroupEventProcessor事件处理组拉取 Canal 数据并分发给适配器支持并发
Loader配置加载器加载 application.yml 和 es/*.yml
Dml 转换DataLoader将 DML 转为目标操作(如 ES Index)支持 insert/update/delete
工作流程1. 从 Canal Server 拉取数据;2. 解析为 Dml 对象;3. 根据配置路由到对应 OuterAdapter;4. 执行写入操作支持多 destination 和多 group

第9章:实战案例

9.1 实时同步 MySQL 到 Elasticsearch

项目内容说明
场景需求将 MySQL 中的商品、订单、用户表实时同步到 ES,用于搜索与分析要求低延迟(秒级)、数据一致
技术架构MySQL → Canal Server → Canal Adapter (ES) → Elasticsearch使用 Canal Adapter 作为中间桥梁

核心配置:

  1. conf/application.yml
canalAdapters:
  - instance: example
    groups:
      - groupId: g1
        outerAdapters:
          - name: es
            hosts: 192.168.1.200:9200
            properties:
              mode: rest

hosts 为 ES 地址

  1. conf/es/product.yml
dataSourceKey: defaultDS
destination: example
groupId: g1
esMapping:
  _index: product
  _type: _doc
  _id: id
  sql: "SELECT id, name, price, category FROM product"
  commitBatch: 3000

sql 必须包含主键映射,commitBatch 控制批量提交大小

同步行为说明
INSERT → ES index自动转换,无需手动写入
UPDATE → ES update
DELETE → ES delete

注意事项:

  • ES 索引需预先创建或启用自动创建
  • 字段类型需匹配(如 price 为 float)
  • 高频更新建议开启 bulk 写入
  • 支持嵌套对象(通过 join 查询)
  • 建议开启 canal.mq.flatMessage=true 提升性能

9.2 构建缓存更新系统(MySQL → Redis)

项目内容说明
场景需求当用户表更新时,自动清除或更新 Redis 缓存,避免脏数据要求高实时性、低延迟
技术架构MySQL → Canal Server → 自定义 Client → Redis使用 Java Client 监听变更

核心代码:

Message msg = connector.getWithoutAck(100, 3, TimeUnit.SECONDS);
for (Entry entry : msg.getEntries()) {
    if (entry.getEntryType() == EntryType.ROWDATA) {
        RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
        String tableName = entry.getHeader().getTableName();
        if ("user".equals(tableName)) {
            for (RowData rowData : rowChange.getRowDatasList()) {
                String id = getColumnValue(rowData.getAfterColumnsList(), "id");
                jedis.del("user:" + id); // 删除缓存
                // 或:jedis.setex("user:" + id, 3600, JSON.toJSONString(rowData)); // 更新缓存
            }
        }
    }
}
connector.ack(msg.getId());

getColumnValue 为自定义方法,建议异步更新缓存

缓存策略说明
删除策略更新后删除缓存,下次读取时重建(推荐,更简单安全)
更新策略直接更新缓存值(需处理并发)

注意事项:

  • 使用连接池(JedisPool)
  • 处理 Redis 网络异常(重试机制)
  • 避免缓存穿透/击穿
  • 可结合 Kafka 异步解耦
  • 生产环境建议引入熔断机制(如 Hystrix)

9.3 数据库审计日志采集

项目内容说明
场景需求记录所有对敏感表(如用户、订单)的增删改操作,用于安全审计要求完整、不可篡改、可追溯
技术架构MySQL → Canal Server → Kafka → Flink → 存储(HDFS/ClickHouse)使用 Kafka 解耦,Flink 实时处理

Canal Server 配置:

canal.serverMode = kafka
canal.mq.topic = audit-log
canal.mq.flatMessage = true
canal.instance.filter.regex = user\\.info,order\\.main

flatMessage 输出 JSON,过滤只关注的表

Kafka 消费者使用 Flink 消费并写入审计库。

审计字段说明
操作类型INSERT/UPDATE/DELETE
操作时间timestamp
操作人可从上下文获取
旧值(before)与新值(after)
IP 地址、客户端信息可在 Flink 中补充上下文信息
存储层说明
短期:MySQL保留 3 个月
长期:ClickHouse压缩比高,查询快
归档:HDFS + Parquet建议按天/月分区

注意事项:

  • 审计日志不可删除(合规性要求高,需定期备份)
  • 记录 Binlog 位点用于溯源
  • 可结合用户行为日志关联分析

9.4 多级级联同步架构设计

项目内容说明
场景需求跨数据中心同步:IDC-A → IDC-B → IDC-C,实现数据逐级分发避免跨地域直接同步,降低延迟与带宽压力

架构:

MySQL (IDC-A)

Canal Server A → Kafka A

Canal Client A → MySQL (IDC-B)

Canal Server B → Kafka B

Canal Client B → MySQL (IDC-C)

每级独立部署 Canal

实现方式说明
级联模式上一级的 MySQL 作为下一级的源库
消息中转(推荐)通过 Kafka 跨地域传输 FlatMessage,解耦更强
关键项说明
位点管理每级独立管理位点,使用 ZooKeeper 或本地文件,避免单点故障
延迟控制每级延迟 < 1s,总延迟 < 3s,使用 canal_parse_delay 监控,优化批大小与网络带宽

注意事项:

  • 避免循环同步(需过滤目标表)
  • 每级需做数据校验
  • 故障时支持断点续传
  • 建议使用 gRPC 提升跨地域传输效率
  • 架构复杂,运维成本高

9.5 高可用双活架构中的 Canal 应用

项目内容说明
场景需求两个数据中心(IDC-A 和 IDC-B)同时提供读写服务,数据双向同步实现 RPO ≈ 0,RTO < 30s

架构设计:

IDC-A: MySQL-A ↔ Canal-A ↔ Kafka-A ↔ Canal-B ↔ MySQL-B :IDC-B
IDC-B: MySQL-B ↔ Canal-B ↔ Kafka-B ↔ Canal-A ↔ MySQL-A :IDC-A

双向同步,互为备份

核心挑战解决方案
循环同步问题(A → B → A 导致数据重复)位点标记法(推荐):在消息头添加 source_id,接收方忽略来自自己的消息
时间戳去重记录最后更新时间,跳过旧数据
唯一键冲突处理使用分布式锁或版本号

Canal 配置要点:

  • canal.instance.filter.black.regex = sync_marker_table(过滤同步标记表)
  • 使用 Kafka 作为中间件,支持跨 IDC 传输
  • 启用 flatMessage 简化处理
切换流程说明
1. 检测主库故障需自动化脚本支持
2. 停止反向同步(B → A)
3. 提升备库为可写
4. 恢复同步(A → B)

注意事项:

  • 避免并发写同一行(需业务层控制)
  • 定期数据一致性校验
  • 使用 ZooKeeper 协调状态
  • 建议结合 OTS(操作时间戳)解决冲突
  • 双活架构复杂,需充分测试