Debezium MySQL CDC:从一致性快照到 GTID 故障切换
一次数据库主节点切换后,CDC 任务很快恢复成 RUNNING,业务 topic 也重新有了消息。半小时后对账却发现:一批订单更新重复出现,另有一张低频表缺了切换窗口内的三次变更。团队看到新主库的 SHOW BINARY LOG STATUS 正常,便把问题归为 Kafka 重复消费;继续追查才发现,连接器保存的是旧节点的 binlog 文件位置,新节点虽然数据追平,却没有同名文件坐标;用于恢复的 GTID 集合又被过滤配置排除了一部分。更隐蔽的是,低频表长期没有捕获事件,连接器 offset 没有随整个 MySQL 实例的日志前进,旧 binlog 已被清理。
Debezium MySQL connector 不是定时比较表内容。它先用一致性快照建立基线,再读取 MySQL binlog 中按提交顺序记录的行变更,解析当时有效的表结构,生成带 key、before/after、操作类型和源坐标的事件。Kafka Connect 保存“读到哪里”,Debezium 的 schema history 保存“那个位置的表长什么样”,MySQL 的 binlog 或 GTID 则提供可继续读取的事实。三者中任何一个被单独重置,都可能得到看似运行、实际重复或缺失的数据流。
先读懂一条变更事件里的恢复坐标
Debezium 事件不是只有 before 和 after。以更新事件为例,key 通常来自主键;value 的 op 表示操作,c、u、d、r 分别常见于创建、更新、删除和快照读取;source 中包含数据库、表、binlog 文件、位置、行号、GTID、connector 版本、事件时间和 snapshot 标记。删除通常先发 delete 事件,再按配置发送同 key、null value 的 tombstone,便于压缩 topic 清除旧 key。
快照事件使用 op=r,它描述 T0 一致性视图中的现存行,不是数据库刚执行了一次 INSERT。下游若把 r 当新增审计事件,会在重做快照时制造业务重复;正确的物化视图通常按主键 upsert,把 r 与 c/u 都视为某个版本的当前值。审计型消费则需要保留源坐标和 snapshot 标记,明确区分“基线读出”和“实时发生”。
MySQL binlog file/position 是某个服务器上的物理坐标。故障切换到另一实例后,文件名和位置通常不可直接比较。GTID 将事务标识为 server_uuid:sequence 集合,使副本能证明自己包含哪些事务;它并不自动证明新节点包含连接器所需的所有历史,也不消除 errant transaction。连接器恢复时必须找到包含已处理 GTID 集合、且保留后续 binlog 的服务器。
Debezium 稳定版 MySQL connector 文档 详细列出事件结构、支持拓扑、snapshot 模式与配置属性。以下配置采用 3.6.0.Final 与 MySQL 8.4 系列语义:该 Debezium 版本按 Apache License 2.0 发布,connector 要求 Java 17+,构建于 Kafka Connect 4.3.0;官方测试矩阵列出 Kafka Connect 3.1+ 与 MySQL 8.0、8.4、9.0、9.1。它不表示任意版本组合都经过同等验证,制品清单仍要锁定 connector、Connect、broker、MySQL、驱动及 converter,并按 3.6 release notes 与 MySQL release notes 做升级回归。MySQL Community Server 的许可与 Debezium 不同,镜像再分发、商业插件和驱动也应单独审查。
让 MySQL 产生可用于 CDC 的日志
connector 至少需要 MySQL 开启 binlog,使用 ROW 格式,并给每个 replication client 唯一 database.server.id。binlog_row_image=FULL 让 update/delete 事件能获得完整的旧值语义,是最容易推理的基线;改成 MINIMAL 会减少日志体积,但未变化列或某些旧值不再完整,下游契约和 converter 必须先验证。
实验实例可以使用下面的 MySQL 配置:
[mysqld]
server-id=101
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL
gtid_mode=ON
enforce_gtid_consistency=ON
binlog_expire_logs_seconds=604800server-id 标识 MySQL 复制节点;Debezium 的 database.server.id 是 connector 作为 replication client 使用的唯一 ID,两者不要混淆。binlog_expire_logs_seconds 示例只表达“必须显式设计保留窗”,不能照抄为生产答案。保留时间至少覆盖允许的最长 connector 停机、快照持续时间、故障发现与人工恢复时间,再留安全余量;日志产生速率高时还要按磁盘容量反算。
创建合成数据库和最小账号:
CREATE DATABASE migration_lab;
CREATE TABLE migration_lab.orders (
id BIGINT PRIMARY KEY,
status VARCHAR(32) NOT NULL,
amount DECIMAL(12,2) NOT NULL,
updated_at TIMESTAMP(6) NOT NULL DEFAULT CURRENT_TIMESTAMP(6)
ON UPDATE CURRENT_TIMESTAMP(6)
);
CREATE USER 'cdc_reader'@'%' IDENTIFIED BY '<cdc-password>';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE,
REPLICATION CLIENT ON *.* TO 'cdc_reader'@'%';
-- 仅当目标环境无法使用全局读锁、快照会退回表锁时追加。
GRANT LOCK TABLES ON *.* TO 'cdc_reader'@'%';权限不是越多越稳。SELECT 用于快照,REPLICATION SLAVE 与 REPLICATION CLIENT 用于 binlog 和复制状态,RELOAD 用于全局读锁流程,SHOW DATABASES 用于元数据发现。官方文档要求在不允许全局读锁、需要表级锁的托管环境额外授予 LOCK TABLES,因此上面的第二条授权不是所有环境的固定基线。若只授予业务 schema 的 SELECT,还要验证 history 捕获与 include list 是否一致。完成初始快照后,能否收回 LOCK TABLES、RELOAD 取决于后续 blocking snapshot、history recovery 和故障重建流程;团队应建立流式读取角色与临时快照授权,而不是让长期账号永久持有不必要能力。
连接前用数据库自身证据确认配置:
SHOW VARIABLES WHERE Variable_name IN (
'log_bin','binlog_format','binlog_row_image',
'gtid_mode','enforce_gtid_consistency','binlog_expire_logs_seconds'
);
SHOW BINARY LOG STATUS;
SHOW GRANTS FOR 'cdc_reader'@'%';预期 log_bin=ON、binlog_format=ROW、gtid_mode=ON,账号权限与设计一致。MySQL 8.4 使用 SHOW BINARY LOG STATUS 检查当前二进制日志坐标;若输出没有日志文件或 Executed_Gtid_Set,先修 MySQL,不要靠重启 connector 碰运气。
把插件和依赖装进 Connect worker
Debezium 官方提供 connector archive、容器镜像以及其他运行形态,安装说明 给出了下载插件并解压到 Kafka Connect plugin.path 的流程。共享环境应把 Debezium MySQL connector、匹配依赖和校验信息构建为不可变镜像;不要在正在运行的 worker 容器里临时下载 JAR。
本地 Compose 可以用 MySQL、KRaft Kafka 和 Debezium Connect 三个服务建立联调环境。下面是单 broker、单 controller 的可运行实验拓扑;kafka:9092 供容器网络访问,localhost:29092 供宿主机命令访问。它不具备 broker 容错能力,不能直接复制为生产集群:
services:
kafka:
image: apache/kafka:4.3.1
hostname: kafka
ports: ["29092:29092"]
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: CONTROLLER://:9093,PLAINTEXT://:9092,HOST://:29092
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,HOST://localhost:29092
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,HOST:PLAINTEXT
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
healthcheck:
test: ["CMD-SHELL", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null 2>&1"]
interval: 5s
timeout: 5s
retries: 20
mysql:
image: mysql:8.4
command:
- --server-id=101
- --log-bin=mysql-bin
- --binlog-format=ROW
- --binlog-row-image=FULL
- --gtid-mode=ON
- --enforce-gtid-consistency=ON
environment:
MYSQL_ROOT_PASSWORD: "lab-root-password"
ports: ["3306:3306"]
healthcheck:
test: ["CMD-SHELL", "mysqladmin ping -h localhost -plab-root-password --silent"]
interval: 5s
timeout: 5s
retries: 30
connect:
image: quay.io/debezium/connect:3.6
environment:
BOOTSTRAP_SERVERS: kafka:9092
GROUP_ID: connect-migration-lab
CONFIG_STORAGE_TOPIC: _connect-migration-configs
OFFSET_STORAGE_TOPIC: _connect-migration-offsets
STATUS_STORAGE_TOPIC: _connect-migration-status
CONFIG_STORAGE_REPLICATION_FACTOR: 1
OFFSET_STORAGE_REPLICATION_FACTOR: 1
STATUS_STORAGE_REPLICATION_FACTOR: 1
ports: ["8083:8083"]
depends_on:
kafka:
condition: service_healthy
mysql:
condition: service_healthy把片段保存为独占实验目录中的 compose.yaml 后执行 docker compose up -d,再用 docker compose ps 与 curl -s http://localhost:8083/connector-plugins 检查三个服务和插件发现。文中的建库、建表与授权 SQL 仍需在 MySQL 就绪后执行,connector 密码应与实验账号一致。本机解压插件更适合看清 plugin.path 和类加载问题;Compose 适合复现网络、重启和多服务依赖;共享 distributed Connect 适合持续同步与 worker 故障迁移;Debezium Server/Engine 适合不需要 Kafka Connect 控制面的嵌入或直达 sink 场景,但其 offset、schema history、高可用和运维模型不同,不能把同一份恢复步骤直接套用。
一份可审查的 connector 配置
先在每台 worker 的 /connector-plugins 中确认 io.debezium.connector.mysql.MySqlConnector 可见,再通过 REST 提交:
{
"name": "mysql-orders-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"tasks.max": "1",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "cdc_reader",
"database.password": "<cdc-password>",
"database.server.id": "5401",
"topic.prefix": "labmysql",
"database.include.list": "migration_lab",
"table.include.list": "migration_lab.orders",
"snapshot.mode": "initial",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-history.mysql-orders",
"include.schema.changes": "true",
"heartbeat.interval.ms": "10000",
"tombstones.on.delete": "true"
}
}把 JSON 保存为 mysql-orders-cdc.json 后执行:
curl -s -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' \
--data @mysql-orders-cdc.json
curl -s http://localhost:8083/connectors/mysql-orders-cdc/statusdatabase.hostname/port/user/password 决定源连接;密码应在共享环境改成 ConfigProvider 引用。database.server.id 必须在同一 MySQL 复制拓扑中唯一,否则服务器会踢掉 ID 冲突的 client。topic.prefix 是事件命名空间和恢复身份的一部分,运行后修改会把事件写到新 topic,并破坏消费者、schema change topic 与历史认知;不要把它当展示名称。
database.include.list 与 table.include.list 是按完整名称匹配的 anchored regular expressions,点号等正则元字符必须按意图转义;误把子串当匹配规则会静默扩大或缩小捕获集合。过滤数据不总等于过滤 schema:schema.history.internal.store.only.captured.tables.ddl 默认是 false,初始快照会保存数据库内非系统表的结构,以便后来扩展捕获。若改成 true 来缩小 history,后来新增表就必须通过增量或 schema snapshot 流程补足结构,不能只改 include list 然后期待自动理解过去 DDL。
MySQL connector 对一条 binlog 通常只有一个 task,因此 tasks.max=1 是事实边界。增加到更大不会把事务日志并行拆分;容量瓶颈需要从快照并行能力、单 task 解码、Kafka producer、网络、消息尺寸和下游消费分别处理。
初始快照怎样和增量日志接上
snapshot.mode=initial 的关键不是“先 SELECT 再读日志”,而是在一致性事务与必要锁语义下记录 binlog 水位、读取 schema 和表数据,再从该水位接续 streaming。快照期间发生的新写入不会靠运气补齐:只要对应 binlog 被保留,快照完成后会从记录位置继续读取。快照行发 op=r,并共享快照源坐标;最后成功状态写入 Connect offset。
官方文档列出的常用模式有不同风险:
| 模式 | 行为 | 适合的决策 |
|---|---|---|
initial | 没有已完成快照状态时取 schema 与数据,然后流式读取 | 新建持续 CDC 的常规入口 |
initial_only | 只做快照,完成后不继续流式读取 | 一次性基线导出验证 |
no_data | 取 schema,不发历史行,随后读取新变更 | 已有独立基线且能证明切点一致 |
when_needed | 没有 offset 或日志位置不可用时自动快照 | 接受自动重建及其源库负载的场景 |
recovery | 用当前源 schema 重建丢失或损坏的 history | connector 停止后完全没有执行任何 DDL 的受控恢复 |
always 每次启动都重新快照,容易造成大量 r 事件和源库负载,不适合作为“防止漏数”的粗暴保险。旧配置名 schema_only_recovery 已弃用,应迁到 recovery;never 已从 3.6 移除,过去用它表达“不快照”的配置应迁到 no_data 并回归启动行为。recovery 的前提是 connector 停止后源库完全没有发生任何 DDL;即使只新增一个看似兼容的可空列,也已经破坏官方恢复前提。执行前应冻结 DDL、停止 connector、记录 offset 与当前 schema、重建 history,验证恢复后再解除冻结。无法证明停机窗口内零 DDL 时,应恢复原 history 备份,或以全新 connector、history 与业务 topic 重新快照和对账。删除 offset 再用 initial 也不是普通重启,而是重建整条基线。
大表快照要同时控制 snapshot.fetch.size、snapshot.max.threads、max.batch.size、max.queue.size、Kafka producer 限制与 Connect 堆内存。MySQL connector 的 snapshot.fetch.size 默认不设置,以流式读取结果;官方特别提醒,显式设置可能让驱动把完整结果集取入内存,因此不能把“增大 fetch”当作通用提速手段。snapshot.max.threads 增加表级并行会提高源库连接、IO、网络和 worker 堆压力;max.queue.size 必须大于 max.batch.size,扩大队列只是在内存中吸收短时波动,不能提高单 task 的持续处理上限。快照期间监控已完成表数、剩余表、持续时间、锁等待、源库 IO、Kafka 发送失败和 binlog 增长;只盯 connector RUNNING 看不到一场正在拖垮主库的全表扫描。
正向实验:证明基线与增量连续
启动 connector 后,在合成表写入两行,再更新和删除:
INSERT INTO migration_lab.orders(id,status,amount)
VALUES (1001,'CREATED',88.50),(1002,'CREATED',19.90);
UPDATE migration_lab.orders SET status='PAID' WHERE id=1001;
DELETE FROM migration_lab.orders WHERE id=1002;若两行在 snapshot 开始前已存在,topic labmysql.migration_lab.orders 中应先看到对应 key 的 op=r;快照后执行的 update 应为 op=u,包含 before.status=CREATED 和 after.status=PAID;delete 应为 op=d,随后可见相同 key 的 tombstone。若写入发生在快照并发窗口,事件组合可能是快照读加后续更新,或直接由 binlog 事件呈现,但按 key 物化后的最终值应与源表一致。
验证不能只数消息。抽取事件 key、op、source.file、source.pos、source.gtid、source.snapshot,按主键重放到一个临时映射,最终应只有 1001=PAID,1002 被删除。再重启承载 task 的 worker,写入 1003,预期从既有 offset 接续而不重新发整表 r。若出现少量重复,要确认它们是否位于上次 offset 提交与崩溃之间,并验证下游幂等是否收敛。
ts_ms 也有两层:事件 envelope 的处理时刻与 source.ts_ms 的数据库变更时刻可用于估算捕获延迟,但时钟漂移会让差值失真。监控延迟时同时校验 MySQL、worker 和 broker 的时间同步,并结合 binlog position/GTID 是否前进。
反向实验:让 schema history 丢失显形
在实验环境完成一次快照和一次 DDL:
ALTER TABLE migration_lab.orders
ADD COLUMN channel VARCHAR(16) NOT NULL DEFAULT 'web';
UPDATE migration_lab.orders SET channel='store' WHERE id=1001;确认事件能按新 schema 解码后,停止 connector。仅在可销毁实验集群删除 schema-history.mysql-orders,保留 Connect offset,再启动 connector。预期不是“自动从当前表结构继续”,而是 task 因无法恢复数据库 schema history 而失败,日志出现 history topic 不存在、无法恢复 schema 或对应位置 schema 缺失的证据。
这个反例说明 offset 与 history 是同一恢复协议的两半:offset 告诉 connector 从哪个旧位置读取,history 告诉它在那个位置怎样解释行。只清其中一个会产生不自洽状态。修复有三条路:从可靠备份恢复原 history topic;证明 connector 停止后完全没有任何 DDL,再按官方 recovery 流程重建;或者创建全新 connector 身份、history topic 与业务 topic,重新快照并对账切换。共享环境禁止边试边删。
内部 schema history topic 必须单分区,以保留全局 DDL 顺序;它只供 connector 使用,业务消费者应读取 topic.prefix 对应的 schema change topic,而不是解析内部格式。它与 Connect 的 config/offset/status topic 不同:history 恢复需要从早期 DDL 顺序重建各位置的表结构,不能按 key 压缩成“最新值”。因此显式使用 cleanup.policy=delete 并把时间、字节保留都设为无限,禁止 compact 或 compact,delete;若组织策略必须设置有限保留,应把可恢复备份与定期受控 recovery 纳入流程,而不是等待旧段被删后才发现。独占实验可在创建 connector 前执行:
docker compose exec kafka /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka:9092 --create \
--topic schema-history.mysql-orders --partitions 1 --replication-factor 1 \
--config cleanup.policy=delete \
--config retention.ms=-1 \
--config retention.bytes=-1生产环境还要设置足够副本、严格 ACL,并把 history 与 Connect offset 一起纳入一致灾备。retention.ms=-1 和 retention.bytes=-1 表示不按时间或总字节主动删除,但仍需监控容量和 broker/topic 级覆盖配置;若必须设置有限 retention,它必须长于最坏停机、审计与回滚窗口,且要有恢复副本。把 history topic 名复用于两个 connector,会让不同源的 DDL 互相污染。
heartbeat 解决的是可见性和低流量位点推进
heartbeat.interval.ms 让 connector 在没有普通变更事件时周期性向 heartbeat topic 发消息。它可以证明 task 仍能向 Kafka 发送,并让消费者观察最后活跃时刻;值为 0 时禁用。heartbeat 不是数据库健康检查:MySQL 连接卡死、producer 重试或 worker 停顿时它也可能停止,需要结合 connector status、JMX 指标和数据库连接观察。
更难的场景是:MySQL 实例有大量变更,但都发生在被过滤的数据库或表,connector 看到 binlog 向前却不产生捕获事件,offset 可能无法按需要持续提交;低频捕获库停机久后,旧 binlog 被清理。Debezium 提供 heartbeat.action.query,可周期性在源库执行一条会进入 binlog 的动作,从而产生可归属于捕获范围的水位。常见设计是在捕获数据库建专用 heartbeat 表并 upsert 时间标识。
CREATE TABLE migration_lab.debezium_heartbeat (
id TINYINT PRIMARY KEY,
touched_at TIMESTAMP(6) NOT NULL
);连接器配置可增加一条完整的合成语句:
{
"heartbeat.interval.ms": "10000",
"heartbeat.action.query": "INSERT INTO migration_lab.debezium_heartbeat(id,touched_at) VALUES (1,NOW(6)) ON DUPLICATE KEY UPDATE touched_at=NOW(6)",
"table.include.list": "migration_lab\\.(orders|debezium_heartbeat)"
}还要执行 GRANT INSERT, UPDATE ON migration_lab.debezium_heartbeat TO 'cdc_reader'@'%';。示例同步扩大 table.include.list,因为 action 表必须进入捕获集合;若下游不需要它的业务事件,可用经过验证的 SMT 丢弃输出,但不能在源过滤阶段让连接器看不见该表。正向证据是 heartbeat 表持续更新、对应 GTID/binlog 与 Connect offset 一起前进;反向实验撤销这两个写权限后,预期 connector 日志出现 action query 授权失败,普通 heartbeat topic 即使仍偶尔有消息,也不能证明源位点被推进。实验结束先移除 connector 配置中的 action query 并恢复原 table.include.list,再撤销写权限并删除 heartbeat 表。不允许 CDC 账号写源库时,应通过扩大 binlog 保留、独立受控作业或架构调整解决,不能假装普通 heartbeat topic 一定推动 MySQL 源坐标。
heartbeat 周期越短,水位观测越灵敏,但会增加源库写、binlog、Kafka 消息和指标量。告警阈值应大于周期加上可接受抖动,并区分“heartbeat topic 无消息”“源 GTID 不前进”“业务表无事件”三种含义。
GTID failover 要证明集合包含关系
使用 file/position 时,connector 通常必须回到原服务器或具有可映射日志的拓扑。开启 GTID 后,连接器能在兼容拓扑中寻找已处理事务之后的位置,但切换成功需要同时满足:
新服务器的已执行 GTID 集合包含 connector 已处理的 GTID。后续需要的事务仍有对应 binlog 可读。server_uuid 来源集合符合 gtid.source.includes / gtid.source.excludes。
新节点 schema、数据与账号权限已经追平。connector 的 topic.prefix、offset 与 schema history 保持不变。
切换前保存 connector offset 和两边 @@GLOBAL.gtid_executed,用 MySQL GTID 集合函数判断包含关系,而不是比较字符串长度:
SELECT @@GLOBAL.server_uuid, @@GLOBAL.gtid_executed;
SELECT GTID_SUBSET('<connector-processed-gtid-set>', @@GLOBAL.gtid_executed)
AS processed_transactions_exist;预期返回 1 才说明新节点至少执行过 connector 已处理集合。它仍不证明所需 binlog 未过期,因此还要检查日志列表和副本追平状态。gtid.source.includes 适合多源拓扑中只接受可信 origin;配置错误会让连接器忽略必要事务来源或把本地 errant transaction 纳入恢复。include 与 exclude 不能靠模糊正则随意维护,变更必须有拓扑 owner 评审和样例 GTID 测试。
planned failover 的稳妥顺序是:让候选节点追平;暂停或停止 connector 并记录最终水位;证明 GTID 包含和日志保留;切换 database.hostname 或稳定代理端点;启动 connector;观察新事件的 server_id、GTID、file/pos 与重复窗口;做分块对账后再释放旧节点。非计划故障无法保证干净停止,因此默认接受至少一次重放,并让下游按主键与版本收敛。
代理或 DNS 能隐藏主机变化,却不能替代 GTID 证明。如果代理把 connector 导向落后的副本,TCP 连接成功只会更快地产生错误恢复。CDC 端点的健康检查应包含复制追平和日志可用性,不只检查 3306 端口。
至少一次是默认契约,EOS 有严格条件
Debezium 官方明确给出默认至少一次语义:正常情况下不漏变更,但故障点附近可能重复。原因是 source task 从 MySQL 读取、向 Kafka 发送和提交 Connect offset 不是天然的单一数据库事务;若消息已写 Kafka、offset 尚未提交时进程崩溃,恢复会从旧 offset 重读。下游必须把事件 key、源坐标或业务版本纳入幂等设计。
Kafka Connect 从支持 KIP-618 的版本开始为 source connector 提供 exactly-once 能力。Debezium EOS 配置说明 要求 Kafka Connect distributed 模式且版本至少为 3.3.0,所有 worker 设置:
exactly.once.source.support=enabledconnector 同时设置:
{
"exactly.once.support": "required",
"transaction.boundary": "poll"
}还需要 Kafka transaction 支持、正确的 transactional ID ACL、兼容的 connector 声明和相同 worker 配置。required 会在 connector 不支持时拒绝启动,比静默降级更适合受控环境。Debezium MySQL 在官方支持列表中,但这只约束 source 记录与 offset 写入 Kafka 的事务边界,不等于“源 MySQL 到任意目标数据库端到端绝对一次”。
业务消费者若未启用 read_committed,可能看见 aborted transaction;sink 若不具备事务或幂等写,重试仍会在目标端重复;业务副作用、外部 API 与多 topic 消费也各有边界。官方还提示 Kafka 事务协议与 exactly-once 正确性存在已知研究和开放问题。因此架构决策应写成可验证条件:在故障注入中是否观察到重复、消费者隔离级别是否正确、目标端是否幂等、事务超时与 ACL 是否可控。不要用一个 enabled 配置替代端到端证明。
实践中,主键 upsert 加可比较的业务版本通常比依赖全链路 EOS 更容易恢复。GTID 是事务身份集合,file/position 是单个 MySQL 实例上的物理坐标;把二者拼成字符串或元组,不能得到跨节点全序,也不能据此用“大于”覆盖“小于”。目标端需要全序更新时,应优先使用源表单调 version、受约束的业务序列,或在同一拓扑内经过证明的提交序;源坐标用于去重、审计和断点证明,跨节点并发或 errant transaction 则按业务冲突策略收敛。delete 也要保存墓碑版本,防止延迟 update 复活已删记录。没有主键的表应在 CDC 上线前补稳定键或明确只做追加审计。
Schema 演进、敏感数据与消息契约
DDL 会进入 binlog,Debezium 解析后更新内存 schema 并写 internal history,同时可向 schema change topic 发业务可见事件。schema change 消息格式在官方文档中仍标为 incubating,不能把其当前 JSON/Avro 形状固化为永久外部契约。新增可空列通常容易兼容;重命名列在事件层常表现为删旧加新;缩小精度、修改时区语义、字符集或 collation、改变主键都会影响 key 和下游状态。DDL 发布必须把 connector 解析能力、converter/schema registry 兼容策略、消费者升级顺序与回滚一起评审。
在线 DDL 工具可能创建影子表、触发器并交换表名。若 include list 排除了辅助表,connector 可能缺少完成演进所需的 DDL/history;官方 MySQL 文档专门提醒 gh-ost、pt-online-schema-change 辅助表捕获问题。团队应在合成库完整演练所用 DDL 工具,确认 schema history、业务 topic 和下游映射,再决定用 SMT 过滤不需要的辅助表事件,而不是在源端发现阶段直接失明。
CDC 复制的是行级事实,密码哈希、证件号、备注、软删除内容都可能进入 Kafka,并在 topic、重试、DLQ、日志和测试导出中形成多份副本。最优先在 column.include.list 缩小字段;需要掩码时使用受评审 converter/SMT,并验证 key、schema 和删除语义。连接器配置、REST 响应和异常 trace 不得记录明文密码,Kafka ACL 要分别保护业务 topic、schema change、heartbeat、schema history 和 Connect 内部 topic。
源账号按 connector 或安全域隔离,网络只允许 worker 到指定 MySQL,TLS 验证服务端身份;证书轮换先在并行连接验证,再滚动 worker。插件运行在 Connect JVM 内,能接触解析后的凭证和数据,供应链审核、制品签名、SBOM 和版本 owner 都是数据权限的一部分。
容量瓶颈通常先出现在快照或大事务
稳态 CDC 的容量由源端 binlog 产生速率、单 task 解析能力、消息膨胀倍数、Kafka 分区与 producer 吞吐、消费者处理能力共同决定。before/after、schema、headers 和 JSON 编码可能让消息远大于原行;大字段和宽表会触发 producer max.request.size、broker message.max.bytes 或消费者限制。只调大一端会把失败推到下一端。
大事务在提交后形成一段密集事件,可能占满 max.queue.size、推高延迟并造成长时间 poll。max.batch.size 较大可提高吞吐却增加内存和重放批次;poll.interval.ms 影响空闲轮询;Kafka producer 的压缩、batch 与 linger 改变网络和延迟。调参必须用真实行宽分布和合成大事务压测,观察 JVM heap、队列剩余容量、binlog lag、发送延迟和失败重试。
快照成本与全表行数、索引、并发写、buffer pool、网络和 Kafka 写入有关。把快照放在业务低峰仍可能因缓存污染拖慢线上查询。更安全的架构是从已追平、开启 binlog 的副本读取,但副本必须保证所需事务和日志,且快照一致性、GTID failover 与读权限都经过验证。托管数据库允许的锁与复制权限不同,不能默认副本方案总可用。
成本治理以资源量而非价格描述:每日变更行数乘平均事件尺寸得到 Kafka 入口量;再乘副本、保留与下游副本估算存储;快照产生一次全量放大;schema history、heartbeat 和内部 topic 量小但不可随意删除;观测若采集完整 payload 会产生高额存储与敏感数据风险。按 connector 建 lag、事件速率、错误率、快照进度和 topic 增长预算,超预算先找大事务、字段膨胀或失速消费者。
故障证据要能指向恢复动作
| 现象 | 判断证据 | 常见根因 | 安全恢复 |
|---|---|---|---|
| 启动即报 binlog 格式错误 | SHOW VARIABLES、task trace | 未开 ROW 或日志关闭 | 修 MySQL 配置后重启失败 task |
server id 冲突后反复断连 | MySQL error log、replication client | 两个 connector 使用同一 ID | 分配唯一 ID,核对是否产生重放 |
| snapshot 很久不结束 | snapshot 指标、锁等待、源 IO | 大表、慢网络、fetch/queue 不匹配 | 降低负载或改用验证过的副本,不删 offset |
| 重启提示 binlog 不可用 | offset 的 file/pos/GTID、日志列表 | 停机超过保留窗 | 从一致快照重建或恢复日志,禁止前跳 |
| DDL 后解析失败 | history topic、DDL、connector 日志 | history 缺失或 DDL 不支持 | 恢复 history 或隔离新链路重快照 |
| failover 后重复 | GTID 集合、切换前后 offset | 非干净停止或新节点坐标映射 | 下游幂等收敛,验证无缺口 |
| status 绿但 heartbeat 停止 | heartbeat、JMX、Kafka producer 指标 | 无捕获事件、发送阻塞或 task 卡顿 | 区分空闲与阻塞,再恢复连接 |
不要在第一时间清空 offset 和 history。先导出 connector 配置与 offset、保存 status trace、记录 MySQL GTID/binlog 列表、抓取 snapshot/streaming 指标,再判断哪层断裂。证据完整时通常可以选择“原位续跑、受控重放、history 恢复、全新链路重建”之一;证据被重启和删除覆盖后,只剩代价最大的重新快照与全量对账。
项目接入要围绕幂等和可回放设计
消费者不要只反序列化 after。先稳定解析 key、op、source、snapshot 与 delete/tombstone;建立 schema 兼容测试;把未知字段向前兼容;将原始事件与业务投影分层。物化数据库写入使用主键 upsert 和版本条件,处理成功后再提交消费 offset。跨表业务一致性若依赖 MySQL 事务,可启用并消费 transaction metadata,按事务边界聚合,但仍要设计超时、超大事务和部分 topic 不可用时的策略。
connector 配置进入 Git 时只保存 secret 引用。CI 可调用 /connector-plugins/{class}/config/validate 校验字段,再检查 topic.prefix、history topic、include list、server ID、snapshot mode、heartbeat 和 EOS 组合;发布系统通过 REST 声明式更新。数据库、数据平台、业务消费者和安全团队分别拥有源权限、运行时、事件契约和数据分级,任何一方单独改 include list 或 DDL 都可能破坏链路。
上线门槛应是一次可回放的故障演练:快照期间持续写入,验证最终表状态;强杀 worker,确认重复可由幂等收敛;让 binlog 保留不足的实验实例产生明确失败;执行计划 DDL;切到包含全部已处理 GTID 的副本;撤销权限并恢复。每次演练保留源坐标、topic offset、目标校验和错误 trace,而不是只截图 RUNNING。
清理、回滚和长期退出
配置错误但尚未产生业务事件时,可以停止并删除 connector,删除合成业务 topic、schema history topic 和实验账号。已经进入共享 topic 后,直接重建同名 connector 会混合两条时间线;应先暂停消费者,决定是按 key 重放收敛、从新 topic 重新建表,还是恢复旧配置与 offset。
实验清理示例:
curl -s -X PUT http://localhost:8083/connectors/mysql-orders-cdc/stop
curl -s http://localhost:8083/connectors/mysql-orders-cdc/offsets
curl -s -X DELETE http://localhost:8083/connectors/mysql-orders-cdc
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--delete --topic labmysql.migration_lab.orders
bin/kafka-topics.sh --bootstrap-server localhost:9092 \
--delete --topic schema-history.mysql-orders最后在 MySQL 执行 DROP USER 'cdc_reader'@'%';,按是否仍需合成数据决定删除 migration_lab。这些删除只适用于独占实验对象。生产退出先冻结 schema 与 connector 配置,记录最终 GTID/file/pos,等待消费者追平并完成分块 checksum 与业务不变量校验,保留约定回切窗口,再撤销账号、ACL、secret、topic 和监控。schema history 与 offset 的保留期应满足审计和回滚,但归档中不能含明文凭证。
回滚 connector 版本时,先确认旧版能理解新版本写入的 offset、history 和事件 schema;若不确定,使用隔离 worker group、独立 history 与业务 topic 做并行验证。数据库 failback 也要重新证明 GTID 包含关系,不能因为旧主机恢复就把 DNS 指回去。真正可治理的 CDC 不是“任务能自动重启”,而是每次启动、快照、DDL、切换、重放和退出都能说清当前坐标、结构版本、重复边界与责任人。
