第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 Slave | Canal Server 启动后,伪装成 MySQL 的从库(Slave),向主库发起 dump 协议请求,获取 Binlog 流。 | 需在 MySQL 创建具有 REPLICATION SLAVE 和 REPLICATION 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 Admin | Web 管理后台,用于集中管理多个 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-bin | log-bin = mysql-bin | 启用二进制日志,文件前缀为 mysql-bin | 必须开启,Canal 依赖 Binlog 获取数据变更 |
| binlog-format | binlog-format = ROW | 设置 Binlog 格式为 ROW 模式 | 必须为 ROW,STATEMENT 或 MIXED 不支持精确行级解析 |
| server-id | server-id = 1 | 设置唯一服务器 ID,用于主从复制标识 | 每个 MySQL 实例必须唯一,否则无法开启 Binlog |
| expire-logs-days | expire-logs-days = 7 | Binlog 文件保留天数 | 根据 Canal 消费速度设置,避免日志过早清理导致位点失效 |
| binlog-row-image | binlog-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 已启用且格式正确 | 返回值应为 ON 和 ROW |
2.2 ZooKeeper 搭建与配置(集群模式)
| 配置项 / 命令 | 语法 / 示例 | 用途 | 注意事项 |
|---|
| tickTime | tickTime=2000 | ZooKeeper 心跳间隔(毫秒) | 通常设为 2000,不宜过大或过小 |
| dataDir | dataDir=/var/lib/zookeeper | 数据存储目录,包含 myid 和快照 | 目录需存在且可写,每个节点路径不同 |
| clientPort | clientPort=2181 | 客户端连接端口 | Canal Server 和 Client 将连接此端口 |
| initLimit | initLimit=10 | Follower 初始连接 Leader 的超时时间(以 tickTime 为单位) | 一般为 1015,表示 2030 秒 |
| syncLimit | syncLimit=5 | Follower 与 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 列表中定义 |
| 启动 ZooKeeper | bin/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.port | canal.port = 11111 | Canal Server TCP 服务端口 | 多实例部署时需区分端口 |
| canal.zkServers | canal.zkServers = 192.168.1.10:2181,192.168.1.11:2181 | 指定 ZooKeeper 集群地址 | 若使用 HA 模式,此项必填 |
| canal.instance.mysql.host | canal.instance.mysql.host = 192.168.1.100 | MySQL 主库地址 | 支持域名或 IP |
| canal.instance.mysql.port | canal.instance.mysql.port = 3306 | MySQL 端口 | 默认 3306 |
| canal.instance.connectionCharset | canal.instance.connectionCharset = UTF-8 | 数据库连接字符集 | 建议设为 UTF-8,避免乱码 |
| canal.instance.master.journal.name | canal.instance.master.journal.name = mysql-bin.000001 | 指定初始解析位点 | 首次启动可不填,从最新位置开始;恢复时需手动指定 |
| canal.instance.master.position | canal.instance.master.position = 1234 | | |
| canal.instance.dbUsername | canal.instance.dbUsername = canal | 连接 MySQL 的用户名密码 | 必须具有 REPLICATION 权限 |
| canal.instance.dbPassword | canal.instance.dbPassword = canal | | |
| canal.instance.filter.regex | canal.instance.filter.regex = db1\\.table1,db2\\..* | 表级过滤正则表达式 | 支持库名、表名匹配,. 需转义为 \\. |
| 启动命令 | sh bin/startup.sh | 启动 Canal Server | 日志位于 logs/canal/canal.log 和 logs/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.manager | canal.admin.manager = 192.168.1.200:8089 | 指定 Admin 服务地址 | 部署后填写实际 IP 和端口 |
| canal.admin.port | canal.admin.port = 8089 | Admin 服务端口 | 默认 8089 |
| canal.admin.user | canal.admin.user = admin | 登录账号密码(明文存储) | 建议首次登录后修改 |
| canal.admin.passwd | canal.admin.passwd = admin | | |
| 启动 Admin | sh bin/startup.sh local | 本地模式启动 Admin | 生产环境建议外接数据库 |
| 访问地址 | http://<ip>:8089 | Web 管理界面 | 默认账号密码:admin/admin |
| 关联 Deployer | 在 Admin 界面中添加 Server 节点 IP 和端口 | 实现集中管理 instance 配置 | 需确保网络互通 |
2.5 启动与验证 Canal 服务
| 操作 / 命令 | 语法 / 示例 | 用途 | 注意事项 |
|---|
| 启动 Canal Server | sh 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 是否注册到 ZooKeeper | HA 模式下应看到 /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.newSingleConnector | CanalConnector.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.newClusterConnector | CanalConnectors.newClusterConnector(String zkServers, String destination, String username, String password) | 创建集群连接器,通过 ZooKeeper 发现 Server | CanalConnector connector = CanalConnectors.newClusterConnector("192.168.1.10:2181", "example", "canal", "canal"); | zkServers 为 ZooKeeper 集群地址 |
connector.connect() | void connect() | 连接并自动选择可用 Server | connector.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 数据结构
| 结构 / 字段 | 说明 | 示例代码片段 | 注意事项 |
|---|
| Entry | Binlog 事件的最小单位,包含 header 和 entryType | for (Entry entry : msg.getEntries()) { if (entry.getEntryType() == EntryType.ROWDATA) { ... }
} | 只处理 ROWDATA 类型 |
| Header | 包含日志文件名、位置、时间戳、serverId、schemaName、tableName | String 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、isNull | for (Column col : afterColumns) { System.out.println(col.getName() + ": " + col.getValue());
} | isUpdate 表示该列是否被修改(UPDATE 时有效) |
第4章:Canal 配置详解
4.1 instance.properties 配置详解
| 配置项 | 示例值 | 用途 | 注意事项 |
|---|
canal.instance.mysql.host | 192.168.1.100 | MySQL 主库 IP | 支持域名 |
canal.instance.mysql.port | 3306 | MySQL 端口 | |
canal.instance.master.journal.name | mysql-bin.000001 | 起始 Binlog 文件名 | 首次可不填 |
canal.instance.master.position | 1234 | 起始 Binlog 位置 | |
canal.instance.master.timestamp | 1672531200000 | 起始时间戳(毫秒) | 优先级高于 position |
canal.instance.master.serverId | 1 | MySQL server-id | |
canal.instance.dbUsername | canal | 连接用户名 | |
canal.instance.dbPassword | canal | 连接密码 | 明文存储,注意权限 |
canal.instance.connectionCharset | UTF-8 | 字符集 | 推荐 UTF-8 |
canal.instance.filter.regex | db1\\.table1,db2\\..* | 白名单过滤 | 支持正则 |
canal.instance.filter.black.regex | test\\.temp.* | 黑名单过滤 | 优先级高于 white |
canal.instance.tsdb.enable | true | 是否启用时间序列数据库 | 用于解决 DDL 问题 |
4.2 canal.properties 全局配置说明
| 配置项 | 示例值 | 用途 | 注意事项 |
|---|
canal.id | 1 | 当前 Canal 实例 ID | 集群中唯一 |
canal.zkServers | 192.168.1.10:2181 | ZooKeeper 地址 | HA 模式必填 |
canal.port | 11111 | TCP 端口 | |
canal.metrics.pull.port | 11112 | gRPC 端口 | |
canal.admin.manager | http://192.168.1.200:8089 | Admin 地址 | 启用 Admin 时填写 |
canal.admin.port | 11110 | Admin 端口 | |
canal.admin.user | admin | Admin 用户名 | |
canal.admin.passwd | admin | Admin 密码(明文) | 建议修改 |
canal.instance.global.mode | spring | 全局 instance 模式 | |
canal.instance.global.lazy | false | 是否延迟加载 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.zkServers | zk1:2181,zk2:2181 | 指定 ZooKeeper 集群 | 多个 Server 注册至此 |
canal.instance.global.spring.xml | classpath:spring/default-instance.xml | instance 加载方式 | |
canal.serverMode | tcp 或 kafka 或 rocketmq | 输出模式 | HA 通常配合消息队列 |
| cluster 模式连接 | CanalConnectors.newClusterConnector(...) | Client 自动发现可用 Server | 避免单点故障 |
| 故障转移 | 自动 | 当前 Server 宕机,Client 切换到其他 Server | 依赖 ZooKeeper 选举机制 |
4.5 日志与监控配置
| 配置项 | 文件路径 | 用途 | 注意事项 |
|---|
canal.log.level | INFO | 设置日志级别 | 可设为 DEBUG 调试 |
canal.log.dir | logs/ | 日志目录 | |
| 主日志 | logs/canal/canal.log | Canal Server 启动与运行日志 | 查看是否启动成功 |
| Instance 日志 | logs/example/example.log | 每个 instance 的解析日志 | 查看 Binlog 解析情况 |
| 错误日志 | logs/canal/canal.pid + .err | 异常堆栈 | |
| JMX 监控 | com.alibaba.canal:type=destination | JMX MBean 暴露 | 可集成 Prometheus |
canal.metrics.enable | true | 是否启用指标收集 | |
canal.metrics.pull.port | 11112 | 指标拉取端口 | |
第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) | 获取一批未确认消息,返回 batchId | Message 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) | 回滚到指定 batchId | connector.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.regex | instance 级白名单 | db1\\.t1,db2\\..* | |
| 黑名单过滤 | canal.instance.filter.black.regex | 排除某些表 | test\\.temp.* | 优先级高于白名单 |
| 通配符说明 | \\. 表示点,.* 表示任意字符 | 正则匹配 | ^user_\\d+$ 匹配 user_1, user_2 | 注意转义 |
| 注意事项 | 过滤在 Server 端生效,修改后需重新 connect/subscribe | | | 建议先在测试环境验证正则 |
5.5 数据转换与自定义处理器
| 方法 / 类 | 说明 | 代码示例 | 注意事项 |
|---|
| 自定义 EventSink | 继承 AbstractEventSink,重写 onEvent | public 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.serverMode | kafka | 设置输出模式为 Kafka | 必须配置 |
canal.mq.servers | 192.168.1.10:9092 | Kafka broker 地址 | |
canal.mq.retries | 3 | 发送失败重试次数 | |
canal.mq.batchSize | 16384 | 批量发送大小 | |
canal.mq.topic | canal-topic | 默认 Topic | |
canal.mq.dynamicTopic | db1.user:topic1,db2.order:topic2 | 动态 Topic 路由 | 支持正则 |
canal.mq.flatMessage | true | 使用 FlatMessage 格式 | 推荐开启 |
| Kafka 消费者 | 使用 KafkaConsumer<String, String> | 消费 JSON 消息 | 需解析 FlatMessage 结构 |
| 注意事项 | 确保 Kafka 可写,Topic 需提前创建或开启自动创建 | | 建议监控 Kafka 延迟 |
6.2 Canal + RocketMQ 集成实践
| 配置项 | 示例值 | 用途 | 注意事项 |
|---|
canal.serverMode | rocketmq | 输出到 RocketMQ | |
canal.mq.servers | 192.168.1.10:9876 | NameServer 地址 | |
canal.mq.producerGroup | canal-producer | 生产者组名 | |
canal.mq.topic | canal-rocket-topic | 默认 Topic | |
canal.mq.dynamicTopic | db1\\..*:topic_db1 | 动态 Topic | |
canal.mq.accessChannel | local 或 remote | 访问模式 | |
canal.mq.flatMessage | true | 使用 JSON 格式 | |
| RocketMQ 消费 | 使用 DefaultMQPushConsumer | 订阅 Topic | 需处理消息顺序性 |
| 注意事项 | 依赖 canal.mq.rocketmq.core 模块,确保 NameServer 可达 | | RocketMQ 需开启 ACL 时配置密钥 |
6.3 Canal Adapter 数据同步到 Elasticsearch
| 配置文件 | 路径 | 说明 | 示例内容 | 注意事项 |
|---|
application.yml | conf/application.yml | 启用 ES 模式 | canalAdapters: - instance: example groups: - groupId: g1 outerAdapters: - name: es hosts: 192.168.1.200:9200 properties: mode: rest | hosts 为 ES 地址 |
es/*.yml | conf/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 或 RestTemplate | Request 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 或 Lettuce | Jedis 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 failed | MySQL 网络不通、用户权限不足、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 expired | ZK 会话超时,网络抖动 | 检查 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.parallel | CPU 核数 ~ 2×核数 | 并行解析 Binlog | 开启需设置 parallel = true |
| MQ 发送批次 | canal.mq.batchSize | 8192 ~ 32768 | 提高 Kafka/RocketMQ 吞吐 | 需平衡延迟与吞吐 |
| 内存队列容量 | canal.instance.memory.bufferSize | 16384 ~ 32768 | 缓冲 Binlog 事件 | 过小可能导致阻塞 |
| 网络参数 | canal.tcp.soTimeout | 60000 | Socket 超时(毫秒) | 防止连接挂起 |
| 建议 | 监控 GC 与内存使用,压测调优参数 | | | 调优需结合业务场景与硬件资源 |
7.4 监控指标(JMX、Prometheus)
| 指标类别 | JMX MBean 名称 | Prometheus 路径 | 关键指标 | 说明 |
|---|
| 实例状态 | destination=example | /metrics | running、dumpRunning、parseRunning | 查看 instance 是否运行 |
| 解析延迟 | ParseCounter | canal_parse_delay{destination="example"} | delay(毫秒) | Binlog 时间与解析时间差,核心延迟指标 |
| 消费位点 | MetaCounter | canal_mq_position{topic="..."} | position、filename | 当前消费的 Binlog 位置 |
| 消息吞吐 | SinkCounter | canal_sink_qps{destination="..."} | QPS、TPS | 每秒处理的消息数 |
| 连接状态 | ClientIdentity | canal_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 转换为 Entry | parseRowsEvent(...) | 核心转换逻辑 |
EventStore | 内存存储解析后的 Entry | put(Entry)、get(...) | 默认实现为 MemoryEventStoreWithBuffer |
| 工作流程 | 连接 MySQL → dump Binlog → 解析为 Entry → 存入 EventStore → 通知 EventSink | | 解析失败会重试 |
8.3 Entry 与 Protocol 结构解析
| 结构 | 说明 | 字段示例 | 注意事项 |
|---|
| Entry | 顶层结构,包含 header 和 storeValue | entryType: ROWDATA, storeValue: bytes | storeValue 为 PB 序列化数据 |
| Header | 元信息 | logfileName, position, timestamp, schemaName, tableName, eventType | 用于定位和路由 |
| RowChange | 行变更封装 | eventType: INSERT, rowDatasList | 需反序列化 storeValue |
| RowData | 单行数据 | beforeColumns, afterColumns | INSERT 无 before,DELETE 无 after |
| Column | 列信息 | name, value, isUpdate, isNull | isUpdate 仅 UPDATE 有效 |
| Protocol Buffer | 使用 Google Protobuf 定义 .proto 文件 | entry.proto, logproxy.proto | 编译生成 Java 类 |
8.4 Client-Server 通信协议实现
| 协议 | 核心类 | 说明 | 注意事项 |
|---|
| TCP 协议 | CanalServerWithNetty、SimpleCanalConnector | 基于 Netty 实现 TCP 长连接 | 自定义二进制协议 |
| gRPC 协议 | GrpcCanalServer、GrpcCanalConnection | 基于 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 作为中间桥梁 |
核心配置:
conf/application.yml:
canalAdapters:
- instance: example
groups:
- groupId: g1
outerAdapters:
- name: es
hosts: 192.168.1.200:9200
properties:
mode: rest
hosts 为 ES 地址
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(操作时间戳)解决冲突
- 双活架构复杂,需充分测试