新闻详情

新闻详情

首页 / 资讯中心 / 详情

大数据架构深度解析—Flink CDC 实战:MySQL 到数据仓库的实时同步完整方案

发布时间:2026/9/25 20:38:37来源:尧图网络
大数据架构深度解析—Flink CDC 实战:MySQL 到数据仓库的实时同步完整方案
一、为什么是 Flink CDC数据同步的老三样方案问题定时 sqoop 全量抽取T1时效性差大表慢Canal Kafka 手写消费链路长要维护 MQ代码多双写业务代码同时写库和数仓侵入业务极易数据不一致Flink CDC 的核心优势一个引擎吃下捕获 计算 写出不需要中间 MQ。MySQL(Binlog) ──▶ Flink CDC Source ──▶ 清洗/打宽 ──▶ Sink(StarRocks/Doris/Hudi) (Exactly-once)底层用的是 Debezium 捕获 Binlog但 Flink 把它封装成标准的表 API/SQL你写 SQL 就行。二、架构与组件┌──────────┐ Binlog ┌──────────────────────────────┐ ┌──────────────┐ │ MySQL │ ──────────▶ │ Flink 作业 │ ─▶│ StarRocks │ │ (开启binlog)│ │ Source(CDC) → Transform → Sink│ │ (数仓) │ └──────────┘ └──────────────────────────────┘ └──────────────┘ │ Checkpoint (Exactly-once)关键角色Sourcemysql-cdcconnector基于 Debezium捕获 INSERT/UPDATE/DELETE 三种变更Transform普通 Flink SQL做字段映射、过滤、打宽SinkStarRocks/Doris/Hudi 的 connector支持 upsert三、环境准备3.1 MySQL 开启 Binlog# my.cnf server-id 1 log-bin mysql-bin binlog_format ROW # 必须 ROWCDC 才能解析行级变更 binlog_row_image FULL expire_logs_days 7并给 Flink 用的账号授权CREATE USER flink_cdc% IDENTIFIED BY cdc_pass; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO flink_cdc%; FLUSH PRIVILEGES;3.2 下载 connector把这两个 jar 放到 Flink 的lib/目录版本对齐 Flink 1.18flink-sql-connector-mysql-cdc-3.1.0.jar flink-connector-starrocks-1.2.9.jar四、实战MySQL → StarRocks 实时同步4.1 建 MySQL 源表-- Flink SQL CREATE TABLE mysql_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username flink_cdc, password cdc_pass, database-name shop, table-name orders, server-time-zone Asia/Shanghai, scan.startup.mode initial -- 先全量快照再增量 binlog );scan.startup.mode initial是重点Flink CDC 会先全量扫一遍历史数据再无缝切到增量 Binlog不用你手动做存量增量。4.2 建 StarRocks 目标表CREATE TABLE sr_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10,2), status STRING, update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector starrocks, jdbc-url jdbc:mysql://sr-fe:9030, load-url sr-be:8040, database-name shop, table-name orders, username root, password , sink.properties.format json, sink.properties.strip_outer_array true );4.3 同步 SQL核心一行-- 建视图做简单清洗过滤无效订单CREATE VIEW v_orders AS SELECT id, user_id, amount, status, update_time FROM mysql_orders WHERE amount 0; ​ -- 写入数仓 INSERT INTO sr_orders SELECT * FROM v_orders;提交./bin/sql-client.sh -f sync_orders.sql # 或打包成 jar 用 ./bin/flink run 提交生产推荐后者4.4 保障 Exactly-once# flink-conf.yaml execution.checkpointing.interval: 30000 # 30s 一次 checkpoint execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 600000 state.backend: rocksdb # 大状态用 rocksdb state.checkpoints.dir: hdfs:///flink/ckStarRocks Sink 配合 checkpoint 做两阶段提交保证故障恢复后不重复不丢失。五、性能实测环境MySQL 8.04c8g、Flink 1.183×TaskManager每 4 槽、StarRocks 3.x3 BE。用sysbench模拟订单表持续写入指标实测值源端写入速率1000 TPS端到端同步延迟 P500.8 s端到端同步延迟 P991.9 sCheckpoint 耗时1.2 s全量快照阶段吞吐12 万行/min1000 万行表约 14 min 完成故障恢复kill TM从最近 CK 恢复无数据丢失对比传统 T1 批处理最快次日才能查实时性从天降到秒。六、踩坑记录问题现象解决Binlog 不是 ROW 格式CDC 启动报错binlog_formatROW必须改大表全量快照锁表业务写入被阻塞用scan.incremental.snapshot.chunk.size调小分片或低峰期源端 DDL 变更作业挂掉开启schema.change.enabledtrue3.x 支持或手动改表结构后重启Checkpoint 频繁失败状态太大超时换 rocksdb 后端 调大 timeoutSink 报主键冲突重复 upsert确认目标表 PRIMARY KEY 与源一致时区错乱时间字段差 8 小时加server-time-zone且 Flink 设同区反压传回 Source同步延迟飙升优化 Sink 写入批次batch 参数七、总结Flink CDC 用一套 SQL 把捕获→计算→写出做成端到端实时管道告别 T1 和手写 Canalinitial启动模式自动先全量后增量存量数据不用单独处理Checkpoint 两阶段提交保障 Exactly-once故障可恢复生产建议打包 jar 提交而非 sql-client并配 rocksdb 状态后端下一篇实时数据有了消息队列怎么选型周四我们用 Pulsar vs Kafka 把消息中间件讲透。往期回顾Kafka 深度解剖 2消费者组再均衡 Rebalance 全流程StarRocks 实时数仓搭建比 ClickHouse 更适合多维分析的场景Spark 3.5 AQE 调优10 个生产环境案例让作业提速 3-10 倍
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

