最佳实践 · 大数据运维

待环境验证

Flink CDC Schema 事件、目标表演进与恢复边界

辨别 Pipeline 结构事件、演进策略、类型和路由契约,联查目标 DDL 与后续数据,验证 checkpoint 恢复不能撤销外部结构变化的边界。

Apache FlinkApache Flink CDCMySQL大数据
阅读导引 · 理解后再操作

这篇知识解决什么问题

从一条 SchemaChangeEvent 追踪到目标 DDL 与新结构数据,判断演进策略、类型映射和路由是否保持业务契约,并验证恢复时外部结构变化的处理。

行为模式决定不同取舍

exception、evolve、try_evolve、lenient、ignore 处理结构事件的方式不同;默认 lenient 不能推导源目标结构完全一致。

键和字段契约跨越连接器边界

目标支持加列不代表支持任意类型与主键变化,多源路由还需防止来源内唯一键在目标发生冲突。

进入文章正文
关联架构图解7 个组件 · 点击展开

CDC 图将结构协调与数据 sink 分开,目标既接收数据也接收经 MetadataApplier 申请的结构变化。checkpoint 控制线表达 Flink 恢复协作,不提供目标 DDL 回滚;实际异步 DDL 完成状态需在目标系统补证。

查看场景架构与实施步骤
参考架构 · 非实时拓扑

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,确认外部影响。

故障域与操作边界

版本与接入方式

图对应 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 Pipeline API,运行时基线为 Flink 1.20,源端示例为 MySQL 8.0.x。3.4 官方兼容矩阵列出 Flink 1.19/1.20;连接器支持的数据类型、DDL 和提交保证仍需按实际 sink 版本核验。本文是待环境验证的参考,所有示例均未执行,状态 PENDING。兼容矩阵

Pipeline 传递表身份、数据事件及结构事件。使用 mysql-cdc 的 Flink SQL 表并不自动等于部署了这一整套 Pipeline DDL 演进机制,静态 SQL schema、查询计划和目标写入能力需要另外确认。先记录接入方式和连接器制品,避免把一个产品名当作统一能力开关。

Schema 事件怎样进入数据流

DataChangeEvent 表示表身份、操作类型和变化前后数据;SchemaChangeEvent 描述建表、加列、改列类型等结构变化。对于新表,CreateTableEvent 先于其数据事件;发生结构变化时,相关 SchemaChangeEvent 先于使用新结构的数据。事件顺序使下游有机会先准备 schema,再解释随后的记录。

Sink 不只有写数据的一侧:EventSink 负责数据,MetadataApplier 负责应用目标结构变化。框架协调刷新与结构变更,但目标系统是否支持某种 DDL、DDL 是否异步、权限与类型映射是否满足,都属于连接器和目标系统的真实能力。不能由“收到 SchemaChangeEvent”推导目标表已经完成变化。Pipeline API 与事件模型

五种行为模式不是同一种宽容

  • exception: 遇到结构变化直接失败,适合需要先评审变更的契约,但会暂停同步。
  • evolve: 尝试把结构变化应用到下游,应用失败会触发作业失败恢复;不会自动补齐目标不支持的 DDL。
  • try_evolve: 尝试应用并容忍不支持的变化,后续记录可能需要转换;不能把继续运行当作无损兼容。
  • lenient: 3.4 默认模式,采用更保守的结构调整以避免丢弃原有数据,例如可能保留旧列并增加新列;目标结构未必与源端完全相同。
  • ignore: 不向下游应用结构变化,不会自动让旧目标正确接纳任意新字段。

模式写在 pipeline.schema.change.behavior。默认值仍需要在变更单中显式记录,特别是升级后不能只验证作业能启动。lenient 默认对删表、截断表等破坏性变化的处理,也要结合事件级配置核验。结构演进策略

一份可评审的事件过滤片段

下面只展示策略片段,需合并到已核验的完整 Pipeline 与具体 sink 配置,不能直接运行。

sink:
  include.schema.changes:
    - create.table
    - add.column
  exclude.schema.changes:
    - drop.column
    - drop.table
    - truncate.table
pipeline:
  schema.change.behavior: evolve

这份策略只允许所列结构事件;exclude 优先于 include。它用于评审“目标保留历史结构、仅开放建表和加列”的契约,不意味着其他源端 DDL 会被阻止执行。比如源端改列类型被过滤,后续记录仍可能与目标旧类型不相容。因此源库变更流程还需提前阻断未获准 DDL,而不能指望过滤器替代数据库权限或发布审批。

类型与业务键必须逐项核对

验收字段不能只比较名称,还要记录精度、scale、符号、字符编码、空值、默认值、时间语义和目标列能力。MySQL SQL source 的 BIGINT UNSIGNED 可能映射成 DECIMAL(20,0),不能机械套到只支持有符号 64 位整数的目标字段;SQL source 的类型映射也不能原样当作所有 Pipeline sink 的映射表。

选取接近边界的正负数、超长字符串、NULL 与毫秒/微秒时间样本,比较传输前后的值。rename、扩宽、缩窄和主键变化要分别演练;某目标支持加列,不代表支持所有这些变化。默认值在历史行与新增行上的效果也可能不同,应以目标系统的读取结果验收。MySQL source 类型映射

