适用版本与数据范围
本文采用 Flink CDC 3.4.x、Flink 1.20 和 MySQL 8.0.x,讨论有稳定主键的 InnoDB 单表初始同步与持续变更。官方兼容矩阵列出 CDC 3.4.x 对应 Flink 1.19 与 1.20,不应直接套到 Flink 2.x。CDC Pipeline 与 SQL/DataStream source 是不同接入方式,本文 YAML 仅用于说明 Pipeline source 的局部配置。连接器兼容矩阵
操作前记录 CDC、Flink、MySQL、JDBC 驱动和 sink 连接器版本,确认驱动依赖实际存在。官方文档已核对,以下命令未在业务环境执行,验证状态为 PENDING。整库多表事务一致性、主键变更和无主键表需单独方案,不由本篇单表验收推导。
全量与增量为什么不能简单拼接
如果先扫完整张表,再从扫描结束的 binlog 位置开始消费,扫描期间已经读过的行发生的更新可能被遗漏;若从扫描开始位置重新消费,又需要处理扫描结果和日志重叠。Flink CDC 增量快照将表按 chunk 拆分,使并行读取、重叠修正和恢复进度具备明确边界。
每个 chunk 记录 LOW binlog 位置,读取该范围的行到缓冲,再记录 HIGH,并使用 LOW 到 HIGH 之间对应主键的变化修正该 chunk 结果。此处 LOW/HIGH 表示日志位置,不是 Flink 的事件时间 Watermark;一个是快照衔接边界,另一个用于事件时间推进,故障定位不能混用。MySQL CDC source 原理
快照完成与持续日志的切换
chunk 可以并行读取并按 chunk 进度 checkpoint;所有快照 chunk 完成后,还需等待相应进度进入已完成 checkpoint,再进入后续 binlog 读取。日志阶段使用单个 binlog reader,不能用提升 source 并行度直接推断日志阶段吞吐按比例增长。
不同 chunk 的读取时间不相同,因此其阶段性输出不是同一墙钟时刻的数据库全局快照。业务需要在追赶到明确日志边界后对账;只看“快照扫描结束”就对外开放跨表报表,可能把过渡状态误当作最终一致结果。增加下游并行度也不会自动保留原数据库事务在多个目标表中的原子可见性。
用有限检查确认源端前提
在已认证的 MySQL 会话中,对一个指定源执行以下只读查询。它只读取变量与一个已知表定义,不扫描业务全表;将 cdc_lab.orders 替换为已批准的实际表。
SELECT VERSION(), @@GLOBAL.log_bin, @@GLOBAL.binlog_format,
@@GLOBAL.binlog_row_image, @@GLOBAL.gtid_mode,
@@GLOBAL.binlog_expire_logs_seconds;
SHOW CREATE TABLE cdc_lab.orders;确认 binlog 开启、ROW 格式及 FULL 行镜像,核对主键、字符集、时区、源表过滤和权限。增量快照不依赖全局读锁,不表示对源库没有影响;并行扫描仍会产生查询、网络和缓冲成本。日志保留预算至少覆盖初始快照、最长故障停顿、追赶时间及余量;配置的保留秒数不是“所需位置一定还在”的证据,还需核对实际可用日志。MySQL 8.0 二进制日志变量
限定单表与分块资源
下列为受控 Pipeline 文件的配置片段,省略连接与认证部分,不能直接提交;凭据通过实际部署已支持的方式配置,不假定 YAML 自动展开任意环境变量。
source:
type: mysql
tables: cdc_lab.orders
server-id: 5401-5410
scan.startup.mode: initial
scan.incremental.snapshot.chunk.size: 2048
pipeline:
parallelism: 2server-id 范围需要在复制拓扑和并发 CDC 作业之间唯一,并容纳并行 reader;示例数字必须先经环境登记。chunk.size 是试点值,需结合行宽、键分布和源端负载调整,不能只按行数估算内存。稳定主键通常比会被更新的非主键分块列更容易保证范围一致性;热点键和大行可能让个别 chunk 成为尾部瓶颈。Pipeline MySQL 配置
启动模式与恢复状态分别核对
initial 用于新作业先快照再跟日志;latest-offset 跳过已有行,只跟随后续变化;snapshot 只读快照后结束。specific-offset 或 timestamp 需要对应历史日志仍可用,不能靠修改启动模式找回已清理的日志。恢复已有作业时应以加载的状态、split 进度和日志位置为准,不把新作业初始化参数当作覆盖恢复位点的开关。
GTID 可帮助识别事务与故障切换位置,但新源仍必须包含恢复所需的事务及可读取日志,相关复制配置也要匹配。仅修改 hostname 不构成无损切换证明。对低更新表,heartbeat 有助于位点推进;它不会延长源端日志保留,也无法弥补已经发生的日志缺口。
分层观察停滞在哪里
先按数据库和表发现实际 metric ID,再观察 isSnapshotting、isStreamReading 与剩余/完成 split 数量,避免复制其他部署的指标前缀。快照剩余数量不变时,关联最慢 chunk、查询耗时和源库压力;快照已完成却未进入日志阶段时,检查完成 checkpoint 和下游反压。
若已进入日志阶段,结合 source 位置、checkpoint 时间线、sink 提交延迟与业务更新时间判断追赶。Flink checkpoint 完成只能作为状态恢复的证据之一,source 的一致性保证不自动转化为外部 sink 的 exactly-once。append-only 目标、upsert 目标及事务型目标应分别验证重放效果,可联读 Flink 检查点恢复。
验收、误区与停止条件
使用专用小表和固定主键集合,覆盖快照期间更新、删除、新增以及跨 chunk 边界的样本。在约定日志边界追平后比较主键集合、最终值和删除结果,再做一次隔离失败恢复,核对 checkpoint、恢复位置与目标重复行为。不要用源表 COUNT 与目标 COUNT 相等替代值、空值和删除语义验证。
所需 binlog 已丢失、键范围出现无法解释的漏行、源库延迟超预算或下游重复副作用无法补偿时,应停止扩大试点。保留原状态和日志证据,选择独立目标重新初始化并核对切换窗口;删除 checkpoint 后接回原目标可能重复写入,改 latest-offset 会丢掉初始化缺口。结构变更应继续执行 Schema 与 sink 演进 的独立验收。