简介: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版主版本对齐。另外,如果你要写自定义UDF,Maven的settings.xml里配一下阿里云仓库镜像,依赖下载速度会快很多,避免反复等超时:
<mirror> <id>aliyunmaven</id> <mirrorOf>central</mirrorOf> <url>https://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.interval=60s \ -D execution.checkpointing.mode=EXACTLY_ONCE \ -D execution.checkpointing.state-backend=filesystem \ -D execution.checkpointing.unaligned=true这段命令的作用是让SQL客户端驱动Flink作业运行,-f指定SQL文件,-D参数在启动时覆盖默认配置。execution.checkpointing.interval设60秒,意味着每60秒做一次Checkpoint,这是实时湖仓同步作业的常见节奏,太频繁会增大数据库和OSS压力,太少则故障恢复时丢数据窗口变大。EXACTLY_ONCE保证端到端不重不丢。unaligned=true开启非对齐Checkpoint,在状态较大或源端压力不均时能明显减少Checkpoint超时,代价是恢复时状态加载时间略长。
作业上云之后,日常巡检也可以脚本化。Flink的JobManager暴露标准REST API,查询作业是否运行中的最简脚本如下:
#!/bin/bash # 作业状态巡检:配合crontab每5分钟执行一次 JOB_MANAGER_URL="http://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连接串里显式加useSSL=false&serverTimezone=Asia/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 dt='2024-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.unaligned=true,它能跳过对齐等待,直接保存当前状态,对CDC同步这种延迟敏感的作业效果明显。
5.4 全量转增量时丢数据:位点衔接的黑匣子
现象:大表首次同步,全量阶段数据条数对得上,增量阶段开始后,下游出现部分数据缺失,检查源库binlog发现有些变更没被消费。
原因:全量快照是按主键分片并行读取的,分片之间没有强一致快照,某分片读完后立刻切到binlog位点,但另一个分片还在读旧数据,中间产生的变更恰好落在空档。server-id配置不当会放大这个问题,多个并行分片使用同一ID时,MySQL会踢掉旧连接。
解决:给每个并行度分配独立的server-id区间;首次同步不要开太高的并行度,1~2个并行度最稳;如果数据量过大,先用scan.startup.mode=initial跑通小表,再对核心大表单独建作业。全量结束后的增量断点很难肉眼发现,建议在同步作业里加一行只有增量阶段才更新的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参数检查、建表模板、提交脚本、巡检脚本各司其职,新人照着走一遍最多半天就能上手。实时湖仓没有想象中那么玄,也没那么省事,把基础链路和避坑参数吃透,剩下的都交给时间。希望帮到你。
本文还有配套的精品资源,点击获取