路由合并会改变正确性问题

Pipeline route 可以把源表映射到其他目标表,也可让多个源表汇入同一目标。此时同名列的类型演进需要协调,源表内唯一主键不一定在合并目标中唯一。例如两个租户各有 order_id=100,若目标只用 order_id 做 upsert,会相互覆盖;应把来源身份纳入可验证的键设计。

记录每条 route 的匹配范围和目标表名,使用已知表名验证正反样本,避免正则范围意外覆盖其他库。多源结构不一致时,作业继续运行可能只表示完成了某种兼容转换,不表示业务字段仍完整。不要用目标总行数增长掩盖键覆盖与字段丢失。路由规则

用有限证据定位演进卡点

先固定一个源表、一个结构事件时间窗和一个目标表,读取已知表定义。以下只读示例只访问一个 MySQL 表;目标是其他数据库时采用其对应有限元数据查询。

SHOW CREATE TABLE cdc_lab.orders;
SELECT COLUMN_NAME, COLUMN_TYPE, IS_NULLABLE, ORDINAL_POSITION
FROM information_schema.COLUMNS
WHERE TABLE_SCHEMA = 'cdc_lab' AND TABLE_NAME = 'orders'
ORDER BY ORDINAL_POSITION
LIMIT 100;

依次检查 source 是否捕获事件、过滤策略是否允许、MetadataApplier 是否申请并完成 DDL、后续数据是否成功提交。将首个异常、事件类型、目标 DDL 任务、重启次数和 checkpoint 时间线对齐。源端结构已变而 sink 报转换错误时,应先查旧目标结构与映射,不能仅增加 checkpoint timeout 或无限重启。

DDL 与 checkpoint 的恢复边界

Flink 状态回退不等于目标数据库结构回退。一个 DDL 可能已在目标完成,而作业尚未形成新的可恢复 checkpoint;恢复旧状态时,连接器需要能处理已存在结构、重放事件及数据映射。对异步 DDL,还需确认目标任务最终状态,不能只看提交请求返回。

同样,事件刷新和演进协调不是跨数据库分布式事务,无法自动撤销已发布的列变更或外部读取结果。升级连接器、改变路由或更换 schema 行为前,应在独立目标用真实状态副本演练兼容性,保留旧制品、配置、目标结构和恢复位置。快照与增量衔接可参考 CDC 一致性验证

验收、停止与回退

在隔离小表依次演练允许的加列、未允许的改类型、明确拒绝的删除类变化以及 DDL 前后一次恢复。每轮分别比较源 schema、目标 schema、代表性历史行与新增行、主键删除语义和错误记录;验证多源路由时增加相同主键的来源冲突样本。

出现未预期的字段丢失、键覆盖、目标 DDL 状态不明或不兼容恢复时,停止扩大同步范围。恢复旧配置不能撤销已提交 DDL;需要向前修复兼容结构或在独立目标重新同步、对账后切换。不要直接对生产目标执行反向 DROP,也不要删除 checkpoint 试图清除结构冲突。交接应列明已经发生的外部变化和未验证项。

参考资料

从现象到判断

先收集证据,再缩小范围。以下是判读路径,不代表已经确认根因或获准变更。

  1. 源表已经加列,目标仍旧结构或后续记录报类型错误。

    只读核对
    逐段核对事件捕获、include/exclude、MetadataApplier 结果与目标表定义。
    如何判读
    先确定事件在哪一层被阻断,不能仅延长 checkpoint 超时。
  2. try_evolve 后作业持续运行,但某字段出现空值或转换差异。

    只读核对
    对照源值、事件 schema、目标类型和边界值样本。
    如何判读
    容忍结构变化不表示信息无损,继续运行不是验收结果。
  3. 从旧 checkpoint 恢复后反复遇到已存在列或结构不匹配。

    只读核对
    对齐已完成外部 DDL、恢复状态与连接器处理重复事件的能力。
    如何判读
    Flink 状态回退不会撤销目标结构,可能需要向前兼容修复或独立重建。
常见误区与判断边界 2 项

把事件过滤当源库 DDL 保护

过滤器不会阻止源端执行改类型或删列;源库发布流程仍需约束未允许变化。

多源合并沿用来源内主键

不同来源的同值键可能互相覆盖,必须验证目标键包含足够来源身份。

交接时应留下的证据

作为记录提纲使用,不是自动检查结果;未取得的证据应标记缺口,并注明负责人。

  • 归档 Pipeline/SQL 接入方式、制品版本及源目标表定义。
  • 记录演进模式、事件过滤优先级和支持的 DDL 清单。
  • 保留边界值、历史行、新行及来源键冲突样本。
  • 列出已完成外部 DDL、一次恢复结果和不能随配置撤销的变化。

记录需包含环境、版本、时间与时区;分享前脱敏,不附访问令牌、密码或完整业务敏感数据。

继续阅读与资料核对

补充相关主题,再结合当前环境的实施记录形成结论。

返回原理导读

DOUYA OPS ECOSYSTEM

贡献你的经验,帮助更多运维人

把故障复盘、标准流程和最佳实践沉淀为可检索、可复用的知识内容。