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

资讯详情

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

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

大数据架构深度解析—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 倍
返回列表