恒美微站
首页
关于我们
建站服务
主题模板
案例展示
资讯中心
联系我们
5小时搭建实时湖仓:Flink CDC同步MySQL到数据湖实战
首页
资讯中心
/
5小时搭建实时湖仓:Flink CDC同步MySQL到数据湖实战
5小时搭建实时湖仓:Flink CDC同步MySQL到数据湖实战
发布时间:2026/10/1 15:48:38
简介5小时玩转阿里云实时计算Flink实时湖仓课程的配套原始业务数据脚本面向大数据与实时计算学习者适合正在学习阿里云Flink实时湖仓搭建、希望获得可运行示例数据的开发者。资源包共含4个文件由两个SQL脚本和两个TXT说明组成SQL脚本提供模拟业务数据及建表语句TXT文件则涵盖ECS环境安装JDK、ZooKeeper、Kafka、MySQL等组件的命令参考便于快速还原课程实验环境。整个压缩包仅623KB轻量易用不占用存储空间。目前已有241人学习下载适合作为Flink实时湖仓入门与实操训练的辅助材料。通过对照脚本与数据读者可以跳过繁琐的环境准备直接聚焦实时计算链路验证也可结合课程内容深入理解原始业务数据如何进入Kafka或MySQL为后续Flink SQL开发提供数据基础。1. 5小时能到什么程度实时湖仓不是“数据搬过来”这么简单你在RDS MySQL里把一条订单状态改成“已发货”几秒之后数仓里的明细表已经跟着变了报表不用拉批处理。这就是实时湖仓要解决的问题而阿里云实时计算Flink版是目前把这条链路做得最“短平快”的托管方案。标题里这个以“原始业务数据脚本.zip”结尾的工程包对应的就是一套从业务库同步原始数据到湖存储的脚本集合常见内容包括binlog准备语句、Flink SQL建表模板、作业提交脚本和巡检脚本。5小时这个时间承诺对已经会用SQL、但没碰过流计算的人是可以实现的前两小时搭环境和理解链路中间两个小时跑通第一条同步作业最后一小时踩完坑。这篇文章适合两类人被要求“把业务数据实时同步进数仓”的后端工程师以及想评估Flink实时湖仓值不值得投入的数据团队。先说清楚组件选型再给可复制的命令和参数最后是翻车记录。2. 实时湖仓的组件选型为什么这套Flink组合在阿里云上最省心2.1 Flink在实时湖仓里的角色不是“搬运工”很多人第一次接触Flink实时湖仓时习惯把Flink理解成一个“更快的DataX”。这个类比在入门阶段能帮助理解但到了设计阶段会误导选型。DataX这类离线同步工具是“拉一次、落一次”的批处理源端没有增量概念目标端也不关心数据是不是按主键变更的。Flink在实时湖仓里做的是三件事读取变更流、维护状态、以Exactly-Once语义写出。这意味着它不只是把数据从一个地方搬到另一个地方而是要在内存里记住“这个订单当前是什么状态”才能在下一条更新到来时决定怎么改写目标存储。实时湖仓这个叫法本质是把数据湖的存储成本和数据仓库的查询效率合并。原始业务数据落到OSS或Hive表之后仍然保留最细粒度的明细后续再通过批处理或流式任务分层加工。所以Flink这层要承担的是把MySQL的binlog、Kafka的消息或者日志文件以接近实时的速度变成“可被查询的湖上表”。阿里云实时计算Flink版把这块的运维压力接了过去你在SQL编辑器里写作业平台负责拉起JobManager和TaskManager、做Checkpoint、处理故障恢复。相比自建集群省掉的不只是机器成本还有“Flink安装配置到部署”那一整条血泪路径。2.2 CDC、连接器与存储先把三条链路画清楚实时湖仓的链路可以拆成三段每一段都有对应的组件选择。第一段是接入层最常见的是Flink CDC直接监听MySQL或PostgreSQL的binlog把增删改都解析成事件流。第二段是计算层就是Flink本身它负责把CDC事件流做清洗、补字段、按主键去重然后交给下游。第三段是存储层常见选择是OSS上的Parquet文件加Hive/iceberg元数据或者直接用Paimon这类湖格式。阿里云上的典型组合是“RDS MySQL 实时计算Flink版 OSS Hive元数据服务”这条链路能覆盖从原始业务数据到ODS层的大部分场景。选型时容易忽略的是“原始业务数据脚本”这几个字里隐含的需求它要求同步任务保留业务库的原始结构不做过多加工。这决定了写入格式的选择。如果目标只是留底、供后续回溯Parquet加分区表就够了如果下游还要对这份数据做流式读取或UPSERT那就要上Iceberg或Paimon这类支持ACID的湖格式。Flink官方连接器对Hive和OSS的支持最成熟配置样例也多团队没有湖格式经验时从Hive表加Parquet起步踩坑成本最低。2.3 阿里云实时计算Flink版 vs 自建集群取舍和成本要不要直接用阿里云实时计算Flink版这是第一个需要拍板的问题。自建集群的优势是版本完全可控能装各种自定义插件但代价是你要自己处理Flink的HA、Checkpoint存储、监控告警和版本升级。实时计算Flink版虽然是托管的但Flink SQL、连接器、UDF这些开发方式跟开源版基本一致换成本地集群时作业代码可以平迁区别主要在部署和运维层面。我一般用三个问题来判断是否适合托管第一团队有没有专职的Flink运维人员第二作业数量是否长期超过三个第三是否已经使用了阿里云RDS和OSS。如果三题里有两题答案是“是”直接选托管版更划算。成本上按CU计费新手期用4~8个CU跑CDC同步作业足够比自建三台ECS加云盘的成本低而且省掉了值班成本。需要留意的是托管版的连接器版本跟随平台不能像自建那样随便改Flink小版本遇到连接器BUG时通常只能等平台侧修复这是托管换省心必须接受的约束。3. 动手之前把RDS的binlog、RAM权限和本地CLI一次配齐3.1 开通RDS MySQL与实时计算Flink版最小配置清单正式开始之前先把需要开通的资源列出来。RDS MySQL选择5.7或8.0均可存储空间按业务量预估但要确认实例规格支持binlog基础版也支持只是建议用高可用版切换时不会断binlog。实时计算Flink版在阿里云控制台搜索“实时计算”就能找到开通时选择工作空间所在地域建议与RDS、OSS在同一地域否则公网流量和延迟都会成为隐患。OSS Bucket也提前建好用来存Checkpoint和最终数据文件。权限方面实时计算Flink版需要被授权读取RDS binlog、写入OSS、访问Hive元数据。最省事的做法是在RAM里创建一个角色把AliyunRDSFullAccess、AliyunOSSFullAccess和AliyunDLFFullAccess这几个系统策略挂上然后把这个角色绑定到Flink工作空间的默认服务角色上。这一步在控制台里是可视化操作但权限漏掉一个作业跑起来就会报AccessDenied而且报错信息经常藏在Checkpoint失败里排查起来很费劲。3.2 原始业务库的binlog准备CDC能不能干活全看它Flink CDC读MySQL的前提是binlog格式必须是ROW且binlog_row_image要设置成FULL。很多MySQL实例默认是MIXED格式只能拿到SQL语句拿不到变更前后的完整数据CDC作业会直接报解析错误。先登录RDS执行下面这条SQL确认配置SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format; SHOW VARIABLES LIKE binlog_row_image;如果binlog_format不是ROW需要修改RDS的参数组并重启实例。阿里云RDS提供了一个更方便的做法在控制台的“参数设置”里直接改binlog_format和binlog_row_image提交后实例会自动重启。修改完binlog参数之后建议顺手把binlog的保存时长调到至少24小时避免Flink作业暂停期间binlog被清理导致恢复时找不到位点。然后创建CDC专用账号不要把高权限管理账号直接填在Flink作业里。给最小权限就够了CREATE USER flink_cdc% IDENTIFIED BY your_password; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO flink_cdc%; FLUSH PRIVILEGES;SELECT权限用于全量阶段读取历史数据REPLICATION SLAVE和REPLICATION CLIENT用于拉取binlog。注意CDC账号需要能访问所有要同步的表如果库比较多建议直接GRANT到库级别而不是逐表授权否则作业启动时每张表都要验证权限失败信息会非常零散。3.3 本地脚本环境用阿里云CLI和Maven镜像把路铺平很多人在这一步浪费过时间Flink SQL写好了但作业脚本在本地没法验证语法只能一遍遍传到平台上跑。我建议本地装三样东西阿里云CLI、JDK 8或11、Flink客户端。阿里云CLI用于调用OpenAPI查询作业状态和触发作业JDK是跑Flink SQL客户端的前提Flink客户端版本尽量跟实时计算Flink版主版本对齐。另外如果你要写自定义UDFMaven的settings.xml里配一下阿里云仓库镜像依赖下载速度会快很多避免反复等超时mirror idaliyunmaven/id mirrorOfcentral/mirrorOf urlhttps://maven.aliyun.com/repository/public/url /mirror阿里云CLI的配置也很直接执行一次即可aliyun configure --mode AK --profile flink \ --access-key-id LTAI5tXXXXXXXXXX \ --access-key-secret your_secret \ --region cn-hangzhouAccessKey建议用RAM子账号生成只授予实时计算和RDS的只读权限不要用主账号Key。配置完成后用aliyun flink list-workspaces这类命令验证连通性能返回工作空间列表就说明CLI通路没问题。这套本地环境搭好之后后面所有脚本都可以在本地先试跑再提交到线上。4. 把同步作业一次跑通Flink SQL加脚本是最小可行方案4.1 在Flink SQL上建CDC源表核心参数逐行拆实时计算Flink版的SQL开发面板本质上就是Flink SQL跟开源版语法一致。第一步是创建源表对应MySQL里的业务表。以订单表orders为例建表语句如下CREATE TABLE mysql_orders ( id BIGINT PRIMARY KEY NOT ENFORCED, user_id BIGINT, amount DECIMAL(10, 2), order_status STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3) ) WITH ( connector mysql-cdc, hostname rm-xxxx.mysql.rds.aliyuncs.com, port 3306, username flink_cdc, password your_password, database-name trade_db, table-name orders, server-id 5400-5404, scan.startup.mode initial, scan.incremental.snapshot.chunk.key-column id );这里重点说三个参数。server-id是Flink CDC连接器伪装成MySQL从库时使用的ID范围值5400-5404表示分配5个ID给并发分片多并行度时必须提供区间否则多个分片共用同一ID会在MySQL侧冲突。scan.startup.mode设为initial表示作业启动时先做全量快照再无缝切到增量binlog这是首次同步最常用的模式但要注意全量阶段会占用源库IO大表建议放到低峰期。scan.incremental.snapshot.chunk.key-column建议指定主键或唯一键默认会选第一个主键但如果主键是UUID字符串分片效率很差换成长整型id会快很多。4.2 建Hive目标表分区提交是数据入表的关键原始业务数据落地到Hive表时通常按日期分区。目标表DDL这样写CREATE TABLE ods_orders ( id BIGINT, user_id BIGINT, amount DECIMAL(10, 2), order_status STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3) ) PARTITIONED BY (dt STRING) WITH ( connector hive, sink.partition-commit.trigger partition-time, sink.partition-commit.delay 1 min, sink.partition-commit.policy.kind metastore,success-file, sink.partition-commit.watermark-time-zone Asia/Shanghai );分区提交是Hive Sink最容易出问题的环节。数据写进临时目录后要等分区提交策略确认“这个分区可以对外可见了”才真正变成Hive里的分区。partition-time触发方式依赖事件时间水印新增数据的时间字段要能映射到分区字段delay设1分钟是给迟到数据一点缓冲。metastore,success-file表示提交分区时同时写元数据和success标记文件下游只要看到success文件就知道分区完整。最后那个watermark-time-zone必须设成业务时区否则水印按UTC计算分区提交时间会差8小时——这个问题排查起来非常像玄学。DML语句反而最简单把源表和目标表接起来即可。单表同步不需要JOIN也不需要过滤条件INSERT INTO ods_orders SELECT id, user_id, amount, order_status, create_time, update_time, DATE_FORMAT(update_time, yyyy-MM-dd) FROM mysql_orders;这里要意识到一个设计取舍原始业务数据脚本通常只做原样落地不做数据清洗。所以这个SELECT里唯一的加工是把update_time转成分区字段字符串业务逻辑全部留给下游的数仓层处理。这符合ODS层的定位。4.3 提交流程脚本化提交、巡检、重启三件套Flink SQL作业在阿里云托管版里的正式提交通常是控制台操作但如果要把这套流程交给团队复用写成脚本更靠谱。本地开发调试时用Flink自带的SQL客户端跑是最快的验证方式。把上面的建表和DML写进一个cdc_orders.sql文件然后执行# 本地验证SQL作业生产环境在阿里云控制台以作业形式提交 $FLINK_HOME/bin/sql-client.sh \ -f /app/flink_jobs/cdc_orders.sql \ -D execution.checkpointing.interval60s \ -D execution.checkpointing.modeEXACTLY_ONCE \ -D execution.checkpointing.state-backendfilesystem \ -D execution.checkpointing.unalignedtrue这段命令的作用是让SQL客户端驱动Flink作业运行-f指定SQL文件-D参数在启动时覆盖默认配置。execution.checkpointing.interval设60秒意味着每60秒做一次Checkpoint这是实时湖仓同步作业的常见节奏太频繁会增大数据库和OSS压力太少则故障恢复时丢数据窗口变大。EXACTLY_ONCE保证端到端不重不丢。unalignedtrue开启非对齐Checkpoint在状态较大或源端压力不均时能明显减少Checkpoint超时代价是恢复时状态加载时间略长。作业上云之后日常巡检也可以脚本化。Flink的JobManager暴露标准REST API查询作业是否运行中的最简脚本如下#!/bin/bash # 作业状态巡检配合crontab每5分钟执行一次 JOB_MANAGER_URLhttp://localhost:8081 JOB_ID$1 curl -s $JOB_MANAGER_URL/jobs/overview | \ jq -r --arg jid $JOB_ID \ .jobs[] | select(.jid$jid) | \(.state) \(.name) \(.tasks.total - .tasks.running) tasks not running脚本输出作业状态和异常Task数量。在托管版里JobManager地址通常不直接暴露可以通过控制台查看或配置公网Endpoint原理一致。如果返回的state不是RUNNING就需要看日志做恢复。作业失败时最常用的处理是先停止再根据Checkpoint恢复不要直接改SQL重启否则可能从最早位点重新消费造成重复数据。5. 原始数据同步的避坑记录5个翻车现场与参数修正5.1 Flink的JDBC连接器异常驱动版本和时区一起背锅现象作业运行几小时后日志里出现Communications link failure或者Access denied for user但连接串和账号密码检查过都没问题。重启后能恢复过一阵又断。原因这个问题我在MySQL 8.0实例上遇到过两次。第一层原因是JDBC驱动版本过低MySQL 8.0默认的caching_sha2_password认证插件在旧驱动下无法完成握手第二层原因是连接串里没指定时区MySQL 8.0要求显式设置serverTimezone否则驱动用JVM默认时区常在特定时间点触发连接重置。实时计算Flink版的默认连接器版本一般能兼容自建或写自定义Sink时最容易踩。解决在JDBC连接串里显式加useSSLfalseserverTimezoneAsia/Shanghai同时把连接器版本升级到8.x。这里有个经验同步作业的连接池参数别用默认值connection.max-retry-timeout设到60秒以上可以避免源端短暂闪断导致作业直接失败。5.2 Flink sink Hive表数据不入表分区提交策略的坑现象作业状态是RUNNING源表数据一直在更新但查Hive分区时发现要么没有新分区要么分区下的文件是0字节。目标表看起来一切正常就是“数据不入表”。原因Hive Sink写出的数据停留在临时目录分区提交策略没有真正触发。常见原因有三个分区时间字段取的列不对导致水印算不出来sink.partition-commit.delay设置太长或太短sink.partition-commit.policy.kind漏掉了metastore只写了success-file导致数据文件到位了但元数据没注册。解决先把delay调到1分钟以内做验证确认分区能出来再把分区提交的watermark-time-zone设成Asia/Shanghai。查数据时别用SHOW PARTITIONS直接用SELECT COUNT(*) FROM ods_orders WHERE dt2024-06-01验证更直接。5.3 Checkpoint一直失败状态后端和OSS限流现象作业能跑但每隔一段时间就checkpoint declined或checkpoint expired任务频繁重启数据延迟越来越高。原因实时计算Flink版默认把Checkpoint存到OSS如果业务表数据量大、状态大同时Checkpoint间隔设置过短比如10秒大量小文件写入OSS会触发限流。另一个隐蔽原因是源表没有主键CDC连接器必须依赖全局状态做去重状态无限膨胀。解决把execution.checkpointing.interval调到60秒以上确认源表有主键并在DDL里声明PRIMARY KEY NOT ENFORCED在OSS侧查看Checkpoint目录的文件大小如果单个Checkpoint超过500MB优先优化SQL而不是扩容。一个容易被忽略的做法是开启execution.checkpointing.unalignedtrue它能跳过对齐等待直接保存当前状态对CDC同步这种延迟敏感的作业效果明显。5.4 全量转增量时丢数据位点衔接的黑匣子现象大表首次同步全量阶段数据条数对得上增量阶段开始后下游出现部分数据缺失检查源库binlog发现有些变更没被消费。原因全量快照是按主键分片并行读取的分片之间没有强一致快照某分片读完后立刻切到binlog位点但另一个分片还在读旧数据中间产生的变更恰好落在空档。server-id配置不当会放大这个问题多个并行分片使用同一ID时MySQL会踢掉旧连接。解决给每个并行度分配独立的server-id区间首次同步不要开太高的并行度1~2个并行度最稳如果数据量过大先用scan.startup.modeinitial跑通小表再对核心大表单独建作业。全量结束后的增量断点很难肉眼发现建议在同步作业里加一行只有增量阶段才更新的update_time字段用下游最大值与源库对比判断是否追平。5.5 业务表加列后作业崩溃schema变更的处理策略现象业务方给表加了一个字段Flink作业直接报Table schema mismatch整个链路中断。原因Flink SQL作业在启动时把源表结构注册成静态Schema源端结构变化后无法自动感知。如果目标Hive表也同步改了连接器解析binlog时发现字段数与Schema不符直接抛异常。解决最简单的是停作业、更新DDL、从最近一次Checkpoint恢复。更省事的是减少人工介入把这类变更统一收口到“只保留原始数据”的同步层源表DDL里加一列_raw_data STRING保存整行JSON后面所有字段解析都在下游做。若需要自动化可以关注Flink CDC Pipeline部署方案它以YAML描述整条链路对schema变更的容错更好团队有精力再评估。6. 从“能跑”到“好用”验证、火焰图与血缘的收尾6.1 用火焰图定位反压别急着加并行度作业延迟升高时第一反应是加资源但很多时候瓶颈不在CPU而在某个算子背压。Flink火焰图能看到每个算子的CPU消耗热点先看JobManager里TaskManager的火焰图如果某个Source算子CPU跑满是读取端压力大如果Sink算子CPU跑满是写入端慢。提高并行度前先确认是不是目标端Hive分区提交卡住了小文件合并把sink.partition-commit.delay调大有时比加并行度更有效。6.2 验证原始数据没丢没重两条SQL就够了验证同步质量用最朴素的办法源库数和目标表数分别统计。Flink SQL对源表的实时count代价高可以在业务低峰期跑一次离线比对日常用目标表的update_time最大值间接判断如果这个值距当前时间超过5分钟说明链路有延迟。另外在Hive表上定期跑SELECT dt, COUNT(*) FROM ods_orders GROUP BY dt对照源库每天的总行数偏差超过阈值就触发告警这套验证在数据量不大时足够可靠。6.3 血缘与元数据别到最后补课实时湖仓的数据血缘跟离线数仓一样重要否则半年后没人说得清这张ODS表是从哪个业务库来的。用OpenMetadata这类元数据工具能直接抓Flink作业的血缘关系同步任务交付时把血缘信息一并入库后续排查“某张报表的数据为什么不对劲”会省下大量时间。这里有一条我踩过的教训自定义Data Source或Data Sink要克制先确认官方连接器覆盖不了再动手写否则每个自研连接器都会成为血缘断点和升级包袱。反过来看5小时把链路跑通只是一个开始。我现在的习惯是把所有脚本沉淀在团队仓库里binlog参数检查、建表模板、提交脚本、巡检脚本各司其职新人照着走一遍最多半天就能上手。实时湖仓没有想象中那么玄也没那么省事把基础链路和避坑参数吃透剩下的都交给时间。希望帮到你。本文还有配套的精品资源点击获取