Flink CDC 组件
最近在梳理实时数据链路时,重新整理了一遍 Flink CDC 的相关知识。这类"把数据库变更实时搬出来"的需求几乎每个数据团队都会遇到,值得单独写一篇记录。
背景
做数据同步这件事,最朴素的方案是定时跑批:每天凌晨全量拉一遍业务库,写进数仓。这种方式实现简单,但延迟以小时甚至天计,数据量大了以后对源库的压力也不小。业务方一旦提出"报表要看实时的",跑批方案就顶不住了。CDC 正是为了解决这个问题而出现的技术路线,而 Flink CDC 则是把 CDC 能力和流计算引擎结合得比较彻底的一个实现。
Flink CDC(Change Data Capture)是 Apache Flink 生态中的一个重要组件,用于实时捕获数据库中的数据变更,并将这些变更转换为数据流进行处理。在现代数据架构中,CDC 技术被广泛用于构建实时数据同步、实时数据仓库以及事件驱动系统。
与传统的数据同步方式(如定时全量同步)相比,CDC 能够通过读取数据库的 Binlog / WAL 等变更日志 来捕获数据变化,从而实现低延迟、低侵入的数据同步。
什么是 CDC
CDC(Change Data Capture)即 变更数据捕获,它的核心思想是:
当数据库中的数据发生 INSERT / UPDATE / DELETE 操作时,将这些变化记录下来,并实时传递到下游系统。
这里可以稍微展开一下它的实现机制。CDC 大体分两类:一类是基于查询的,靠轮询时间戳或自增主键字段来发现新数据,实现简单但抓不到 DELETE,也容易漏掉两次轮询之间被多次修改的记录;另一类是基于日志的,直接消费数据库为主从复制而维护的变更日志——MySQL 的 Binlog、PostgreSQL 的 WAL 都属于这类。日志里天然记录了每一行数据变更前后的完整状态,所以基于日志的 CDC 能拿到全量的增删改事件,且对源库几乎没有额外查询压力。Flink CDC 走的就是日志这条路。
在实际的数据架构中,CDC 常用于:
- 数据库 → 数据仓库(实时数仓)
- 数据库 → Kafka(事件流)
- 数据库 → 搜索引擎(如 Elasticsearch)
- 数据库 → 缓存系统(如 Redis)
因此,CDC 是现代 实时数据平台(Real-time Data Platform) 的关键基础能力之一。
Flink CDC 组件的工作原理是通过捕获源数据系统的变更日志,将其转化为数据流并进行实时处理。这个过程可以在不影响源系统的情况下进行,因为CDC组件只是读取源系统的日志,而不会对源系统进行写操作。从数据库的视角看,CDC 客户端和一个普通的从库没有本质区别——它伪装成一个复制客户端,向主库订阅变更日志,再把日志事件解析成结构化的变更记录交给 Flink 处理。
Flink CDC组件的优点在于它可以支持多种数据源,包括关系型数据库(如MySQL,Oracle等)和NoSQL数据库(如MongoDB,Cassandra等)。此外,它还能够支持多种数据格式,如JSON,CSV等。
另一个容易被低估的优点是它和 Flink 本身的整合深度。变更流进入 Flink 之后,就是一条普通的数据流,可以直接套用 Flink 的窗口计算、多流 Join、状态管理和 checkpoint 容错机制。也就是说,"捕获变更"和"处理变更"在同一个引擎里完成,不需要再单独维护一套 Kafka Connect 集群做中转,链路更短,运维对象也更少。
使用Flink CDC组件可以实现实时的数据同步和数据流处理。例如,我们可以使用CDC组件将MySQL数据库中的数据同步到Kafka,然后使用Flink进行实时数据处理。这样就可以实现数据的实时同步和处理,从而提高数据的效率和准确性。
Flink CDC组件
从整体上看,一条 Flink CDC 链路的分工大致是:Source 负责"进",Sink 负责"出",中间由 Flink 的算子和状态机制完成计算。Flink CDC组件内置了以下几个组件:
- Source:该组件用于从数据源中读取数据变更日志,并将其转换为Flink数据流。通常一个 CDC Source 会先对存量数据做一次快照读取,再无缝衔接到增量日志消费,这样下游拿到的就是"全量 + 增量"的完整数据。
- Debezium Connector:该组件是一个开源的CDC工具,可以连接多种数据源(如MySQL、PostgreSQL、MongoDB等)并捕获数据变更日志。Flink CDC 在底层复用了 Debezium 的日志解析能力,把它以内嵌方式跑在 Flink 作业里,省去了独立部署 Debezium 服务的成本。
- Sink:该组件用于将Flink数据流写入目标数据源中,例如Kafka、HDFS、Elasticsearch等。Sink 配合 Flink 的 checkpoint 机制,可以在作业失败重启时避免数据丢失。
- State:该组件用于处理有状态的数据流,例如,如果需要将两个数据流进行合并,则需要使用State组件来存储中间状态。CDC 场景下 State 尤其重要——比如维表关联、按主键去重、把 UPDATE 事件还原成最新镜像,都依赖状态存储。
- Table API:该组件提供了一个SQL-like的API,可以方便地进行数据流处理和查询。对于大多数同步类需求,用 Flink SQL 声明一张 CDC 源表和一张目标表,写一条 INSERT INTO 就能跑通整条链路,不需要写 Java 代码。
踩坑与注意
实际使用中有几个点值得提前留意:
-
源库的日志配置要先确认。以 MySQL 为例,Binlog 需要开启且格式为 ROW,否则拿不到行级别的变更前后镜像;连接账号也需要具备读取 Binlog 相关的复制权限。
-
日志保留时间不能太短。如果作业停了一段时间再恢复,而这期间的 Binlog 已经被数据库清理掉,作业就找不到断点位置,只能重新做全量快照。
-
下游要能处理 UPDATE 和 DELETE。CDC 流是一条包含回撤语义的变更流,如果 Sink 端(比如只支持追加写入的存储)消化不了删除和更新事件,链路设计上就要额外处理。
CDC 作业本质上是一个长期运行的流作业,checkpoint 一定要配置好。没有 checkpoint 的 CDC 作业一旦失败,既丢断点又可能丢数据。
小结
CDC 解决的是"数据库里的变化如何低延迟、低侵入地流出来"这个问题,基于日志的实现是目前的主流路线。Flink CDC 的价值在于把日志捕获(借助 Debezium)和流式计算(Flink 引擎本身)收拢到了一个框架里:Source 读变更、Sink 写目标、State 撑起有状态计算、Table API 降低使用门槛。对于想搭实时数仓或者做异构数据同步的团队,它是一个链路短、组件少、值得优先考虑的选项。
评论 / COMMENTS