场景目标
为一个有限源表建立源位置、Flink 状态、结构事件和目标结果的证据链,验证更新、删除、结构变化与恢复行为,明确源日志和外部 DDL 的不可回退边界。
大数据运维 · 高级
以 MySQL 单表为试点,核对 CDC 版本、增量快照与日志衔接、Schema 事件及目标提交,通过有限样本与隔离恢复验证同步边界。
为一个有限源表建立源位置、Flink 状态、结构事件和目标结果的证据链,验证更新、删除、结构变化与恢复行为,明确源日志和外部 DDL 的不可回退边界。
准备 CDC/Flink/MySQL/驱动/sink 版本、Pipeline 或 SQL 接入方式、源端受限查询权限、唯一 server-id 范围、Job ID、checkpoint 存储及独立目标。REST 示例要求 FLINK_REST_URL 为环境已配置认证的 HTTPS 网关,JOB_ID 为单个真实作业,命令只执行有限 GET。
以 MySQL 单表为试点,核对 CDC 版本、增量快照与日志衔接、Schema 事件及目标提交,通过有限样本与隔离恢复验证同步边界。
点击组件,在图下方查看职责;连线编号对应流向解读。小屏可横向滚动,或直接展开文字说明。
source 从已知表读取 chunk,并在正确日志边界继续变更流。
新表及新结构的数据前具有相应 schema 事件,按 Pipeline 模型协调。
结构处理与数据写入协同,仍受实际 sink 能力限制。
通过目标连接器申请结构变化;异步任务完成状态与权限要到目标核验。
实际 append、upsert 或事务协议决定重放及可见性,不能统一标成 exactly-once。
协调 source 状态快照与恢复,持久进度包含 split 和日志位置。
sink 是否利用 checkpoint 协调提交取决于连接器;此连线不承诺跨系统事务。
保存 checkpoint ID、完成时间、失败及实际恢复输入。
对照有限主键、最终值、删除行为及已完成 DDL,确认外部影响。
图对应 CDC 3.4.x Pipeline 与 Flink 1.20;兼容矩阵另支持 1.19,不直接用于 Flink 2.x 或静态 SQL schema 的 DDL 结论。
并行 chunk 不是全库同一时刻视图,source 可恢复不等于目标 exactly-once;Flink checkpoint 不能撤销已完成的目标 DDL。
示例只读诊断,变化与失败恢复只在有限隔离试点验证;所需日志丢失、键覆盖或 DDL 状态不明时停止扩量。
图解是本站基于官方资料整理的逻辑参考;实施前仍需核对实际部署版本、组件支持范围与变更审批。
采用 Flink CDC 3.4.x、Flink 1.20 与 MySQL 8.0.x,优先选择有稳定主键的 InnoDB 单表。CDC 3.4.x 官方支持 Flink 1.19/1.20;Pipeline 与 SQL source 的结构演进能力分开核验。
本流程先进行只读诊断,再在独立目标演练一次快照期间变化、结构变化和失败恢复。预计时间仅覆盖有限试点,不包含大表全量初始化。文档及参考图均待环境验证,不提供在线体验。
登记 CDC 3.4.x 与 Flink 1.19/1.20 兼容组合,本试点采用 Flink 1.20、MySQL 8.0.x。记录具体驱动、sink 制品、接入方式、源表负责人、目标表和主键删除语义;选择独立目标与有限数据范围。
版本和依赖可对应到实际制品;Pipeline 与 SQL source 的结构演进边界已写明。
依赖不匹配或目标写入语义不清时停止部署,只保留清单和待补项。
在已认证的 MySQL 会话中只读检查一个源与一个已知表,将示例表名替换为批准对象。核对 ROW、FULL、权限、server-id 登记及真实可用日志,保留预算应覆盖快照、停顿和追赶。
SELECT VERSION(), @@GLOBAL.log_bin, @@GLOBAL.binlog_format,
@@GLOBAL.binlog_row_image, @@GLOBAL.binlog_expire_logs_seconds;
SHOW CREATE TABLE cdc_lab.orders;源表有稳定主键,必需日志及行镜像满足连接器要求;实际日志覆盖窗口与配置秒数分别记录。
日志缺口、键不稳定或权限不足时暂停初始化,不通过改 latest-offset 跳过缺口。
在受控配置评审中限定单表、唯一 server-id 范围和初始并行度。以下为不完整 Pipeline 片段,需合并已核验的连接、认证和 sink 配置;不假定 YAML 自动展开环境变量。
source:
type: mysql
tables: cdc_lab.orders
server-id: 5401-5410
scan.startup.mode: initial
scan.incremental.snapshot.chunk.size: 2048
pipeline:
parallelism: 2表匹配没有扩大范围;server-id 与其他复制客户端不冲突,chunk 预算考虑行宽与源库负载。
源库延迟或缓冲压力超预算时暂停试点并恢复上次配置;已写入数据需独立对账。
按表发现 split 进度指标,关联最慢 chunk、已完成 checkpoint 和后续日志读取。LOW/HIGH 是 binlog 位置,不是事件时间 Watermark;快照结束后等待相应 checkpoint 完成再衔接日志。对一个 Job 只读采集一次,10 秒超时后停止并记录错误。
curl --fail --silent --show-error --max-time 10 \
"${FLINK_REST_URL}/jobs/${JOB_ID}/checkpoints"snapshot 与 binlog 阶段可解释,split 进度和 checkpoint 同窗对应;采样结果不含凭据。
checkpoint 持续失败或 sink 反压时停止扩量,保留状态和日志,不清理 checkpoint。
针对单表列出允许和拒绝的结构事件,记录 schema.change.behavior 与 include/exclude,排除规则优先。跟踪 source 事件、MetadataApplier、目标 DDL 状态和后续数据,分别检查类型、空值与默认值。
允许的变化能在目标确认完成;被过滤或不支持的变化有已说明行为,源库发布另有约束。
遇到未知类型损失、目标 DDL 状态不明或重复失败时停止扩大范围;回退 Flink 配置不会撤销外部 DDL。
明确目标是事件追加、主键 upsert 还是事务提交;选定少量业务 ID 比较更新、删除和重复重放。多源 route 合并时加入来源相同主键样本;若目标为 Kafka,再单独核对键、分区和事务可见性。
目标行为与声明契约一致,source 恢复保证和外部提交保证没有合并为一个配置标签。
发生键覆盖、漏删或重复副作用时保留原始事件并停止试点;不重置消费位点盲目重放。
在专用小表将样本控制为约 100 个主键,覆盖快照期间少量插入、更新、删除和允许的加列。追到约定日志边界后比较主键及最终值,再做一次可撤销的隔离故障恢复,记录恢复 checkpoint、源位置及已完成目标 DDL。
正常与恢复后的值、删除、结构和重复行为符合预期;单次演练仅证明该版本和窗口。
恢复不相容时保留原状态、制品与独立目标,选择向前兼容修复或重新初始化,不接回生产目标试错。
归档版本、配置、Job ID、源位点、checkpoint、目标结构和有限对账结果,说明 binlog 保留余量、未覆盖故障和下一次采样时间。切换入口前独立确认新写入归属和业务补偿路径。
每项结论有真实样本和时间窗;未进行环境演练的内容保持 PENDING,不将图示当作部署证据。
任一关键一致性检查失败时保持原入口;配置回退不撤销已写数据或 DDL,按记录的外部变化处理。