大数据运维 · 高级

Flink CDC 快照、结构演进与同步恢复

以 MySQL 单表为试点,核对 CDC 版本、增量快照与日志衔接、Schema 事件及目标提交,通过有限样本与隔离恢复验证同步边界。

Apache FlinkApache Flink CDCMySQL大数据

场景目标

为一个有限源表建立源位置、Flink 状态、结构事件和目标结果的证据链,验证更新、删除、结构变化与恢复行为,明确源日志和外部 DDL 的不可回退边界。

环境要求

准备 CDC/Flink/MySQL/驱动/sink 版本、Pipeline 或 SQL 接入方式、源端受限查询权限、唯一 server-id 范围、Job ID、checkpoint 存储及独立目标。REST 示例要求 FLINK_REST_URL 为环境已配置认证的 HTTPS 网关,JOB_ID 为单个真实作业,命令只执行有限 GET。

参考架构 · 非实时拓扑

Flink CDC 快照、结构事件与目标提交架构

以 MySQL 单表为试点,核对 CDC 版本、增量快照与日志衔接、Schema 事件及目标提交,通过有限样本与隔离恢复验证同步边界。

  • 数据 / 请求
  • 控制 / 管理
  • 观测 / 查询

点击组件,在图下方查看职责;连线编号对应流向解读。小屏可横向滚动,或直接展开文字说明。

Flink CDC 快照、结构事件与目标提交架构:组件关系图以 MySQL 单表为试点,核对 CDC 版本、增量快照与日志衔接、Schema 事件及目标提交,通过有限样本与隔离恢复验证同步边界。 MySQL 源表与日志 → CDC Source:快照与日志事件;CDC Source → 结构事件协调:数据与结构事件;结构事件协调 → 数据 Sink:协调后的事件;结构事件协调 → 独立目标系统:MetadataApplier / DDL;数据 Sink → 独立目标系统:数据写入与提交;Flink 状态恢复 → CDC Source:检查点与恢复位置;Flink 状态恢复 → 数据 Sink:检查点协作;Flink 状态恢复 → 有限验收证据:完成与恢复证据;独立目标系统 → 有限验收证据:结构与业务对账。箭头说明见下方流向解读。
逻辑参考图,待环境验证。以 CDC 3.4.x / Flink 1.20 / MySQL 8.0.x 为基线,展示 Pipeline 单表同步与恢复责任;不代表跨表全局快照、跨数据库事务或任意 sink 的 exactly-once。

MySQL 源表与日志

数据 / 制品

源表快照和持续 binlog 必须能够衔接;GTID 与保留配置不能代替实际日志可用性证据。

查看关联工具
全部组件职责 7 个组件
MySQL 源表与日志
源表快照和持续 binlog 必须能够衔接;GTID 与保留配置不能代替实际日志可用性证据。工具介绍 MySQL 源表与日志
CDC Source
按 chunk 修正 LOW/HIGH 之间的变化,完成快照 checkpoint 后衔接持续 binlog;LOW/HIGH 不是事件时间 Watermark。工具介绍 CDC Source
结构事件协调
依据演进策略协调新结构事件和数据;目标 DDL 通过 sink 的 MetadataApplier 实现,能力按连接器核验。工具介绍 结构事件协调
Flink 状态恢复
逻辑上合并展示 checkpoint 协调与存储依赖,记录 source 进度及相关算子状态;目标提交保证另验,外部 DDL 不随状态回退。工具介绍 Flink 状态恢复
数据 Sink
根据实际连接器写入目标;主键、重放、删除与提交语义需要逐项验证,不由 source 一致性自动推导。工具介绍 数据 Sink
有限验收证据
固定一个源表和有限主键集合,对齐 checkpoint、目标结构及更新删除结果,交接停止条件。
独立目标系统
演练使用独立目标;数据与结构变化可能已对外完成,恢复旧 Flink 状态不会自动撤销它们。
流向解读 9 条连接
  1. 1

    MySQL 源表与日志 CDC Source

    数据 / 请求 · 快照与日志事件

    source 从已知表读取 chunk,并在正确日志边界继续变更流。

  2. 2

    CDC Source 结构事件协调

    数据 / 请求 · 数据与结构事件

    新表及新结构的数据前具有相应 schema 事件,按 Pipeline 模型协调。

  3. 3

    结构事件协调 数据 Sink

    数据 / 请求 · 协调后的事件

    结构处理与数据写入协同,仍受实际 sink 能力限制。

  4. 4

    结构事件协调 独立目标系统

    控制 / 管理 · MetadataApplier / DDL

    通过目标连接器申请结构变化;异步任务完成状态与权限要到目标核验。

  5. 5

    数据 Sink 独立目标系统

    数据 / 请求 · 数据写入与提交

    实际 append、upsert 或事务协议决定重放及可见性,不能统一标成 exactly-once。

  6. 6

    Flink 状态恢复 CDC Source

    控制 / 管理 · 检查点与恢复位置

    协调 source 状态快照与恢复,持久进度包含 split 和日志位置。

  7. 7

    Flink 状态恢复 数据 Sink

    控制 / 管理 · 检查点协作

    sink 是否利用 checkpoint 协调提交取决于连接器;此连线不承诺跨系统事务。

  8. 8

    Flink 状态恢复 有限验收证据

    观测 / 查询 · 完成与恢复证据

    保存 checkpoint ID、完成时间、失败及实际恢复输入。

  9. 9

    独立目标系统 有限验收证据

    观测 / 查询 · 结构与业务对账

    对照有限主键、最终值、删除行为及已完成 DDL,确认外部影响。

