Flink CDC 3.7 新特性:已有目标表的 Schema 扩展机制

背景与问题

在 CDC 数据同步场景中,上游业务表的结构变更是常态。典型的场景包括:上游 MySQL 表在运行过程中新增字段,而下游 Paimon 或 Fluss 目标表已经由其他任务或外部流程提前创建,但其 Schema 比上游更窄。
例如,某业务线的 MySQL 订单表初始结构如下:
对应的 Paimon 目标表在一周前由另一个批处理任务提前创建,其结构相同。随后业务升级,MySQL 表新增了两个字段:
此时 MySQL 表的 Schema 已包含 source_system 和 event_time_ms,但 Paimon 目标表仍只有最初的四列。当新的 Flink CDC pipeline 启动时,初始 CreateTableEvent 到达 Paimon Sink,由于目标表已存在,Sink 会按目标表的现有 Schema 继续处理后续数据。结果是:上游新增的两列数据被直接忽略,既无报错也无告警,source_system 和 event_time_ms 的数据永久丢失。
notion image
在 Flink CDC 3.7 之前,这是框架层面的一个空白:运行时 schema.change.behavior 可以处理运行中的 AddColumnEvent,但如果初始建表时目标表就已经比上游窄,数据会被静默丢弃。
FLINK-40647 针对这一问题引入了已有目标表的 Schema 扩展能力,通过 existing-table.schema-expansion.mode 选项,使框架在初始建表阶段能够以可配置的方式与现有目标表进行 Schema 对齐。

新增选项:existing-table.schema-expansion.mode

该选项在 Sink 配置中声明,例如:
需要注意以下约束:
  • 该选项属于框架级选项,由 Flink CDC 直接消费,不会作为连接器属性透传给 MetadataApplier。
  • 不应将其配置在 catalog.properties.* 层级下,否则不会生效。
  • 模式值必须使用字符串形式,例如 "EXPAND";布尔值会被解析器拒绝。
  • 仅对流式作业生效;在 Batch 模式下,该选项会被忽略并输出告警日志。

四种工作模式

下图展示了四种模式在初始 CreateTableEvent 阶段的基本处理路径:
notion image
模式
目标表已存在时的行为
目标表不存在时的行为
失败处理
DISABLED(默认)
不检查、不扩展,使用 Sink 原生行为
Sink 自行建表
不适用
CHECK
校验上游列能否被目标表容纳,且主键一致;不执行任何 DDL
作业失败
任何不兼容、读取失败或能力缺失都会以聚合错误使作业失败
TRY_EXPAND
检查并尽力执行安全 DDL,随后回读验证目标表 Schema
Sink 自行建表
缺少扩展能力时作业失败;机制自身失败仅记录日志并回退到 Sink 原生行为
EXPAND
检查并执行安全 DDL,随后回读验证目标表 Schema
Sink 自行建表
任何不兼容、DDL 不支持、执行或验证失败都会使作业失败
DISABLED:保持原有行为
默认模式,与 3.7 之前的行为一致。适合上下游 Schema 已通过其他机制保证一致、无需框架干预的场景。
CHECK:校验优先,不执行 DDL
该模式适用于目标表 Schema 由外部系统(如 DBA、数据平台或 DDL 编排工具)管理的场景。CDC pipeline 不应修改下游表结构,但仍需防止上游 Schema 变宽导致的数据静默丢失。
在 CHECK 模式下:
  • 框架对比上游 pipeline Schema 与目标表 Schema;
  • 发现缺失列、类型不兼容或主键不一致时,以聚合错误终止作业;
  • 错误信息列出每一处差异,并附带方言无关的 ALTER TABLE 修复模板;
  • 不执行任何 DDL,因此不受 include.schema.changes 和 Sink DDL 能力的影响;
  • 独立于 schema.change.behavior,即使后者设为 IGNORE 或 EXCEPTION 也会执行。
TRY_EXPAND:尽力扩展,失败降级
该模式适合希望自动完成 Schema 扩展、同时允许在异常情况下降级到 Sink 原生行为的场景。其执行流程为:
  1. 探测目标表现有 Schema;
  1. 对可安全扩展的差异生成 DDL,例如将缺失的普通非键列添加为 nullable 列;
  1. 执行 DDL 后回读目标表 Schema 进行验证;
  1. 若扩展失败,记录日志并将原始 CreateTableEvent 交给 Sink 处理。
需要指出的是,TRY_EXPAND 仅在确认连接器支持扩展能力后才会吞掉机制自身的失败,它不会屏蔽 Sink 自身 Schema 处理抛出的错误,也不保证扩展失败后所有上游列都能落入目标表。
EXPAND:严格模式
与 TRY_EXPAND 执行相同的检查与 DDL 流程,但任何失败都会导致作业失败。适合对 Schema 一致性要求严格的 pipeline。

扩展范围与安全约束

该机制对可执行的操作有明确限制:
  • 新增缺失列:将普通非主键列以 nullable 形式追加到目标表;
  • 安全拓宽列类型:在目标系统支持的前提下,对普通非主键列进行安全类型拓宽;
  • 不修改表键:不会新增、删除或修改主键与分区键;
  • 主键比较时忽略分区列:由于 Paimon 等连接器会将分区列合并到存储主键中,框架在比较主键时会排除两侧的分区列,避免误判。
当主键不一致时:
  • CHECK 与 EXPAND 会直接使作业失败;
  • TRY_EXPAND 会跳过扩展,将 CreateTableEvent 交由 Sink 处理,最终由连接器自身的键校验机制决定行为。

与 schema.change.behavior 的关系

existing-table.schema-expansion.mode 与 schema.change.behavior 作用于不同阶段:
  • schema.change.behavior 控制运行时收到的 Schema 变更事件(如 AddColumnEvent、AlterColumnTypeEvent)的处理方式;
  • existing-table.schema-expansion.mode 控制初始 CreateTableEvent 遇到已有目标表时的行为。
CHECK 模式独立于 schema.change.behavior,即使后者设为 IGNORE 或 EXCEPTION 也会执行;TRY_EXPAND 和 EXPAND 在 IGNORE 或 EXCEPTION 模式下会跳过框架侧的初始处理。

配置示例

以下为一个完整的 YAML 配置示例:

选型建议

根据数据治理模式与一致性要求,可按以下原则选择模式:
  • 目标表 Schema 由外部系统管理:使用 CHECK,使 CDC 成为 Schema 一致性的校验者而非修改者。
  • 希望自动兼容、可接受机制级失败降级:使用 TRY_EXPAND,并配合日志与监控观察扩展结果。
  • 对 Schema 一致性要求严格:使用 EXPAND,任何不兼容都直接失败。
  • 无需改变现有行为:保持默认 DISABLED。

总结

FLINK-40647 填补了 Flink CDC 在初始建表阶段 Schema 对齐方面的空白。过去,框架只能在运行时通过 schema.change.behavior 应对 Schema 变更,而如果目标表在初始阶段就与上游不一致,数据可能被静默丢弃。通过 existing-table.schema-expansion.mode,用户可以根据自身的数据治理需求选择校验、尽力扩展或严格扩展策略。
当前 Paimon 与 Fluss 连接器已实现该能力。在 Flink CDC 3.7.0 中,建议相关用户根据实际场景启用该选项,并注意将模式值配置为字符串,同时确保作业运行在 streaming 模式下。
Loading...