十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Flink CDC

Flink CDC Flink CDC 是什么Flink CDC是基于 Apache Flink 的 CDC 数据采集工具/连接器主要用于实时捕获数据库中的INSERT、UPDATE、DELETE变化并把这些变化交给 Flink 进行清洗、转换、关联和写入下游系统。一句话理解Flink CDC 用 Flink 实时读取数据库变更日志并进行流式处理它常用于把 MySQL、PostgreSQL、Oracle、SQL Server 等数据库中的数据同步到Kafka、Iceberg、Doris、StarRocks、Elasticsearch、Hive 等1. Flink CDC 的基本架构以 MySQL 同步到 Iceberg 为例MySQL | | 读取 binlog v Flink CDC Source | | 生成 INSERT / UPDATE / DELETE 事件 v Flink 流处理 | | 清洗、过滤、去重、字段转换、关联 v Iceberg / Kafka / Doris完整链路可能是MySQL 订单库 | v Flink CDC | -- 全量快照读取已有订单 | -- 增量阶段持续读取 binlog | v Kafka | v Flink 下游任务 | -- Iceberg 原始 CDC 表 -- Doris 当前状态表 -- Redis 实时缓存2. Flink CDC 到底捕获什么假设 MySQL 中有一张订单表CREATETABLEorders(order_idBIGINTPRIMARYKEY,user_idBIGINT,statusVARCHAR(20),amountDECIMAL(18,2),update_timeTIMESTAMP);INSERT新增数据INSERTINTOordersVALUES(1001,2001,CREATED,99.90,NOW());Flink CDC 会产生一条新增事件概念上类似{op:c,order_id:1001,user_id:2001,status:CREATED,amount:99.90}UPDATE修改数据UPDATEordersSETstatusPAIDWHEREorder_id1001;更新事件通常包含修改前和修改后的数据{op:u,before:{order_id:1001,status:CREATED},after:{order_id:1001,status:PAID}}在 Flink Table/SQL 中更新也可能表现为一组UPDATE_BEFORE UPDATE_AFTERchangelog 事件而不一定直接表现为一个 JSON 的before/after对象。DELETE删除数据DELETEFROMordersWHEREorder_id1001;Flink CDC 会产生删除事件{op:d,order_id:1001}3. Flink CDC 的核心流程全量快照 增量读取Flink CDC 通常不是一上来就只读 binlog而是分两个阶段工作。阶段一全量快照 Snapshot先读取数据库当前已经存在的数据MySQL 现有 10 亿条订单 | v Flink CDC 读取历史数据 | v 同步到下游这一阶段解决的是下游系统如何先拥有一份完整的历史数据阶段二增量读取 Incremental全量读取过程中数据库仍然可能继续发生变化。因此 Flink CDC 需要记录一个可靠的 binlog 位点然后从这个位点开始持续读取增量全量数据 从一致 binlog 位点开始的增量数据 完整且连续的数据最终流程是Snapshot 阶段 | | 读取已有数据同时记录日志位点 v Incremental 阶段 | | 持续读取 binlog/WAL v 实时 CDC 流这个过程非常关键。如果全量和增量衔接错误就可能出现全量期间产生的更新被遗漏某些数据重复同步下游最终状态不一致。4. Flink CDC 和 Debezium 有什么关系可以这样理解Debezium主要负责从数据库日志中捕获变更Flink CDC把 CDC 采集能力集成到 Flink 中让数据可以直接进入 Flink 流处理任务。传统架构可能是MySQL - Debezium - Kafka - Flink - Iceberg使用 Flink CDC 后也可以是MySQL - Flink CDC Source - Flink 处理 - Iceberg / Kafka / DorisFlink CDC 底层很多数据库连接器使用了 Debezium 的日志解析能力但它不只是简单的日志采集还提供了Flink Checkpoint全量与增量切换流式转换多表同步并行快照与 Flink Sink 的集成。5. MySQL 使用 Flink CDC 需要什么条件以 MySQL 为例通常需要开启 binlog[mysqld] log-binmysql-bin binlog_formatROW binlog_row_imageFULL server-id1其中log-bin开启 binlogbinlog_formatROW使用行级日志CDC 更容易准确获取行变化binlog_row_imageFULL尽量记录完整行数据server-idCDC 客户端需要使用唯一的复制实例 ID。创建复制权限示意如下CREATEUSERflink%IDENTIFIEDBYpassword;GRANTSELECT,RELOAD,SHOWDATABASES,REPLICATIONSLAVE,REPLICATIONCLIENTON*.*TOflink%;具体权限要根据数据库版本、安全规范和部署方式调整生产环境不建议直接授予过大的权限。6. 一个简单的 Flink CDC SQL 示例下面是一个概念性的 MySQL CDC SourceCREATETABLEorders(order_idBIGINT,user_idBIGINT,statusSTRING,amountDECIMAL(18,2),update_timeTIMESTAMP(3),PRIMARYKEY(order_id)NOTENFORCED)WITH(connectormysql-cdc,hostnamemysql-host,port3306,usernameflink,passwordpassword,database-nametrade_db,table-nameorders,scan.startup.modeinitial);其中scan.startup.mode initial表示先做全量快照再继续读取增量日志。也可能根据场景使用initial 先全量再增量最常见 latest-offset 从当前最新日志位置开始不读取历史 snapshot 只做快照完成后结束具体支持的启动模式和参数要以当前 Flink CDC 版本为准。查询这张 CDC 表时SELECT*FROMorders;得到的不一定只是普通的追加流而可能是包含新增、更新和删除的动态表 Changelog。7. Flink CDC 写入 Iceberg 时怎么处理常见有两种设计。方式一保存原始 CDC 事件order_id | op | status | event_time ---------|----|----------|----------- 1001 | c | CREATED | 10:00:00 1001 | u | PAID | 10:01:00 1001 | u | FINISHED | 10:10:00优点保留完整变更轨迹便于审计可以重新回放方便排查数据异常。适合写入 Iceberg 原始明细表。方式二维护当前状态表order_id | user_id | status | amount | update_time ---------|---------|----------|--------|------------ 1001 | 2001 | FINISHED | 99.90 | 10:10:00下游根据主键进行更新、删除或MERGE。优点查询当前状态方便缺点是需要正确处理重复事件乱序事件删除事件并发写入小文件和 Compaction。实际生产中经常同时保存原始 CDC 事件表 当前最新状态表8. Flink CDC 最需要注意的几个问题8.1 重复消费任务故障恢复后某些事件可能被重新处理。因此需要设计幂等逻辑例如业务主键order_id 事件版本source binlog position / transaction id / update version处理原则可以是同一个主键只接受版本更大的事件不能只依赖 Flink CDC 自动保证业务上绝对不重复。8.2 乱序事件例如先收到10:02 SETTLED后收到10:01 FILLED如果直接覆盖旧事件可能把新状态改回去。因此应使用源端日志位点、递增版本号或可靠的更新时间进行判断。8.3 数据库压力全量快照可能对源库产生读取压力生产上需要考虑是否从只读副本读取快照并行度分片大小业务高峰期限制binlog 保留时间CDC 任务延迟监控。8.4 Schema 变更源库新增字段时需要同步考虑源数据库 Schema - CDC 消息 Schema - Flink 表结构 - Iceberg Schema - Doris / Kafka 下游 SchemaIceberg 支持 Schema Evolution但不代表所有下游都能自动兼容。字段删除、重命名和类型修改尤其需要谨慎。8.5 删除语义数据库删除一条记录后下游不一定应该直接物理删除。可以根据业务选择真正删除 写 is_deleted true 保留删除事件 写入审计表金融、订单和账户类系统通常需要保留删除、撤销或作废事件便于审计和对账。9. Flink CDC 是 Exactly-Once 吗更严谨的说法是Flink CDC 可以结合 Flink Checkpoint、数据库日志位点和下游事务提交机制实现较强的一致性保障但不能简单地说“用了 Flink CDC整个业务链路就绝对 Exactly-Once”。需要分别看数据库日志读取 - Flink 状态和 Checkpoint - Kafka / Iceberg / Doris Sink - 下游业务更新例如Flink 可以从一致的 Checkpoint 恢复CDC 可以从保存的 binlog 位点继续读取Iceberg 可以原子提交 Snapshot但外部 Redis、HTTP 接口或非事务数据库仍然需要自行保证幂等。所以面试时不要只回答“Exactly-Once”而要补充业务主键 版本判断 幂等 Sink Checkpoint 对账机制10. Flink CDC 和定时全量同步的区别对比项定时全量同步Flink CDC数据读取反复扫描整张表读取变化日志实时性通常分钟级或小时级通常秒级或更低取决于链路源库压力可能较大增量阶段压力较小更新删除需要自行比较原生捕获变更事件历史初始化可以直接完成需要 Snapshot 阶段复杂度相对简单需要处理位点、乱序、重复和 Schema 变更适用场景数据量小、实时性要求低实时数仓、实时同步、实时风控面试时可以这样回答Flink CDC 是基于 Flink 的变更数据捕获框架通常通过读取 MySQL binlog、PostgreSQL WAL 等数据库事务日志捕获 INSERT、UPDATE 和 DELETE并以 Changelog 流的形式交给 Flink 处理。它一般先执行一次全量快照再从一致的日志位点切换到增量读取从而避免数据遗漏。Flink CDC 可以把数据实时写入 Kafka、Iceberg、Doris 等系统。生产上需要重点关注全量增量衔接、Checkpoint 恢复、重复和乱序事件、主键幂等、删除语义、数据库压力、binlog 保留时间以及 Schema 变更。对于 Iceberg可以同时保存原始 CDC 事件和当前状态表分别服务于审计回放和实时查询。最简单的记忆方式Flink CDC负责从数据库捕获变化 Flink负责处理和转换变化 Kafka负责传输变化 Iceberg负责保存变化和历史版本 Doris/Redis负责快速查询当前结果
返回列表