Canal:MySQL Binlog 订阅、位点恢复与增量同步治理
一次商品索引切换后,值班群里出现了很矛盾的证据:MySQL 中价格已经改成 199,搜索结果仍是 229;Canal 进程和 Elasticsearch 写入程序都活着,监控里甚至还有持续增长的消费条数。工程师重启 Canal 后,旧价格突然被覆盖成新价格,又有少量商品出现重复刷新。表面上是“索引延迟”,真正的问题却横跨了四个状态:MySQL 的 binlog 已经走到哪里,Canal 解析到哪里,客户端确认到哪里,下游真正落盘到哪里。
这类事故提醒我们,CDC 不是“启动一个监听器”便结束。一个可接管业务的链路必须回答:历史数据如何形成基线,增量从哪个一致水位接上;DDL 发生后按哪一版表结构解释行事件;生产者失败时位点能否继续推进;重启或主库切换会重放多少数据;最终由谁消除重复并证明没有漏数。Canal 很适合做 MySQL binlog 的增量订阅与消费,但它不是自带全量校验、目标端事务和自动切流的通用迁移平台。
先看清 Canal 在链路中的位置
Canal 模拟 MySQL replica 的 dump 协议,从源库接收 binary log,再把字节流解析为结构化 Entry。官方仓库把它定义为“MySQL 数据库增量日志解析,提供增量数据订阅和消费”,采用 Apache License 2.0,仓库当前没有归档标记,发布页的稳定基线是 v1.1.8。实施时应锁定发布包或镜像标签,不要追随 latest,也不要把 README 中列出的 MySQL 兼容范围直接外推到 MySQL 8.4、9.x 或未验证的云数据库兼容层。
一套常见链路包含以下对象:
canal-server 负责连接 MySQL、解析 binlog、维护事件缓存,并向 TCP 客户端或 MQ 生产者交付事件。destination 是逻辑订阅名,例如 catalog;客户端、日志目录、位点和 HA 都围绕这个名字识别同一条订阅。instance 是 destination 的运行配置,conf/catalog/instance.properties 决定源库、复制身份、过滤规则和表结构历史。
TCP 客户端通过 getWithoutAck、ack、rollback 控制消费确认;MQ 模式则由 Canal 内置生产者把消息交给 Kafka、RocketMQ 等系统。client-adapter 是可选的下游适配层,不等于 Canal 核心,也不能让任意目标端自动获得幂等、校验和无损 DDL 演进。
图中最重要的不是箭头,而是三类不同的状态。源位点证明“读到哪里”,消费确认表示“可以释放到哪里”,目标端幂等记录证明“业务效果做到哪里”。任何一个状态丢失,都可能表现为重复、空洞或无法判断。Canal 的价值集中在中间的数据面;全量快照、目标端约束、差异修复和切流状态机仍需由迁移方案补齐。
把 MySQL 源端变成可解释的日志源
Canal 官方 QuickStart 要求自建 MySQL 开启 binlog、使用 ROW 格式,并为 Canal 分配不冲突的 replica server id。生产上还应优先使用 binlog_row_image=FULL:MINIMAL 虽能减少日志量,却可能让更新和删除事件缺少未变化列,下游若按完整行覆盖就会产生错误。
下面使用合成库 lab_catalog 和占位密码。修改持久化参数后需要按数据库变更流程重启 MySQL;仅执行 SET GLOBAL 不会改写已经建立连接的会话,也不会形成重启后的稳定配置。
# my.cnf
[mysqld]
server_id=101
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL
binlog_expire_logs_seconds=604800CREATE DATABASE IF NOT EXISTS lab_catalog
CHARACTER SET utf8mb4 COLLATE utf8mb4_0900_ai_ci;
CREATE USER IF NOT EXISTS 'canal_cdc'@'localhost'
IDENTIFIED BY '<CANAL_DB_PASSWORD>';
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.*
TO 'canal_cdc'@'localhost';
CREATE TABLE lab_catalog.product (
id BIGINT PRIMARY KEY,
name VARCHAR(120) NOT NULL,
price_cents INT NOT NULL,
version BIGINT NOT NULL
);
INSERT INTO lab_catalog.product VALUES (1001, 'desk-lamp', 22900, 1);
SHOW VARIABLES WHERE Variable_name IN
('log_bin', 'binlog_format', 'binlog_row_image', 'server_id');
SHOW GRANTS FOR 'canal_cdc'@'localhost';预期证据是 log_bin=ON、binlog_format=ROW、binlog_row_image=FULL,并且复制账号只有读取与复制相关权限。MySQL 8.4 的术语虽然已经改用 source/replica,静态权限名仍是 REPLICATION SLAVE:它允许复制客户端请求源端更新;REPLICATION CLIENT 则允许读取 binary log 状态和日志清单。两者都不是 REPLICATION_SLAVE_ADMIN 的同义替换,后者是管理复制通道的动态权限,Canal 不需要靠它读取 binlog。账号语法、TLS 和托管数据库裁剪后的权限应按目标实例及 MySQL 权限定义 核对。不要为了省事授予 ALL ON *.*;这会让一次配置泄漏直接升级为全库写权限事故。
binlog_expire_logs_seconds 不是越大越安全。保留窗口至少应覆盖“最长允许停机 + 最长人工发现时间 + 恢复和回放时间”,同时受磁盘容量约束。保留七天只是实验值,不能复制成生产结论。若 Canal 落后到所需文件已被清理,重启不会凭空找回历史,只能从备份、归档日志或新全量基线重新接续。
主库可能切换时,GTID 要从 MySQL 拓扑开始启用,不能只把 Canal 的布尔值改成 true。所有候选主都应启用 gtid_mode=ON、enforce_gtid_consistency=ON,复制节点还要启用 log_replica_updates=ON,否则在副本执行过的事务未必进入其 binary log。MySQL 要求先满足 GTID 一致性再进入 ON;已有拓扑应按官方在线迁移状态机逐级切换,不能直接把下面三行覆盖到生产并重启。
[mysqld]
gtid_mode=ON
enforce_gtid_consistency=ON
log_replica_updates=ONSHOW VARIABLES WHERE Variable_name IN
('gtid_mode', 'enforce_gtid_consistency', 'log_replica_updates');
SELECT @@GLOBAL.gtid_executed, @@GLOBAL.gtid_purged;三项变量应分别为 ON、ON、ON,并保存 gtid_executed 与 Canal 已确认水位用于切换比较。gtid_purged 增长本身不表示已丢数据,但它说明哪些事务已不在本机 binary log 中;仍未被 Canal 消费的 GTID 一旦进入 purged 集合,就不能靠自动定位补回。
用发布包或容器启动一个独立 destination
发布包更适合第一次学习,因为目录、日志和状态文件都可直接观察。到 Canal Releases 选择锁定版本,下载 canal.deployer,不要把示例里的版本变量留给启动脚本动态解析。
CANAL_VERSION=1.1.8
curl -fL -o canal.deployer.tar.gz \
"https://github.com/alibaba/canal/releases/download/canal-${CANAL_VERSION}/canal.deployer-${CANAL_VERSION}.tar.gz"
mkdir -p ./canal-home
tar -xzf canal.deployer.tar.gz -C ./canal-home
cp -R ./canal-home/conf/example ./canal-home/conf/catalogconf/canal.properties 负责 server 级设置,conf/catalog/instance.properties 负责这条订阅。下面保持 TCP 模式,便于观察客户端确认协议。${CANAL_DB_PASSWORD} 应在部署时由秘密管理系统渲染到只对运行账号可读的临时配置中;示例中的占位符不会被 Canal 自动解析。
# conf/canal.properties
canal.port = 11111
canal.destinations = catalog
canal.auto.scan = true
canal.instance.global.spring.xml = classpath:spring/file-instance.xml
canal.instance.transaction.size = 1024
# conf/catalog/instance.properties
canal.instance.mysql.slaveId = 9101
canal.instance.master.address = localhost:3306
canal.instance.master.journal.name =
canal.instance.master.position =
canal.instance.master.timestamp =
canal.instance.gtidon = false
canal.instance.dbUsername = canal_cdc
canal.instance.dbPassword = <CANAL_DB_PASSWORD>
canal.instance.connectionCharset = UTF-8
canal.instance.filter.regex = lab_catalog\\.product
canal.instance.filter.black.regex =
canal.instance.tsdb.enable = true
canal.instance.get.ddl.isolation = true几个字段会直接改变故障形态。slaveId 必须与 MySQL 拓扑中的 server id 以及其他 Canal 实例不同,否则源库会踢掉重复复制连接。filter.regex 是经过 Java Properties 解码后再编译的 Java 正则,所以配置文件里要写两个反斜杠 lab_catalog\\.product,解码后的有效正则才是 lab_catalog\.product。journal.name、position 和 timestamp 留空表示由元数据或源端当前位置决定,手工填写是高风险重放操作。get.ddl.isolation=true 让 DDL 独立成批,能降低消费者并发执行 DML 时跨过结构变化的概率,但不会替下游执行兼容性评审。
当源端已经完成 GTID 改造,可把 file/position 配置切换为 GTID 模式。master.gtid 是显式起点,不是每次启动都要粘贴的“当前 GTID”;已有 destination 应优先恢复其持久化 meta,只有新建、隔离回放或经审批重定起点时才填写它。
canal.instance.gtidon = true
canal.instance.master.journal.name =
canal.instance.master.position =
canal.instance.master.timestamp =
canal.instance.master.gtid =正向切换实验先记下 Canal 最后确认的 GTID 集合 <CANAL_ACKED_GTID_SET>,在候选主执行下面的包含关系检查;结果为 1 才说明候选主执行历史覆盖了 Canal 已确认事务。还要比较旧主切换水位,确认候选主没有缺少已提交事务,并验证所需后续 GTID 尚未被清理。
SELECT GTID_SUBSET(
'<CANAL_ACKED_GTID_SET>', @@GLOBAL.gtid_executed
) AS candidate_contains_checkpoint;反向实验使用一台故意延迟、使包含关系返回 0 的候选副本。自动化控制器必须阻断切换,Canal 不应降级为“从候选主当前位点继续”;否则进程和 lag 都可能正常,缺失事务却永远不会出现。GTID 只把跨日志文件的事务身份稳定下来,producer 已成功而 meta 尚未刷新时仍可能重放,下游幂等没有因此消失。
不要靠肉眼猜反斜杠是否被吞掉。下面直接读取实际 instance.properties,同时验证目标名命中、把点替换成其他字符后不命中:
cd ./canal-home
jshell <<'EOF'
import java.nio.file.*;
import java.util.*;
import java.util.regex.*;
var p = new Properties();
try (var r = Files.newBufferedReader(Path.of("conf/catalog/instance.properties"))) {
p.load(r);
}
var regex = p.getProperty("canal.instance.filter.regex");
System.out.println(regex);
System.out.println(Pattern.matches(regex, "lab_catalog.product"));
System.out.println(Pattern.matches(regex, "lab_catalogXproduct"));
/exit
EOF预期三行核心证据依次是 lab_catalog\.product、true、false。若最后一行也是 true,说明点号仍被当成任意字符;若第二行是 false,则规则过窄或转义层数错误。这个检查证明的是配置解码与正则边界,启动后仍应以 instance 日志和一条合成 DML 证明订阅真正生效。
启动后先看两层日志,而不是只做端口探测:
cd ./canal-home
sh bin/startup.sh
tail -n 80 logs/canal/canal.log
tail -n 80 logs/catalog/catalog.logserver 日志出现监听端口,只证明进程启动;instance 日志出现连接和 parser 启动,才证明订阅开始。若报 Access denied,核对账号的 host 匹配和认证方式;若报 server id 冲突,修改 slaveId 后重新建立连接;若报找不到 binlog 文件,先停止盲目重启并确认保存位点与 SHOW BINARY LOGS 的交集。
容器方式适合共享测试环境,但必须显式挂载配置、日志和状态目录,并锁定镜像标签。官方提供 Docker QuickStart,而生产镜像还要经过组织自己的来源、签名、SBOM 和漏洞门禁。容器里的 localhost 指向容器自身;连接宿主机 MySQL 时应使用受控网络入口,不能把数据库端口直接暴露到公网。
从 Entry 协议读懂 DML、DDL 与事务边界
TCP 默认交付 protobuf 对象。一个 Message 包含多个 Entry;Entry Header 带有 schema、table、event length、execute time、logfile name 和 offset,StoreValue 可解析为 RowChange。DML 的 RowDatas 保存 before/after columns,DDL 则携带 isDdl、sql 和事件类型。官方 ClientAPI 给出了连接、订阅、获取与确认接口。
下面的 Java 片段省略项目构建文件,但保留了正确的确认顺序。真正项目应把 database.table + 主键 + 源位点 或业务版本写入幂等表,并与目标数据更新放在同一目标端事务中;只有事务提交成功才 ack。
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress("localhost", 11111),
"catalog", "", "");
connector.connect();
connector.subscribe("lab_catalog\\.product");
connector.rollback();
while (!Thread.currentThread().isInterrupted()) {
Message message = connector.getWithoutAck(200, 1, TimeUnit.SECONDS);
long batchId = message.getId();
if (batchId == -1 || message.getEntries().isEmpty()) {
continue;
}
try {
targetTransaction(() -> {
for (CanalEntry.Entry entry : message.getEntries()) {
applyIdempotently(entry); // 目标写入与去重记录同事务
}
});
connector.ack(batchId);
} catch (Exception failure) {
connector.rollback(batchId);
throw failure;
}
}ack 的语义是通知 Canal 这批可释放,不是证明 Elasticsearch、缓存和审计库同时成功。如果先 ack 再写目标,进程在两者之间退出会永久漏数;如果先写目标再 ack,进程在两者之间退出会重放,因此下游必须幂等。rollback(batchId) 让未确认批次再次可取,也不会撤销已经写入的目标数据。
事务边界也不能想当然。Canal 1.1.8 的键名是 canal.instance.transaction.size,默认值为 1024;单个事务超过这个事件数阈值时可能被拆成多批交付,官方 AdminGuide 明确说明此时无法保证一次获取就是完整事务。消费者若必须按源事务原子生效,应检测 TransactionBegin/End、缓冲完整事务并设置大小上限;超大事务进入隔离队列,而不是无限堆内存。把它误写成 transactionn.size 不会得到预期覆盖,排障时应从启动使用的最终配置反查实际值。
做一次正向实验,让位点、事件和目标结果对得上
先在 TCP 消费者运行时提交一组可辨识事务:
START TRANSACTION;
UPDATE lab_catalog.product
SET price_cents = 19900, version = 2
WHERE id = 1001;
INSERT INTO lab_catalog.product
VALUES (1002, 'monitor-arm', 45900, 1);
COMMIT;
ALTER TABLE lab_catalog.product
ADD COLUMN stock_state VARCHAR(16) NOT NULL DEFAULT 'available';
UPDATE lab_catalog.product
SET stock_state = 'backorder', version = 3
WHERE id = 1001;正确实现应观察到同一事务中的 update 和 insert,以及独立的 DDL 批次,随后是带新列结构的 update。可记录但不泄漏业务值的证据包括:Header 中递增的 logfileOffset、同一事务的边界、DDL 的 isDdl=true 与 SQL 摘要、目标幂等表中的 source position、目标表最终 version=3。不要只用“消息条数增加”判定成功,因为重复事件也会增加计数。
验证可以分三层:
-- 源端业务事实
SELECT id, price_cents, version, stock_state
FROM lab_catalog.product ORDER BY id;
-- 目标端业务事实与幂等水位(示意表)
SELECT id, price_cents, version, stock_state
FROM search_projection ORDER BY id;
SELECT destination, source_file, source_offset, apply_status
FROM cdc_apply_log
WHERE destination = 'catalog'
ORDER BY source_file DESC, source_offset DESC LIMIT 10;预期结果不是要求两个系统在每个瞬间完全相等,而是在定义的延迟预算内,源端 1001 的最终版本和目标一致,DDL 之后的新字段被消费者按已发布契约处理,幂等水位不倒退。若 DDL 消息已出现而后续 DML 解析报列数不匹配,优先检查 TSDB 状态与 DDL 支持,不要直接跳过异常行。
用反向实验暴露“活着但不正确”
最有价值的演练是故意在目标写入成功后、ack 之前终止消费者。重新启动后,同一批会再次出现。幂等实现应把第二次应用识别为重复,目标业务版本不增加;非幂等实现则可能重复发送通知、重复累加库存或反复刷新索引。
演练步骤可以这样设计:
在测试消费者加入故障开关:目标事务提交后打印 TARGET_COMMITTED,随后在调用 ack 前退出。写入一条 version=4 的更新,保存 batch id、源 file/offset 和目标幂等键。重启消费者,确认同一源位点再次到达。
检查 cdc_apply_log 只有一条成功记录,业务表保持版本 4,并在重复计数器上增加一次。
典型证据应同时包含事件类型、位点、表结构版本和目标端最终状态:
TARGET_COMMITTED batchId=<BATCH_ID> source=mysql-bin.<N>:<OFFSET>
PROCESS_EXIT before_ack=true
REDELIVERED batchId=<NEW_BATCH_ID> source=mysql-bin.<N>:<OFFSET>
IDEMPOTENT_SKIP key=lab_catalog.product:1001 version=4
ACKED batchId=<NEW_BATCH_ID>这里 batch id 可以变化,稳定身份应来自源位点、主键和业务版本。反过来,如果第二次消费造成 version 从 4 变成 5,说明下游把“事件到达次数”误当成业务更新次数。此时增加 Canal 重试次数只会扩大副作用,正确修复是建立幂等写入或版本比较。
MQ 模式还应做生产者失败演练:临时把测试 topic 改成无权限名称或阻断 broker 网络,观察 instance 日志、生产错误和源端 lag。恢复后检查是否重发,并用消息 key/源位点去重。不要在业务高峰用真实生产 topic 做破坏性演练。
位点与 TSDB 为什么必须一起恢复
Canal 恢复不只需要 binlog file/position。行事件只携带列序号和值,解析器必须知道该位点对应的表结构;如果今天表有五列,却拿今天的结构解释三天前只有四列的事件,值会错位。Canal 的 TableMeta TSDB 用 DDL 序列和快照重建某一历史位点的结构,官方 TableMetaTSDB 说明默认可使用 H2,也可配置 MySQL 存储。
canal.instance.tsdb.enable=true 开启能力,默认 H2 数据位于 destination 目录。它与 file meta 都不应留在容器临时层。只备份位点而漏掉 TSDB,重启后可能从正确日志位置开始,却按错误 schema 解析;只备份 TSDB 而丢掉位点,则可能从新位置漏过历史,或从旧位置大规模重放。
恢复检查应把四件事放在同一记录里:
destination 与 instance 配置版本;已确认的 binlog file/position 或 GTID 集合;对应的 TSDB 数据及其存储健康;
下游最后成功应用的源位点和幂等状态。
手工删除 meta.dat、H2 文件或 ZooKeeper 节点不是普通“清缓存”。它相当于改写恢复协议,必须先冻结写入、保全旧状态、确定新起点和重放影响。尤其不要看到 DDL 解析异常就删除 TSDB 后直接启动,这会把可诊断的局部问题变成边界不明的数据事故。
全量基线要与增量水位严丝合缝
Canal 核心读取增量日志,不会替目标库生成一致性全量快照。一个新目标若只启动 Canal,只能看到启动后的变化,历史未变化行永远不会出现。正确衔接需要独立快照工具,并记录与快照一致的 binlog 水位。
一种常见流程是:在一致性快照开始时取得 T0 水位;把快照导入目标的 staging 区;让 Canal 从不晚于 T0 的位置持续收集增量;全量导入结束后按主键和版本应用积压事件;完成分块校验后才允许目标接读。具体 dump 工具如何给出一致水位取决于 MySQL 版本、事务表比例和锁策略,不能把 SHOW MASTER STATUS 与任意时刻导出的文件随意拼接。
如果源表没有主键,全量行与增量更新难以稳定对齐;如果快照过程中发生主键更新,下游要能把 before/after 主键解释为删除旧键加写入新键;如果目标导入覆盖了更新后的新值,就会产生“全量反盖增量”。因此 staging 合并应比较源版本、业务更新时间或确定的事件序,而不能简单执行最后到达者覆盖。
退出也要有明确动作:确认不再需要回切后,停止源端写入或记录最终水位,等待 Canal 和下游追平,保存最终差异报告,撤销复制账号,停止 destination,按保留策略清理 topic、状态文件和 staging 表。过早删除 binlog 或幂等记录,会让回切窗口失去证据。
MQ 输出把确认责任移到了另一条链路
当多个团队或多语言消费者共享事件时,可以把 canal.serverMode 改为 Kafka 等 MQ。官方 Canal Kafka/RocketMQ QuickStart 列出 tcp、kafka、RocketMQ、rabbitmq、pulsarmq 等模式,并解释了 flat JSON、批次、分区和 producer 参数。
canal.serverMode = kafka
canal.mq.topic = cdc.catalog
canal.mq.flatMessage = true
canal.mq.canalBatchSize = 50
canal.mq.canalGetTimeout = 100
canal.mq.partitionsNum = 6
canal.mq.partitionHash = lab_catalog.product:$pk$
kafka.bootstrap.servers = localhost:9092
kafka.acks = all
kafka.retries = 5
kafka.max.request.size = 1048576Canal 1.1.8 的发布包把 Canal 自身的 topic、批次和分区策略放在 canal.mq.*,把 Kafka Producer 原生参数放在 kafka.*。因此 broker、确认、重试和单请求上限必须分别写成 kafka.bootstrap.servers、kafka.acks、kafka.retries、kafka.max.request.size;旧写法 canal.mq.servers、canal.mq.acks、canal.mq.retries、canal.mq.maxRequestSize 不能作为 1.1.8 的有效 Kafka 配置。可用下面的静态证据核对最终文件,四个旧键的第二条命令应无输出:
grep -E '^(kafka\.(bootstrap\.servers|acks|retries|max\.request\.size))\s*=' conf/canal.properties
grep -E '^canal\.mq\.(servers|acks|retries|maxRequestSize)\s*=' conf/canal.propertiesflatMessage=true 便于普通 JSON 消费者接入,但会放大消息体并固化字段契约;false 使用 protobuf,需要相匹配的反序列化器。kafka.acks=all 只增强 Kafka 生产确认,不等于端到端 exactly-once。topic 的 min.insync.replicas、副本数、生产者重试、消费者提交和目标端幂等仍需共同设计。
分区策略决定可获得的顺序。按主键 hash 可以保持同一行的更新顺序并提高并行度,却不能保证跨行、跨表事务顺序;单分区能扩大顺序域,但吞吐和恢复速度受限。扩分区还可能改变 key 到 partition 的映射,消费者若把“同 key 永远同分区”当成跨扩容不变量,就会误判顺序。
消息过大时,先定位是大事务、LOB、flat JSON 膨胀还是批次过大。单纯同时提高 broker、producer 和 consumer 的消息上限,会把一次异常行扩散为内存峰值和网络拥塞。更稳妥的策略是限制单批、对大字段建立外部对象引用,或把异常事件送入有访问控制的隔离流程。
重启、主库切换与 HA 的真实边界
单机 file 模式能从持久化 meta 恢复,但进程、磁盘和主机仍是同一故障域。Canal 的 server HA 可借助 ZooKeeper:两台 server 使用相同 destination 名,配置 canal.zkServers 和 classpath:spring/default-instance.xml,每台使用不同 slaveId;同一 destination 正常只有一个 active parser。官方 AdminGuide 的 HA 配置 是核对这些字段的入口。
# 两个 Canal 节点保持相同 destination 和源配置
canal.destinations = catalog
canal.zkServers = localhost:2181
canal.zookeeper.flush.period = 1000
canal.instance.global.spring.xml = classpath:spring/default-instance.xml
canal.instance.tsdb.spring.xml = classpath:spring/tsdb/mysql-tsdb.xml
# 两个节点的 instance.properties 指向同一个受保护的 TSDB
canal.instance.tsdb.enable = true
canal.instance.tsdb.url = jdbc:mysql://mysql-tsdb.internal:3306/canal_tsdb
canal.instance.tsdb.dbUsername = canal_tsdb
canal.instance.tsdb.dbPassword = <CANAL_TSDB_PASSWORD>
# 节点 A 的复制身份
canal.instance.mysql.slaveId = 9101
# 节点 B 使用 9102ZooKeeper 解决的是 Canal server 的 active 选举以及解析/消费游标共享,不会自动共享 TableMeta TSDB。两台可接管节点若各自使用本地 H2,备用节点可能拿旧 schema 解释新位点。热备应像上例一样切换到官方提供的 MySQL TSDB 实现,让同一 destination 的节点读取同一份 DDL WAL 与 checkpoint,并对该库做备份、恢复和可用性监控;密码仍应由秘密系统渲染到仅运行账号可读的配置。等价的可恢复方案是单 active 配合受 fencing 保护的冷备,把 H2 与位点作为一个不可分割的持久卷快照恢复,并在接管前验证快照和所需 binlog 同时存在。把两份独立本地 H2 文件各自称为“HA”不满足恢复条件。
即使 TSDB 已共享,源 MySQL 的 binlog 过期、节点配置漂移、下游写入后未确认、MQ 已收到但目标未落盘仍不会被 ZooKeeper 自动解决。canal.zookeeper.flush.period 也意味着故障瞬间可能从较早游标恢复,产生有限重放;这需要下游幂等吸收。
MySQL 主库切换是另一套 HA。file/position 只在同一日志序列内有意义,新主库上的文件名相同也不代表事件相同。Canal 的非 GTID standby 模式还需要显式配置备用地址、开启检测和心跳切换;官方配置只支持一个 standby,fallbackIntervalInSeconds 会在新日志序列中向前回找,因此设计目标是偏向“不漏”,代价是可能重复。
canal.instance.standby.address = mysql-standby.internal:3306
canal.instance.standby.journal.name =
canal.instance.standby.position =
canal.instance.standby.timestamp =
canal.instance.detecting.enable = true
canal.instance.detecting.sql = select 1
canal.instance.detecting.retry.threshold = 3
canal.instance.detecting.heartbeatHaEnable = true
canal.instance.fallbackIntervalInSeconds = 60select 1 只能证明端点可查询,不能证明候选主已经追平;写心跳表能留下更强的日志连续性证据,却要求额外写权限和心跳表治理。更可靠的生产切换应由数据库控制面先完成提升与 fencing,再让 Canal 重连。可用 GTID 时优先按前述包含关系验证 canal.instance.gtidon 与候选主事务集;不用 GTID 时必须接受回退查找和重复窗口。两种模式都要演练“新主缺少旧主最后事务”和“旧主重新可达形成双写”,而不是让 Canal 的连接成功替数据库拓扑背书。
HA 验证不能止于杀进程后新节点启动。一次能暴露 TSDB 断层的验收应围绕 DDL 前后切换:先由节点 A 解析 ADD COLUMN stock_state 及一条使用新列的 DML,记录该 DDL 源位点、ZooKeeper 游标、共享 TSDB 快照/WAL 水位和目标版本;停止 A 并确认 B 成为 active 后,再执行 ADD COLUMN source_tag 与第二条 DML。最终应同时证明 B 的起始位点没有跳过 A 的已确认水位、两条 DDL 都有可追溯源位点、切换前后的 DML 均按对应列结构解析、日志没有 column size is not match,并且目标端通过幂等应用收敛到两个新字段的最终版本。
反向验收可在隔离环境让 B 使用一份停留在第一条 DDL 之前的本地 H2 副本再触发接管;预期它不能通过上述结构连续性门槛,应出现 schema 缺失、列数不匹配或被治理层阻断,而不是把“进程已监听端口”记成成功。验收记录至少保留故障前最后源位点、ZooKeeper 游标、A/B 使用的 TSDB 身份、新节点开始位点、两次 DDL 位点、重复事件数、目标最终版本和恢复时长。若新节点活了但源位点直接跳到当前,可能已经漏数;若从很早位置开始且幂等表容量不足,则恢复会把目标压垮。
DDL 与 adapter 故障怎样传播到下游
Canal 能解析和输出 DDL,也能用 TSDB 回放历史表结构,但“看见 DDL”不等于“目标自动安全变更”。例如 MySQL 新增 NOT NULL 列带默认值,搜索索引可能只需忽略;同步到另一个关系库时却要先扩展目标结构。重命名列、修改字符集、改变主键或枚举语义,更需要独立的兼容窗口。
消费者应把 DDL 当成受治理的控制事件:保存 SQL 摘要、源位点、schema 版本和处理结论;未知 DDL 默认暂停相关表而不是跳过;先让目标接受新旧两版数据,再发布源端变更;确认旧字段不再出现后才收缩。canal.instance.get.ddl.isolation=true 只是降低批内乱序,不会判断变更是否向后兼容。
client-adapter 可加快 RDB、Elasticsearch 等常见投影的接入,但它仍受目标版本、映射、SQL 模型和幂等语义约束。需要跨库事务、自动 schema registry、全量校验、复杂变换和大量连接器治理时,应比较 Debezium/Kafka Connect、Flink CDC 或托管迁移服务,而不是继续给 adapter 叠加脚本直到变成无人能升级的平台。
选择 Canal 更合理的信号是:源端以 MySQL 为主;团队接受自己实现或治理消费者;需要 TCP 拉取确认或直接投递国内常见 MQ;能够为全量、校验和切流另建流程。若目标是跨多种数据库、统一连接器生命周期、强 schema 契约和大规模任务编排,Canal 只应作为候选数据源组件,而不是平台总控。
权限、敏感数据与网络要按数据面等级管理
复制账号能读取大范围业务表,CDC 消息通常包含整行,敏感级别不低于源数据库。配置文件、线程转储、错误日志、MQ 死信、重放样本和差异报告都可能复制姓名、手机号、令牌或密文。治理不能只保护 MySQL 密码,还要保护“被密码读出来的数据”。
落地时至少执行这些约束:复制账号按表可见性与产品边界拆分;MySQL、ZooKeeper、MQ 和管理端口只在受控网络开放;配置中的数据库和 MQ 凭证由秘密系统注入并定期轮换;日志默认只记录 schema、table、主键散列、位点和错误码,不打印整行;死信与回放 topic 使用更严格的 ACL 和较短保留;Canal Admin 若启用,要接入组织认证、审计和最小管理权限,不能把默认账号暴露在共享网络。
字段脱敏的位置也要慎重。解析后立刻脱敏可以降低传播风险,但会使某些下游失去精确回放能力;让每个消费者自行脱敏则扩大原始数据暴露面。架构上应按用途拆 topic:受控原始流只给必要处理器,面向分析或开发的流在受审计转换器中完成字段裁剪、散列或令牌化。
容量、成本与观测要围绕积压而不是进程存活
Canal 的成本主要来自源库 binlog 生成与保留、解析 CPU、event store 内存、网络、MQ 存储、下游重放和状态数据库。全量阶段虽然不由 Canal 核心完成,却会与 CDC 争夺源库 IO;下游暂停时积压会转化为 binlog 保留压力或 MQ 存储压力。
容量估算可以从峰值变更字节率开始:峰值行变更数 × 平均行事件大小 × JSON/协议膨胀系数 × 保留时长 × 副本系数。还要单独估计超大事务,因为平均值无法解释 ring buffer 被单事务占满、MQ 单消息超限或消费者长时间等待事务结束。canal.instance.memory.buffer.size 要求为 2 的幂,扩大它只是增加瞬时缓冲,不会提升持续小于输入的下游吞吐。
监控至少形成以下关联:源端当前 binlog/GTID 与 Canal 解析位点的差;Canal event store 使用量;TCP 未确认批次或 MQ producer 失败;消息发布延迟;消费者 lag;目标应用延迟与失败数;重复命中数;DDL 隔离队列;binlog 剩余保留窗口。只看 JVM 存活和每秒解析数,会错过“解析很快、目标一直失败”的事故。
告警应带可行动阈值。例如“按当前写入速率,剩余 binlog 磁盘只能支撑两小时”比“磁盘 80%”更能指导决策;“目标应用水位落后源提交水位十五分钟且持续扩大”比“Canal 有延迟”更能区分短峰值与系统性欠容量。
把接入、变更、回滚和退出变成团队协议
每条 destination 都应有 owner、源表清单、数据分级、输出契约、位点存储、下游列表、延迟预算、binlog 保留预算和退出日期。配置进入版本库时只能保存无秘密模板,并通过评审检查过滤正则、server id、DDL 策略、分区键和消息大小。运行时生成的 meta、TSDB 和凭证不能混进 Git。
升级 Canal 时,先在影子 destination 从相同源读取,输出到隔离 topic;比较旧新两路的主键、事件类型、源位点和 schema 解释;覆盖 DML、DDL、大事务、重启、主切换和 producer 失败;确认消费者兼容后再切换。不要让两个实例写同一非幂等目标来“并行验证”。
回滚必须区分程序回滚与数据回滚。程序可以退回旧版本,已错误应用到目标的数据需要按源位点和差异记录修复;如果新版本改变了 JSON 结构或 topic 路由,旧消费者未必能直接接管。保留旧 topic、旧配置、最终水位和转换版本,直到回切窗口关闭。
最终可以用五个问题判断这条链是否成熟:能否从快照水位重建目标;能否证明重启只重放不漏数;能否在 DDL 前后按正确 schema 解码;能否在 producer 或目标失败时阻止错误位点推进;能否撤销账号、清理状态并仍保留审计证据。五个问题都能回答,Canal 才从“能抓 binlog”变成可治理的增量数据能力。