Atlas 300V 24G实战:用昇腾NPU跑通YOLO视频流目标检测 2026/9/25 21:23:34

Atlas 300V 24G实战:用昇腾NPU跑通YOLO视频流目标检测

刚拿到Atlas 300V 24G这块加速卡时,我也怀疑过它到底算不算运算加速卡——毕竟从外观到驱动安装方式,都跟我平时用的显卡完全不一样。但用它在YOLO模型部署上跑完一个完整的目标检测任务之后,我可以明确地说:它确实是AI推理加速卡…

阅读更多 →
比较器的输出电流限流机制:LM311,LM211 2026/9/25 21:23:27

比较器的输出电流限流机制:LM311,LM211

LM311输出限流模式LM311,LM211,LM111 **AD\Test\2026\September\TestLM211OutputCurrentLimit.SchDoc *** 01 【LM311 比较器输出限流】 一、测量电路 这是比较器LM311它内部的参考电路图, 在它的输出端被称之为隔离的三极管, 它可以接地或者是负电源&a…

阅读更多 →
老旧蓄电池站改造难题?不停机加装蓄电池在线监测方案来了✨ 2026/9/25 21:22:42

老旧蓄电池站改造难题?不停机加装蓄电池在线监测方案来了✨

很多已投运的老旧蓄电池站,都面临同一个棘手痛点: 没有蓄电池在线监测装置🔋 电池单体电压、内阻、温度、剩余容量 SOC 全靠人工定期现场巡检。人工巡检不仅耗费大量人力,更存在明显短板:无法实时捕捉单体劣化、内阻飙…

阅读更多 →
工业 UPS 监控升级|14 路自定义干接点告警,通讯 + 硬件双重守护供电安全 2026/9/25 21:22:36

工业 UPS 监控升级|14 路自定义干接点告警,通讯 + 硬件双重守护供电安全

在机房、工厂配电室、轨道交通、医疗实验室等工业场景中,UPS 作为关键供电保障,一旦出现市电中断、电池低压、过载、逆变器故障等问题,如果告警不及时,极易造成设备停机、数据丢失,带来难以估量的损失。很多传统工业 U…

阅读更多 →
IronClaw Review Readiness:以证据驱动的 PR 合并就绪度看板 2026/9/25 21:22:36

IronClaw Review Readiness:以证据驱动的 PR 合并就绪度看板

人工智能AI 应用交互助手AI Agent 【免费下载链接】ironclaw IronClaw is an Agent OS focused on privacy, security and extensibility 项目地址: https://gitcode.com/gh_mirrors/iro/ironclaw 点击查看 免费下载 本文围绕 IronClaw 仓库中的 review-readiness …

阅读更多 →
UltraData-RL-2609 许可证与数据合规须知:Apache 2.0、上游许可与新基准去污完整指南 2026/9/25 21:22:30

UltraData-RL-2609 许可证与数据合规须知:Apache 2.0、上游许可与新基准去污完整指南

UltraData-RL-2609 许可证与数据合规须知:Apache 2.0、上游许可与新基准去污完整指南 【免费下载链接】UltraData-RL-2609 项目地址: https://ai.gitcode.com/OpenBMB/UltraData-RL-2609 UltraData-RL-2609 是 OpenBMB 开源社区发布的可验证奖励强化学习&am…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