深入理解 Flink CDC API:Event 模型、Data Source 与 Data Sink 设计解析
2026/9/17 1:36:01 网站建设 项目流程

深入理解 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 侧的EventSourceProviderMetadataAccessor接口、Data Sink 侧的EventSinkProviderMetadataApplier接口。读完本文,你将能够说清楚"一条变更是如何从源头变成 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 变更)都继承它;
  • DataChangeEventSchemaChangeEvent:两个具体分支,分别描述数据行级变更与表结构变更。

每个变更事件都包含其归属的 Table ID 与事件载荷(payload)。按载荷类型不同,事件被划分为以下两大类。

1.1 DataChangeEvent:描述数据行级变更

DataChangeEvent 表示外部系统的数据变更(INSERT、UPDATE、DELETE 等),由 5 个字段组成(见源码中字段定义 L51-L63):

字段类型含义
tableIdTableId事件归属的表标识
beforeRecordData变更前镜像(pre-image)
afterRecordData变更后镜像(post-image)
opOperationType变更操作类型
metaMap<String, String>可选的变更元数据,源码注释举例:MySQL binlog 文件名与位置(file name, pos)

操作类型由枚举 OperationType 预定义为 4 种:INSERTUPDATEREPLACEDELETE。各类型下 before/after 的取值约定如下:

  • Insert:新增数据,before = nullafter = 新数据
  • Delete:删除数据,before = 被删除的数据after = null
  • Update:更新已有数据,before = 变更前数据after = 变更后数据
  • Replace:源文档此处留白,但可从源码确认其语义——replaceEvent(TableId, RecordData after)工厂方法(DataChangeEvent.java#L143-L146)构造出的事件before = nullafter = 完整新数据,即"用一条完整新记录整体替换"的变更。

DataChangeEventSerializable类,不提供公开构造器,而是提供了一组静态工厂方法作为统一入口(L99-L154):insertEventdeleteEventupdateEventreplaceEvent,每种都额外提供了带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自身并不携带表结构信息。这样做可以减小数据变更事件的体积、降低序列化开销,但代价是事件"不自描述"。那么框架如何知道该如何解释这些数据变更?

解决方案是对事件的流转顺序提出强制约定:

  1. 如果某张表对框架来说是新的,那么在任何DataChangeEvent之前必须先发出一个CreateTableEvent
  2. 如果某张表的 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 定位为EventSourceMetadataAccessor的"工厂"。落到源码,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 接口由EventSinkMetadataApplier组成,分别负责写入数据变更事件和应用 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:适配 FlinkSink(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 体系可以概括为:

  1. Event 模型ChangeEvent携带TableId,分为DataChangeEvent(tableId + before + after + op + meta 五字段,四种操作类型)与SchemaChangeEvent(AddColumn / AlterColumnType / CreateTable / DropColumn / RenameColumn 等);核心流转约束是"新表先CreateTableEvent、变 Schema 先SchemaChangeEvent",保证任何DataChangeEvent被处理前 Schema 已知;
  2. Data SourceDataSource是工厂,返回EventSourceProvider(复用 FlinkSource/SourceFunction,读取变更并转成 Event 下发)与MetadataAccessorlistNamespaces/listSchemas/listTables/getTableSchema四个元数据方法),并可选声明元数据列与并行 Schema 变更能力;
  3. Data SinkDataSink对称地返回EventSinkProvider(当前支持 Flink Sink V2 API)与MetadataApplierapplySchemaChange+ 支持类型声明),并可通过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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询