适用版本与接入方式
本文依据 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 试图清除结构冲突。交接应列明已经发生的外部变化和未验证项。