资讯动态

Flink CDC 3.5.0 从理论到实践 —— 第 1 章 Flink CDC 概览与发展演进

发布时间:2026/9/30 9:29:16 来源:尧图企业网站定制
Flink CDC 3.5.0 从理论到实践 —— 第 1 章 Flink CDC 概览与发展演进课程定位本系列教程以MySQL 为唯一数据源Sink 覆盖Doris / Paimon / Kafka三大目标从原理到生产落地全链路实战。版本基线Flink CDC 3.5.0 Flink 1.20.x MySQL 8.0/8.4 Doris 4.1 Paimon 1.4.2 Kafka 3.x章节导读1.1 什么是 CDC1.2 主流 CDC 方案对比1.3 Flink CDC 版本演进1.4 Flink CDC 3.5.0 六大核心特性1.5 Flink CDC vs 传统离线数仓 ETL1.6 适用场景与典型架构1.7 本章小结1.1 什么是 CDCCDCChange Data Capture变更数据捕获是一种识别并捕获数据库中数据插入、更新、删除的技术使下游系统能够实时感知源库的数据变化从而驱动数据同步、实时数仓、缓存更新、微服务事件分发等场景。1.1.1 为什么需要 CDC在数据工程演进中传统离线 ETL 通常依赖定时全量拉取或按更新时间增量查询存在几类典型痛点痛点传统方案表现CDC 解决方式时效性差T1 调度数据延迟小时级甚至天级准实时秒级延迟业务变更即时感知对源库压力大全表扫描高峰期影响在线业务增量读取 Binlog无锁、轻量增量识别不准依赖update_time字段并发更新易丢失基于 Binlog 事务日志准确无误无法捕获删除软删除需额外字段硬删除直接丢失Binlog ROW 模式直接捕获 DELETE数据一致性难保证抽取窗口跨事务可能出现幻读事务边界严格对齐Exactly-Once1.1.2 CDC 的两种实现路径CDC 技术实现主要分两类基于查询Query-based周期性执行SELECT拉取变更如基于update_time last_pull_time优点实现简单不依赖数据库日志能力缺点时效性差、无法捕获硬删除、对源库压力较大。基于日志Log-based解析数据库事务日志MySQL Binlog / PostgreSQL WAL / Oracle Redo Log优点准实时、完整捕获 DML含删除、对源库几乎无侵入缺点需要数据库开启日志能力、解析逻辑复杂。Flink CDC 采用基于日志的方式并以 MySQL Binlog 为最主流的数据源。1.1.3 MySQL Binlog 是 CDC 的基石MySQL 的 Binlog二进制日志记录所有改变表数据的 DDL 和 DML 事件是 CDC 的核心数据来源。要让 Flink CDC 正常工作MySQL 端必须满足以下配置参数推荐值说明log_binON开启 BinlogMySQL 8.0 默认开启binlog_formatROW必须为 ROW记录每行变更前后镜像binlog_row_imageFULL记录变更前后完整列值默认 FULLserver-id全局唯一CDC 客户端伪装为 SlaveID 必须唯一Binary Log Transaction CompressionOFFMySQL 8.0.20 引入需关闭最小权限账号CREATEUSERcdc_user%IDENTIFIEDBYCdc.2026;GRANTSELECT,SHOWDATABASES,REPLICATIONSLAVE,REPLICATIONCLIENTON*.*TOcdc_user%;FLUSHPRIVILEGES;小提示REPLICATION SLAVE和REPLICATION CLIENT是读取 Binlog 的必需权限缺失会导致 CDC 启动报错Could not find first log file name in binary log index file。1.2 主流 CDC 方案对比业界主流的开源 CDC 方案有 Canal、Debezium、Flink CDC 三种各有侧重。1.2.1 三大方案速览维度CanalDebeziumFlink CDC出身阿里开源RedHat 开源现属 CNCFApache Flink 顶级项目核心定位MySQL Binlog → Kafka通用 CDCMySQL/PG/Oracle/MongoDB…Flink 生态内端到端 CDC部署形态独立 ServerKafka Connect ConnectorFlink 作业全量初始化支持需配合 Adapter支持snapshot支持增量快照算法无锁下游消费主要写 Kafka需二次消费主要写 Kafka需二次消费直接对接 Flink SinkDoris/Paimon/Kafka/JDBC…计算能力无纯采集无纯采集有流批一体 SQL TransformExactly-Once需下游配合依赖 Kafka 语义Flink Checkpoint 保证生态语言JavaJavaJava/Python/SQL/YAML1.2.2 架构对比图Canal / Debezium 模式采集 二次消费MySQL ──Binlog── [Canal/Debezium] ── Kafka ── [Flink/Spark 消费] ── Doris/Paimon 采集层 缓存层 计算层 存储层特点链路长需维护 Canal/Debezium Kafka 两套组件二次开发多。Flink CDC 模式采集与计算一体MySQL ──Binlog── [Flink CDC Source] ── [Flink Transform] ── [Doris/Paimon/Kafka Sink] 采集计算一体 计算层 存储层特点链路短一个 Flink 作业完成采集、转换、写入无中间组件。1.2.3 选型建议场景推荐方案理由已有 Kafka 生态多下游订阅同一变更Canal / Debezium解耦采集与消费一次采集多端订阅端到端实时同步到 OLAP / 数据湖Flink CDC链路最短无需 Kafka 中转Exactly-Once 强保证需要在采集过程中做 ETL过滤、计算列、路由Flink CDC内置 Transform 能力避免二次作业异构源Oracle/MongoDB/SQL ServerDebezium覆盖数据库最广纯 MySQL 同步到 Doris/PaimonFlink CDC本教程主线本教程聚焦的MySQL → Doris/Paimon/Kafka场景正是 Flink CDC 的核心优势战场。1.3 Flink CDC 版本演进理解版本演进有助于掌握 Flink CDC 的设计哲学与能力边界。1.3.1 演进时间线版本发布时间关键能力架构特征1.02020 年MySQL CDC 初版基于 Debezium 嵌入式引擎全量阶段加锁FLUSH TABLES WITH READ LOCK1.22020 末新增 Postgres / Oracle / MongoDB Connector多源扩展2.02021 年增量快照算法首次引入无锁、并发读取、断点续传2.2-2.42022-2023Schema Evolution 雏形、更多 Source算法优化与稳定性提升3.02024 年Pipeline 架构YAML 定义、整库同步、自动建表从单表 SQL 走向整库 YAML 管道3.1-3.42024-2025Transform 模块、多 SinkDoris/Paimon/Kafka/JDBC、Schema Evolution 增强端到端能力成熟3.5.02025 年稳定性优化、参数体系完善、多 Sink 类型映射完整生产级稳定版本1.3.2 两次关键跃迁跃迁一1.x → 2.0增量快照算法Flink CDC 1.x 全量初始化采用 Debezium 原生的INITIAL模式需要FLUSH TABLES WITH READ LOCK锁库对在线业务影响大且全量阶段单线程读取大表初始化耗时。2.0 引入增量快照算法Incremental Snapshot Algorithm核心思想将全表按主键范围切分为多个Chunk多线程并发读取不同 Chunk无锁基于SELECT ... WHERE pk BETWEEN ? AND ?全量进行的同时持续记录 Binlog 位点全量完成后将 Binlog 位点对齐无缝切换为增量消费。带来的收益无锁、并发、断点续传。跃迁二2.x → 3.0Pipeline 架构2.x 时代 Flink CDC 以SQL / DataStream API为主单表一对一同步写整库需要写大量 SQL维护成本高。3.0 引入PipelineYAML 定义架构source:type:mysqlhostname:127.0.0.1tables:adb.\.*,bdb.user_table_[0-9]sink:type:dorisfenodes:127.0.0.1:8030pipeline:name:MySQL to Doris Pipelineparallelism:4一段 YAML 即可完成整库、多表、自动建表、DDL 同步的端到端同步是 Flink CDC 走向生产易用化的里程碑。1.3.3 3.5.0 的定位3.5.0 是 3.x 大版本线上的稳定演进版本特性集已经成熟参数体系完善是当前推荐的生产部署版本。本教程即以 3.5.0 为基线展开。1.4 Flink CDC 3.5.0 六大核心特性根据官方文档 https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/ 的描述Flink CDC 3.5.0 提供六大核心能力也是本教程后续章节的展开主线。1.4.1 增量快照算法Incremental Snapshot官方原文Flink CDC supports distributed scanning of historical data of database and then automatically switches to change data capturing. The switch uses the incremental snapshot algorithm which ensure the switch action does not lock the database.一句话理解全量历史数据并发扫描无缝切换到增量 Binlog 消费全程无锁。技术要点Chunk 切分按主键范围将表切分为多个分片默认每片 8096 行并发读取无锁设计全量阶段不持有数据库锁不影响在线业务读写无缝切换全量完成时自动对齐 Binlog 位点切换为增量消费不丢不重断点续传Checkpoint 记录已完成 Chunk故障恢复后从断点继续。参数速览参数默认值说明scan.incremental.snapshot.chunk.size8096每个 Chunk 的最大行数scan.snapshot.fetch.size1024单次 SELECT 拉取行数scan.startup.modeinitial启动模式详见第 4 章 MySQL Source 全方位解析。1.4.2 Schema Evolution表结构演进同步官方原文Flink CDC has the ability of automatically creating downstream table using the inferred table structure based on upstream table, and applying upstream DDL to downstream systems during change data capturing.一句话理解上游ALTER TABLE自动同步到下游 Doris / Paimon并自动建表。能力边界自动建表根据上游表结构推断下游 DDL自动创建目标表DDL 同步上游加列、改类型、改长度等ALTER操作自动应用到下游schema-change.enabled默认true可关闭以禁止 DDL 同步。三种 Sink 的支持差异Sink自动建表加列改类型改名删列Doris✅✅部分❌❌Paimon✅✅✅✅✅KafkaN/A字段变更需结合 Schema Registry 或下游容错详见第 8 章 整库同步与 Schema Evolution。1.4.3 Streaming Pipeline流式管道官方原文Flink CDC jobs run in streaming mode by default, providing sub-second end-to-end latency in real-time binlog synchronization scenarios, effectively ensuring data freshness for downstream businesses.一句话理解用一段 YAML 定义端到端同步管道默认流式运行秒级延迟。典型 Pipeline 结构source:# 源type:mysqlhostname:127.0.0.1tables:app.\.*sink:# 目标type:dorisfenodes:127.0.0.1:8030pipeline:# 管道控制name:MySQL to Dorisparallelism:4transform:# 转换可选-source-table:app.ordersprojection:id, name, pricefilter:price 100收益配置即代码整库同步不再需要写几十段 SQL。1.4.4 Data Transformation数据转换官方原文Flink CDC supports data transform operations of ETL, including column projection, computed column, filter expression and classical scalar functions.一句话理解在采集过程中直接做轻量 ETL无需额外 Flink 作业。四种 Transform 能力列投影Projection选择性同步字段如id, name, price计算列Computed Column派生字段如dt DATE_FORMAT(create_time, yyyy-MM-dd)过滤表达式Filter条件性同步如price 100 AND status paid标量函数Scalar FunctionsFlink SQL 内置函数可直接使用。详见第 9 章 Data Transformation 与 ETL。1.4.5 Full Database Sync整库同步官方原文Flink CDC supports synchronizing all tables of source database instance to downstream in one job by configuring the captured database list and table list.一句话理解一个作业同步整个数据库的所有匹配表新增表自动接入。关键参数tables表名正则支持db.\.*整库、db.table_[0-9]分表、[app|web].order_\.*多库schema-change.enabled是否同步 DDL默认 true自动建表下游表不存在时自动创建。示例同步 app 和 web 库的所有以order_开头的表到 Dorissource:type:mysqltables:[app|web].order_\\.*sink:type:dorisfenodes:127.0.0.1:8030详见第 8 章 整库同步与 Schema Evolution。1.4.6 Exactly-Once Semantics精确一次语义官方原文Flink CDC supports reading database historical data and continues to read CDC events with exactly-once processing, even after job failures.一句话理解作业故障重启后不丢数据、不重复数据。实现机制Source 端增量快照算法基于 Flink Checkpoint 记录已读 Chunk 与 Binlog 位点Sink 端Doris基于 Unique 主键模型去重或 Stream Load 两阶段提交Paimon基于主键的 Merge 去重Kafka基于事务或幂等 Key整体Flink Checkpoint 提供 End-to-End Exactly-Once 保障。前置条件必须开启 Checkpoint推荐 EXACTLY_ONCE 模式Sink 必须支持幂等或事务。详见第 10 章 作业运维与监控。1.4.7 特性矩阵速览特性解决的问题关键参数详见章节增量快照无锁全量、并发、断点续传chunk.size/startup.mode第 4 章Schema EvolutionDDL 自动同步schema-change.enabled第 8 章Streaming PipelineYAML 管道定义pipeline.name/parallelism第 5-7 章Data Transformation采集过程 ETLprojection/filter第 9 章Full Database Sync整库多表一次同步tables正则第 8 章Exactly-Once故障不丢不重Checkpoint 配置第 10 章1.5 Flink CDC vs 传统离线数仓 ETL理解 Flink CDC 与传统离线 ETL 的本质差异有助于在架构选型时做出正确决策。1.5.1 架构对比传统离线 ETL如 Spark/Hive 调度MySQL ──(Sqoop/DataX 全量/增量)── HDFS ── Hive ODS ── DWD ── DWS ── ADS ↑ T1 调度数据延迟天级Flink CDC 实时数仓MySQL ──Binlog── Flink CDC ── Doris ODS ──(Flink/Spark)── DWD/DWS ── ADS ↑ 秒级延迟流式持续同步1.5.2 关键差异对照维度传统离线 ETLFlink CDC 实时同步时效性T1小时到天级秒级亚秒级抽取方式全表扫描或基于时间字段增量Binlog 事务日志增量快照无锁对源库压力全表扫描高峰期压力大Binlog 读取几乎无侵入能否捕获删除不能或需软删除字段能Binlog ROW 模式一致性保证难以跨事务边界Exactly-Once事务对齐运维复杂度调度系统 多个作业单个 Flink 作业资源占用离线集群空闲时段高负载流式常驻资源持续占用适用场景历史分析、报表、回流实时大屏、风控、即席查询、事件分发1.5.3 不是替代而是互补Flink CDC并非要完全取代传统离线 ETL二者定位不同实时链路Flink CDC低延迟、高时效服务在线业务、实时大屏、风控预警离线链路Spark/Hive大批量、复杂计算、历史回算、模型训练。生产实践通常是Lambda 架构或Kappa 架构实时层用 Flink CDC 写入 Doris/Paimon 供即席查询离线层基于 Paimon 数据湖做 T1 复杂加工。1.6 适用场景与典型架构1.6.1 四大典型场景场景一实时数仓OLAP 即席查询MySQL ── Flink CDC ── Doris ── BI 报表 / 即席查询特点业务库变更秒级到 Doris分析师即席查询最新数据。本教程第 5 章详解。场景二数据湖入湖ODS 贴源MySQL ── Flink CDC ── PaimonHDFS/S3── Spark/Hive 离线分析特点原始数据落湖保留全量历史供离线数仓分层加工。本教程第 6 章详解。场景三事件分发微服务解耦MySQL ── Flink CDC ── Kafka ── 多下游订阅缓存更新、异步任务、搜索索引特点一次采集多端订阅业务库变更驱动下游微服务。本教程第 7 章详解。场景四一体化实时数仓综合┌── Doris 实时 OLAP MySQL ── Flink CDC ──────┼── Paimon 离线 ODS └── Kafka 事件分发特点一个 CDC Source分流到三个目标覆盖实时/离线/事件三类下游。本教程第 12 章综合实战。1.6.2 与大数据生态的集成Flink CDC 作为 Apache Flink 生态的一部分天然与大数据组件无缝集成下游组件集成方式典型用途DorisDoris Pipeline Sink实时 OLAP 查询PaimonPaimon Pipeline Sink数据湖存储KafkaKafka Pipeline Sink事件分发、缓冲解耦Hive通过 Paimon Catalog 互通离线数仓Spark通过 Paimon/Iceberg 互通批处理加工Hudi / Iceberg类似 Paimon支持 Sink数据湖多选一1.6.3 本教程的架构主线本教程围绕MySQL → Doris / Paimon / Kafka三条主线展开┌──────────────┐ │ Doris │ ← 第 5 章实时 OLAP └──────────────┘ ▲ ┌──────────┐ Flink CDC │ ┌──────────────┐ │ MySQL │ ────────────┼─│ Paimon │ ← 第 6 章数据湖入湖 │ (业务库) │ │ └──────────────┘ └──────────┘ │ ▼ ┌──────────────┐ │ Kafka │ ← 第 7 章事件分发 └──────────────┘三条链路共用同一 MySQL Source通过 Flink CDC 3.5.0 的 Pipeline YAML 统一编排是车联网、电商、金融等行业的典型实时数据架构。1.7 本章小结本章作为全系列教程的开篇回答了四个核心问题CDC 是什么基于数据库事务日志如 MySQL Binlog捕获数据变更的技术解决传统 ETL 时效性差、压力大、无法捕获删除的痛点。为什么选 Flink CDC相比 Canal / Debezium 的采集 二次消费模式Flink CDC 采集与计算一体链路最短Exactly-Once 强保证。Flink CDC 演进到哪了1.x 全量加锁 → 2.0 增量快照无锁 → 3.0 Pipeline YAML 整库同步 → 3.5.0 生产稳定版。3.5.0 有哪些核心能力六大特性——增量快照、Schema Evolution、Streaming Pipeline、Data Transformation、Full Database Sync、Exactly-Once。下一章预告第 2 章《Flink CDC 核心原理深入》将剖析增量快照算法的 Chunk 切分、Binlog 位点对齐、Exactly-Once 落地机制帮助你理解无锁和不丢不重背后的实现原理。参考资料Flink CDC 3.5.0 官方文档https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/Flink CDC GitHubhttps://github.com/apache/flink-cdcMySQL Binlog 文档https://dev.mysql.com/doc/refman/8.0/en/binary-log.htmlDebezium MySQL Connectorhttps://debezium.io/documentation/reference/stable/connectors/mysql.htmlDoris Flink Connectorhttps://doris.apache.org/docs/dev/ecosystem/flink-doris-connector/Apache Paimon 文档https://paimon.apache.org/docs/master/本文是《Flink CDC 3.5.0 从理论到实践》系列教程的第 1 章后续章节将持续更新欢迎关注收藏。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