从架构到实施

  1. 01

    锁定源与初始化范围

    记录兼容组合、源日志、稳定主键和分块资源,确定单表试点。

  2. 02

    联查状态、结构和目标提交

    沿快照衔接、Schema 事件和 sink 契约定位数据停滞或不一致。

  3. 03

    证明有限恢复并交接边界

    通过固定样本与一次隔离恢复对账,保留不能随状态回退的外部变化。

故障域与操作边界

版本与接入方式

图对应 CDC 3.4.x Pipeline 与 Flink 1.20;兼容矩阵另支持 1.19,不直接用于 Flink 2.x 或静态 SQL schema 的 DDL 结论。

一致性与恢复

并行 chunk 不是全库同一时刻视图,source 可恢复不等于目标 exactly-once;Flink checkpoint 不能撤销已完成的目标 DDL。

范围与执行

示例只读诊断,变化与失败恢复只在有限隔离试点验证;所需日志丢失、键覆盖或 DDL 状态不明时停止扩量。

架构依据与版本核对 3 篇官方资料

图解是本站基于官方资料整理的逻辑参考;实施前仍需核对实际部署版本、组件支持范围与变更审批。

方案说明

适用范围

采用 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 的结构演进能力分开核验。

验收范围

本流程先进行只读诊断,再在独立目标演练一次快照期间变化、结构变化和失败恢复。预计时间仅覆盖有限试点,不包含大表全量初始化。文档及参考图均待环境验证,不提供在线体验。

配套知识

官方参考

工具编排

3 个关联工具
  1. Apache Flink CDC快照与变更事件采集检查分块、日志位置、路由和结构演进配置,核对实际连接器能力。
  2. Apache Flink状态恢复与作业运行读取 checkpoint、任务状态和反压证据,使用隔离目标验证恢复。
  3. MySQL源表与二进制日志确认主键、日志保留与源端负载,提供有限一致性样本。

实施步骤

共 8 步
  1. 01

    固定版本与同步契约

    登记 CDC 3.4.x 与 Flink 1.19/1.20 兼容组合,本试点采用 Flink 1.20、MySQL 8.0.x。记录具体驱动、sink 制品、接入方式、源表负责人、目标表和主键删除语义;选择独立目标与有限数据范围。

    验证标准

    版本和依赖可对应到实际制品;Pipeline 与 SQL source 的结构演进边界已写明。

    停止与回退

    依赖不匹配或目标写入语义不清时停止部署,只保留清单和待补项。

    返回步骤起点
  2. 02

    检查源日志与稳定主键

    在已认证的 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 跳过缺口。

    返回步骤起点
  3. 03

    限定采集范围与分块预算

    在受控配置评审中限定单表、唯一 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 预算考虑行宽与源库负载。

    停止与回退

    源库延迟或缓冲压力超预算时暂停试点并恢复上次配置;已写入数据需独立对账。

    返回步骤起点
  4. 04

    核对快照到日志的衔接

    按表发现 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。

    返回步骤起点
  5. 05

    评审结构事件与目标 DDL

    针对单表列出允许和拒绝的结构事件,记录 schema.change.behavior 与 include/exclude,排除规则优先。跟踪 source 事件、MetadataApplier、目标 DDL 状态和后续数据,分别检查类型、空值与默认值。

    验证标准

    允许的变化能在目标确认完成;被过滤或不支持的变化有已说明行为,源库发布另有约束。

    停止与回退

    遇到未知类型损失、目标 DDL 状态不明或重复失败时停止扩大范围;回退 Flink 配置不会撤销外部 DDL。

    返回步骤起点
  6. 06

    验证 sink 的键与提交语义

    明确目标是事件追加、主键 upsert 还是事务提交;选定少量业务 ID 比较更新、删除和重复重放。多源 route 合并时加入来源相同主键样本;若目标为 Kafka,再单独核对键、分区和事务可见性。

    验证标准

    目标行为与声明契约一致,source 恢复保证和外部提交保证没有合并为一个配置标签。

    停止与回退

    发生键覆盖、漏删或重复副作用时保留原始事件并停止试点;不重置消费位点盲目重放。

    返回步骤起点
  7. 07

    进行一次隔离恢复验收

    在专用小表将样本控制为约 100 个主键,覆盖快照期间少量插入、更新、删除和允许的加列。追到约定日志边界后比较主键及最终值,再做一次可撤销的隔离故障恢复,记录恢复 checkpoint、源位置及已完成目标 DDL。

    验证标准

    正常与恢复后的值、删除、结构和重复行为符合预期;单次演练仅证明该版本和窗口。

    停止与回退

    恢复不相容时保留原状态、制品与独立目标,选择向前兼容修复或重新初始化,不接回生产目标试错。

    返回步骤起点
  8. 08

    交接恢复窗口与持续观察

    归档版本、配置、Job ID、源位点、checkpoint、目标结构和有限对账结果,说明 binlog 保留余量、未覆盖故障和下一次采样时间。切换入口前独立确认新写入归属和业务补偿路径。

    验证标准

    每项结论有真实样本和时间窗;未进行环境演练的内容保持 PENDING,不将图示当作部署证据。

    停止与回退

    任一关键一致性检查失败时保持原入口;配置回退不撤销已写数据或 DDL,按记录的外部变化处理。

    返回步骤起点

DOUYA OPS ECOSYSTEM

体验豆芽自研工具与场景能力

部分场景提供体验环境,用于功能验证、测试和技术交流。