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 的数据永久丢失。
在 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 阶段的基本处理路径:
模式 | 目标表已存在时的行为 | 目标表不存在时的行为 | 失败处理 |
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 原生行为的场景。其执行流程为:
- 探测目标表现有 Schema;
- 对可安全扩展的差异生成 DDL,例如将缺失的普通非键列添加为 nullable 列;
- 执行 DDL 后回读目标表 Schema 进行验证;
- 若扩展失败,记录日志并将原始
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...