深入理解 Flink CDC API:Event 模型、Data Source 与 Data Sink 设计解析
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
如果你计划为 Flink CDC 开发自己的连接器,或者希望参与 Flink CDC 的贡献,本文基于官方文档 Understand Flink CDC API,结合当前仓库flink-cdc-common模块的真实源码,带你逐层拆解 Flink CDC 的核心 API 体系:Event 事件模型(DataChangeEvent / SchemaChangeEvent 及其流转规则)、Data Source 侧的EventSourceProvider与MetadataAccessor接口、Data Sink 侧的EventSinkProvider与MetadataApplier接口。读完本文,你将能够说清楚"一条变更是如何从源头变成 Event、流经内部算子、最终写入目标系统的",并知道开发新连接器时需要实现哪些接口。
一、Event:贯穿整条数据管道的统一载体
在 Flink CDC 的语境中,Event 是 Flink 数据流中一种特殊的记录。它描述了源端外部系统捕获到的变更,先被 Flink CDC 内部构建的各个算子处理与转换,最终交给 Data Sink 写入(或应用)到目标端外部系统。
从源码结构看,这一模型落在flink-cdc-common模块的org.apache.flink.cdc.common.event包中,接口层次非常简洁:
- Event:所有事件的顶层接口,标记外部系统流入 Flink CDC 的事件类型;
- ChangeEvent:
Event的子接口,要求携带TableId tableId()——即该事件归属的表。所有变更事件(数据变更 + Schema 变更)都继承它; DataChangeEvent与SchemaChangeEvent:两个具体分支,分别描述数据行级变更与表结构变更。
每个变更事件都包含其归属的 Table ID 与事件载荷(payload)。按载荷类型不同,事件被划分为以下两大类。
1.1 DataChangeEvent:描述数据行级变更
DataChangeEvent 表示外部系统的数据变更(INSERT、UPDATE、DELETE 等),由 5 个字段组成(见源码中字段定义 L51-L63):
| 字段 | 类型 | 含义 |
|---|---|---|
tableId | TableId | 事件归属的表标识 |
before | RecordData | 变更前镜像(pre-image) |
after | RecordData | 变更后镜像(post-image) |
op | OperationType | 变更操作类型 |
meta | Map<String, String> | 可选的变更元数据,源码注释举例:MySQL binlog 文件名与位置(file name, pos) |
操作类型由枚举 OperationType 预定义为 4 种:INSERT、UPDATE、REPLACE、DELETE。各类型下 before/after 的取值约定如下:
- Insert:新增数据,
before = null,after = 新数据; - Delete:删除数据,
before = 被删除的数据,after = null; - Update:更新已有数据,
before = 变更前数据,after = 变更后数据; - Replace:源文档此处留白,但可从源码确认其语义——
replaceEvent(TableId, RecordData after)工厂方法(DataChangeEvent.java#L143-L146)构造出的事件before = null、after = 完整新数据,即"用一条完整新记录整体替换"的变更。
DataChangeEvent是Serializable类,不提供公开构造器,而是提供了一组静态工厂方法作为统一入口(L99-L154):insertEvent、deleteEvent、updateEvent、replaceEvent,每种都额外提供了带meta参数的重载。此外还有面向框架内部场景的静态方法:projectBefore/projectAfter/projectRecords(在 Transform 列投影后替换 before/after 镜像)和route(在路由规则改写 TableId 时重建事件,L201-L208)。开发连接器时应当始终通过这些工厂方法构造事件,而不是自行 new。
1.2 SchemaChangeEvent:描述表结构级变更
SchemaChangeEvent 是一个接口,其载荷描述外部系统中表结构的变化,包含两个关键方法:getType()(返回SchemaChangeEventType枚举值)和copy(TableId newTableId)(在路由改写表 ID 时生成副本)。
源文档列举的 Schema 变更事件类型包括:
AddColumnEvent:新增列;AlterColumnTypeEvent:列类型变更;CreateTableEvent:新表创建;同时用于描述某个DataChangeEvent之前已存在的表结构(即"为后续数据变更声明 Schema");DropColumnEvent:删除列;RenameColumnEvent:列重命名。
进一步从org.apache.flink.cdc.common.event包的源码结构看,实际实现的事件类型比源文档列举的更多,还包括AlterTableCommentEvent(表/列注释变更)、DropTableEvent(删表)、TruncateTableEvent(清空表)以及SchemaChangeEventWithPreSchema(携带变更前 Schema 快照的事件基类)。如果你实现的是需要透传删表或清表语义的连接器,可以对照 SchemaChangeEventType 枚举按需选用。
1.3 Flow of Events:Schema 必须先于 Data
你可能已经注意到,DataChangeEvent自身并不携带表结构信息。这样做可以减小数据变更事件的体积、降低序列化开销,但代价是事件"不自描述"。那么框架如何知道该如何解释这些数据变更?
解决方案是对事件的流转顺序提出强制约定:
- 如果某张表对框架来说是新的,那么在任何
DataChangeEvent之前必须先发出一个CreateTableEvent; - 如果某张表的 Schema 发生了变化,那么必须先发出对应的
SchemaChangeEvent,之后才能发出受影响的DataChangeEvent。
这一约定保证了框架在处理任何数据变更之前,一定已经知晓该表的 Schema。从源码侧也能看到配套机制:event包中的FlushEvent用于在 Schema 变更前后强制上下游同步与 flush;EventDeserializer则负责将反序列化后的事件流按规则重放。源文档中 "Flow of Events" 示意图(上文flow-of-events.png)直观展示了这一顺序:先是CreateTableEvent,然后是若干DataChangeEvent,中途插入 Schema 变更时,SchemaChangeEvent同样位于其影响的数据变更之前。
二、Data Source:EventSource 与 MetadataAccessor 的工厂
源文档将 Data Source 定位为EventSource与MetadataAccessor的"工厂"。落到源码,DataSource 接口(@PublicEvolving)的核心契约就是两个工厂方法:
public interface DataSource { /** Get the EventSourceProvider for reading events from external systems. */ EventSourceProvider getEventSourceProvider(); /** Get the MetadataAccessor for accessing metadata from external systems. */ MetadataAccessor getMetadataAccessor(); /** Get the SupportedMetadataColumn of the source. */ default SupportedMetadataColumn[] supportedMetadataColumns() { return new SupportedMetadataColumn[0]; } @Experimental default boolean isParallelMetadataSource() { return false; } }2.1 EventSourceProvider:把 Flink Source 装进 CDC 框架
EventSourceProvider 是一个 marker 接口。源码注释说明得很直接:它可以复用现有的 FlinkSource或 FlinkSourceFunction实现,未来也可支持 Flink CDC 自己的EventSource实现。仓库中提供了两个具体适配类:
- FlinkSourceProvider:适配 Flink 的
SourceAPI; - FlinkSourceFunctionProvider:适配传统
SourceFunction。
也就是说,EventSource本质上是一个读取外部变更、将其转换为 Event 并向下发射的 Flink Source。你可以参考 Flink 官方文档了解 Flink Source 的内部机制与实现方式(源文档给出的指引)。对连接器开发者而言,这意味着你只需要复用 Flink 社区的 Source 生态,把 Debezium/自研的读取逻辑包成一个Source/SourceFunction,即可接入 Flink CDC 的事件流。
2.2 MetadataAccessor:外部系统的元数据读取器
MetadataAccessor 充当外部系统的元数据读取器,提供 4 个方法(对应"列举 namespace、schema、表,并获取给定 TableId 的表结构"):
| 方法 | 说明 |
|---|---|
List<String> listNamespaces() | 列出所有 namespace;外部系统不支持 namespace 时抛UnsupportedOperationException |
List<String> listSchemas(@Nullable String namespace) | 列出所有 schema;namespace为 null 时列出全部 namespace 下的 schema |
List<TableId> listTables(@Nullable String namespace, @Nullable String schemaName) | 按 namespace 与 schema 列出所有TableId |
Schema getTableSchema(TableId tableId) | 获取给定表的Schema(表结构) |
这套接口与 Flink CDC 配置中的表通配符(namespace.schema.table)直接对应:include-tables匹配依赖listTables的逐级枚举能力,而整库同步时的CreateTableEvent生成则依赖getTableSchema。
此外,DataSource还暴露了两个值得注意的扩展点(源码 L36-L53):
supportedMetadataColumns():声明该源支持的元数据列(如源端 binlog 位置等),默认返回空数组;isParallelMetadataSource()(@Experimental):标记该源是否为"按表并行产生 Schema 变更事件"的源。源码注释明确指出:返回false(默认)会得到兼容 MySQL 这类"单一全局顺序变更流"源的常规算子拓扑;对 MongoDB、Kafka 这类不维护全局顺序 Schema 变更流的源应返回true。注释同时提醒这是实验特性,默认返回false以避免意外行为。
三、Data Sink:EventSink 与 MetadataApplier
与 Data Source 对称,DataSink 接口由EventSink与MetadataApplier组成,分别负责写入数据变更事件和应用 Schema 变更(元数据变更):
public interface DataSink { /** Get the EventSinkProvider for writing changed data to external systems. */ EventSinkProvider getEventSinkProvider(); /** Get the MetadataApplier for applying metadata changes to external systems. */ MetadataApplier getMetadataApplier(); /** 需要按数据变更事件做 Sink 前分区时,用于计算 hash 值 */ default HashFunctionProvider<DataChangeEvent> getDataChangeEventHashFunctionProvider() { return new DefaultDataChangeEventHashFunctionProvider(); } // ... }3.1 EventSinkProvider:当前基于 Flink Sink V2 API
EventSinkProvider 同样是一个 marker 接口,源码注释说明它可以复用现有的 FlinkSink(V2)与SinkFunction实现。EventSink接收上游算子发来的变更事件并应用到外部系统,当前只支持 Flink 的 Sink V2 API。仓库中对应两个适配类:
- FlinkSinkProvider:适配 Flink
Sink(V2); - FlinkSinkFunctionProvider:适配传统
SinkFunction。
因此编写新 Sink 连接器时,推荐路径仍是实现 FlinkSink<InputT, SinkWriterT>(V2 API)并将其包装进FlinkSinkProvider返回。
DataSink上的getDataChangeEventHashFunctionProvider是一个容易被忽略但很实用的扩展点:当需要在写入 Sink 之前按数据变更事件做分区(例如按表或按主键路由到不同下游分区)时,可通过 HashFunctionProvider 自定义 hash 计算逻辑,默认实现为DefaultDataChangeEventHashFunctionProvider;带并行度参数的重载版本在未覆写时会回退到无参版本。
3.2 MetadataApplier:Schema 变更的最终落地者
MetadataApplier 接口继承Serializable, AutoCloseable,当框架从源端收到 Schema 变更事件后,会先做一些内部同步与 flush(即前文事件流约定中 FlushEvent 发挥作用的环节),再通过该 applier 把 Schema 变更应用到外部系统。其接口方法如下:
| 方法 | 说明 |
|---|---|
void applySchemaChange(SchemaChangeEvent) | 核心方法:把给定的 Schema 变更应用到外部系统,可抛出SchemaEvolveException |
setAcceptedSchemaEvolutionTypes(Set<SchemaChangeEventType>) | 设置当前 applier 接受的 schema 演化事件类型(默认 no-op) |
acceptsSchemaEvolutionType(SchemaChangeEventType) | 判断本 applier 是否应处理某类事件,默认返回true |
getSupportedSchemaEvolutionTypes() | 声明下游能处理哪些 Schema 变更事件,默认支持SchemaChangeEventTypeFamily.ALL中的全部类型 |
close() | 关闭 applier 及其底层资源,默认空实现 |
这套"声明支持类型 + 按类型过滤"的设计,使得框架在事件流早期就能根据下游能力决定是否、以及如何执行 schema 演化。仓库中每个 Pipeline 连接器(如flink-cdc-connect/flink-cdc-pipeline-connectors下的 doris、kafka、starrocks、paimon 等模块)都按这一套DataSource/DataSink契约提供了自己的实现,可以挑一个你最常用的连接器源码作为开发参照。
四、小结与源码索引
把三个核心概念串起来,Flink CDC 的 API 体系可以概括为:
- Event 模型:
ChangeEvent携带TableId,分为DataChangeEvent(tableId + before + after + op + meta 五字段,四种操作类型)与SchemaChangeEvent(AddColumn / AlterColumnType / CreateTable / DropColumn / RenameColumn 等);核心流转约束是"新表先CreateTableEvent、变 Schema 先SchemaChangeEvent",保证任何DataChangeEvent被处理前 Schema 已知; - Data Source:
DataSource是工厂,返回EventSourceProvider(复用 FlinkSource/SourceFunction,读取变更并转成 Event 下发)与MetadataAccessor(listNamespaces/listSchemas/listTables/getTableSchema四个元数据方法),并可选声明元数据列与并行 Schema 变更能力; - Data Sink:
DataSink对称地返回EventSinkProvider(当前支持 Flink Sink V2 API)与MetadataApplier(applySchemaChange+ 支持类型声明),并可通过HashFunctionProvider自定义写入前的分区 hash。
关键源码入口索引(均为仓库相对路径):
- Event 顶层接口 / ChangeEvent / DataChangeEvent / OperationType / SchemaChangeEvent
- DataSource / EventSourceProvider / MetadataAccessor
- DataSink / EventSinkProvider / MetadataApplier
在此基础上,配合仓库内各 Pipeline 连接器的完整实现(flink-cdc-connect/flink-cdc-pipeline-connectors/目录)与开发者指南中的 Contribute to Flink CDC,你就具备了开发 Flink CDC 新连接器的完整 API 知识基础。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考