ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink CDC实时数据同步完整指南:3步跑通MySQL到Kafka整库同步链路

Flink CDC实时数据同步完整指南:3步跑通MySQL到Kafka整库同步链路 Flink CDC实时数据同步完整指南3步跑通MySQL到Kafka整库同步链路【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdcFlink CDC 是构建在 Apache Flink 之上的实时数据集成工具它通过捕获数据库变更日志Change Data CaptureCDC提供整库同步、模式演进Schema Evolution与数据转换能力。本文以MySQL 到 Kafka 实时同步这条经典链路为例从架构选型、快速上手讲到了生产加固帮你用一份 YAML 文件把业务库同步延迟从小时级压到秒级。一、方案概览整条链路分三层源库负责暴露变更Flink 计算层负责读、合并与转换目标端负责落地。组件职责技术选型MySQL CDC 源读取全量快照与 binlog输出统一变更事件Flink Source API Debezium 引擎源码见 flink-connector-mysql-cdcPipeline 运行时组装源与汇执行路由、转换与模式演进Flink DataStream 运行时flink-cdc-runtimeKafka Pipeline Sink将事件写入 topic支持按主键哈希分区flink-cdc-pipeline-connector-kafkaFlink 集群Checkpoint 容错保证同步不丢数Flink 1.20 / 2.x可部署于 Standalone/YARN/K8s二、场景与价值先看一个典型场景电商的订单、库存表白天持续变化但分析平台的报表靠每小时的定时 ETL 刷新运营看到的 GMV 始终落后 1~2 小时。这类需求用批量同步很难做好原因在于全量抽取会长时间占用业务库 IO 并锁读资源增量补偿逻辑复杂窗口内极易出现重复或丢失表结构变更又会让按固定字段写的作业直接失败。Flink CDC 的应对方式是把快照 增量做成一条无缝衔接的流水线场景挑战Flink CDC 的应对存量数据量大亿级行全量抽取压垮业务库增量快照Incremental Snapshot按主键分块并行读取不锁表快照期间仍在持续写入同步窗口内数据重复/丢失分块读取与 binlog 回填按位点合并配合幂等写入保证不重不漏上游频繁 DDL下游表结构与事件不匹配模式演进机制自动向下游下发加列、改列等事件表多、同步范围广逐表写作业成本高正则表达式选表一份 YAML 完成整库同步三、快速上手3步跑通第一条实时同步链路1. 准备 MySQL 端开启 binlog 并创建最小权限用户# my.cnf确保以 ROW 格式记录完整行镜像 [mysqld] server-id 1 log-bin mysql-bin binlog_format ROW binlog_row_image FULL-- 只授予 CDC 所需的最小权限 CREATE USER flinkuser% IDENTIFIED BY your_password; GRANT SELECT, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO flinkuser%;CDC 本质上是伪装成一个从库去拉 binlog所以必须有REPLICATION SLAVE/CLIENT权限。2. 准备运行环境启动 Flink 集群并放入连接器 JAR# 解压并启动 Flink 集群开启 checkpoint每 3 秒一次 tar -zxvf flink-2.2.0-bin-scala_2.12.tgz ./bin/start-cluster.sh# 将以下 3 个 jar 放入 Flink CDC 发行包的 lib 目录非 Flink lib cp flink-cdc-pipeline-connector-mysql-*.jar \ flink-cdc-pipeline-connector-kafka-*.jar \ mysql-connector-java-8.0.27.jar $FLINK_CDC_HOME/lib/MySQL 驱动因 GPLv2 协议不在官方预打包范围内需要随作业一起提供checkpoint 是增量快照正确性的前提务必开启。3. 定义 Pipeline 并提交作业source: type: mysql hostname: 127.0.0.1 port: 3306 username: flinkuser password: your_password tables: app_db.\.* # 正则选表整库同步 server-id: 5400-5404 # 每个作业独占一段禁止复用 server-time-zone: UTC sink: type: kafka properties.bootstrap.servers: 127.0.0.1:9092 topic: cdc-mysql-kafka pipeline: name: MySQL to Kafka Pipeline parallelism: 2bash $FLINK_CDC_HOME/bin/flink-cdc.sh mysql-to-kafka.yaml提交成功后Flink Web UI 中可以看到作业先跑完快照阶段、再切换到增量阶段向app_db任意表插入一行Kafka topic 中几秒内就能消费到对应的op: c事件。四、核心实现解析事件在源端内部经历三个阶段SnapshotSplitReader并行读取分块快照 →BinlogSplitReader单读 binlog 并回填快照期间的变更 → 两者合并后经 MySqlRecordEmitter 反序列化为统一事件模型flink-cdc-common 中的DataChangeEvent/SchemaChangeEvent再下发给 Sink。分片策略由 MySqlHybridSplitAssigner 负责把每张表按主键范围切成 chunk并把唯一的 binlog split 交给一个 reader。RecordEmitter 的核心分发逻辑节选protected void processElement(SourceRecord element, SourceOutputT output, MySqlSplitState splitState) throws Exception { if (RecordUtils.isWatermarkEvent(element)) { // 高水位写入 split 状态界定 binlog 回填窗口 splitState.asSnapshotSplitState().setHighWatermark(watermark); } else if (RecordUtils.isSchemaChangeEvent(element)) { // DDL 先落 split 状态再下发供模式演进 splitState.asBinlogSplitState().recordSchema(id, tableChange); emitElement(element, output); } else if (RecordUtils.isDataChangeRecord(element)) { updateStartingOffsetForSplit(splitState, element); // 推进位点供 checkpoint emitElement(element, output); } }这段代码的职责是区分四类事件水位、DDL、DML、心跳DML 每处理一条就推进位点checkpoint 时以位点落盘这正是作业重启不丢数的基础DDL 则被保留并发往下游支撑表结构同步。关键参数参数默认值说明scan.startup.modeinitial首次启动先快照再增量latest-offset只同步新变更server-id5400-6400 随机建议显式指定不重叠的区间避免与其他复制冲突scan.incremental.snapshot.chunk.size8096快照分块行数决定快照并行度粒度scan.snapshot.fetch.size1024快照阶段单次拉取行数调大降低 RTTschema.change.behaviorlenient模式变更策略exception / evolve / try_evolve / lenient / ignore五、生产加固并行度pipeline.parallelism决定快照阶段的并发上限binlog 阶段恒为单 reader因此并行度收益主要在快照期建议与 TaskManager 数量匹配2~8 常见。状态与内存# source 侧配置降低 JM 内存占用、释放空闲 reader scan.incremental.snapshot.metadata.release.enabled: true scan.incremental.close-idle-reader.enabled: true第一个选项在 binlog 阶段释放 JobManager 中缓存的分片元数据快照表多时能显著降低 JM 内存第二个让快照读完的 reader 及时回收。Flink 侧建议state.backend: rocksdb并按快照数据量评估 TaskManager 堆内存。监控指标均暴露为 Flink Metrics可接 Prometheus指标名类型含义numSnapshotSplitsFinished/numSnapshotSplitsRemainingGauge快照分片完成/剩余数量估算全量进度isSnapshotting/isStreamReadingGauge表当前处于快照还是增量阶段snapshotStartTime/snapshotEndTimeGauge快照阶段起止时间用于核算全量耗时currentEmitEventTimeLagGauge事件时间口径的端到端同步延迟fetchDelay累积器binlog 拉取延迟衡量源端读取是否跟得上写入六、排障速查问题原因解决作业反复重启日志报 server-id 被占用与其他 CDC/复制任务 ID 冲突每个作业分配独立server-id区间启动报 offset/position not foundbinlog 已被清理调大binlog_expire_logs_seconds或scan.startup.mode: snapshot重做快照快照阶段慢且业务库负载高并行度为 1 或分块粒度过大调大pipeline.parallelism配合chunk.size/fetch.size上游 DDL 后作业失败schema.change.behavior为exception改为evolve下游自动变更或lenient失败不中断JobManager OOM分片元数据全部驻留 JM开启scan.incremental.snapshot.metadata.release.enabled同步静默停滞源表无变更位点不推进确认heartbeat.interval默认 30s生效检查下游消费七、实践建议与案例某零售企业的订单库20 张表、快照约 5 亿行按本文链路接入 Kafka Doris 后报表数据延迟从2 小时降到 5 秒以内快照阶段在 1.5 小时内跑完且未影响业务高峰上游执行加列 DDL 后下游 Doris 表由模式演进自动跟进未出现人工改表。分阶段实施建议先把生产库只读副本复制到预发环境跑 POC用行数与抽样校验checksum对账先上线 2~3 张核心表观察一周的currentEmitEventTimeLag与 checkpoint 耗时确认稳态再扩展到整库正则选表并定期保存 savepoint确保可回滚升级建立告警作业重启次数、checkpoint 连续失败、事件时间延迟超阈值三者必配。八、展望数据源覆盖持续扩大仓库内已提供 Oracle、OceanBase、SQL Server、PostgreSQL 等源连接器异构源可复用同一套 Pipeline 抽象AI 参与数据转换pipeline-model 模块已支持在管道中调用大模型字段映射与语义转换正从手写 UDF 走向声明式配置湖仓一体落地Iceberg、Paimon、Hudi 等 Sink 均在 pipeline 连接器目录 中实时入湖链路将越来越标准化。延伸阅读MySQL Pipeline 连接器文档Kafka Pipeline 连接器文档MySQL CDC Source 连接器文档Data Pipeline 核心概念Schema Evolution 模式演进说明QuickstartMySQL to Kafka 教程【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进