
简介这是一份面向大数据初学者的 Flink 实战入门资源围绕音乐专辑数据完成从采集、清洗、聚合分析到可视化展示的整套流程。资源包含参考代码其中 flinkProject 对应数据处理部分DrawPic 负责可视化展示可帮助理解 DataStream API、时间窗口、状态管理等核心概念并掌握将分析结果以图表形式呈现的方法。压缩包共86个文件以编译后的 class、xml 配置、csv 数据文件和 html 结果页为主整体仅 2.21MB适合快速下载与本地运行调试。已有563人学习浏览对于想通过实际项目快速上手 Flink 的初学者或数据开发人员来说是一份轻量且完整的练习素材。尽管标签显得随性但代码结构清晰覆盖从数据接入到图表输出的关键环节并提供了常见调试与输出配置思路。1. 基于 Flink 的专辑数据分析为什么难度低但值得做一张 CSV 里躺着一份专辑数据发行日期有斜杠有点号收听量带着千分位逗号流派大小写还不统一。你要回答的问题却很朴素哪些专辑最能打、哪个年代最耐听、哪种流派最能卖。用 Flink 做音乐专辑数据分析展示真正的挑战不在计算而在把脏数据洗成一张可靠的宽表。好在不用像 Hive 项目那样先搭 Hadoop 全家桶本地一个 Flink 集群加 SQL Client 就能跑大部分逻辑用 SQL 写完。适合刚学完 Flink 基础找数据分析项目实践的人也适合想把分析结果接进 Spring Boot 或 BI 工具的人。下文按建模、清洗、指标、落库、展示的顺序把链路和重跑时的坑一次说清。2. 专辑数据建模与清洗CSV 源表到 Flink 宽表怎么搭2.1 先定指标再定字段专辑事实表该有的字段分析展示类项目最忌讳拿到数据就写 SQL。我一般会先问自己要出哪几张图流派 TOP 10 榜单、年代收听趋势、流派占比、销量与评分的关系。图反推字段一张专辑事实表基本十个字段就够了专辑 ID、专辑名、艺人、流派、发行日期、曲目数、总时长、总收听量、总销量、平均评分、评分人数。评分人数要和平均评分一起出现否则「一个人打 5 分」的专辑会把排行榜带偏。这个「先定图再定表」的习惯能让你后面写 Flink SQL 时少走两次弯路。我准备的样例数据长这样注意里面埋了几处脏数据album_id,album_name,artist,genre,release_date,track_count,duration_sec,total_plays,total_sales,avg_rating,rating_cnt 10001,Kind of Blue,Miles Davis,jazz,1959/08/17,9,2730,8234100,520000,4.8,1200 10002,Thriller,Michael Jackson,Pop,1982-11-30,9,2530,10,024,500,66000000,4.7,28000 10003,Nevermind,Nirvana,Grunge,1991-09-24,,1730,6801200,,4.5,900第一行 genre 是首字母小写的jazz日期用斜杠分隔第二行total_plays带千分位逗号还加了引号第三行track_count和total_sales直接是空值。这些不是巧合是我故意放的——后面每一条清洗规则都对应一个真实存在的脏数据场景。数据规模在几十万行级别时本地内存完全跑得动。公开的专辑数据可以参照 Last.fm 或 Million Song Dataset 的字段思路自己造一份不必纠结真实度分析链路跑通比数据本身真实更重要。数据文件的命名也建议统一比如albums_20250601.csv这样批作业读哪个文件一目了然后面做调度和重跑时能少很多口舌。2.2 本地集群与 SQL Client十分钟跑起 Flink 批环境Flink 不需要 Hadoop 环境下载二进制发行版解压就能用。打开终端执行tar -xzf flink-*.tgz cd flink-* ./bin/start-cluster.sh # 启动本地集群 ./bin/sql-client.sh # 进入 SQL 客户端start-cluster.sh会在本机拉起一个 JobManager 和一个 TaskManager默认内存各 1 GB 左右练习场景够用。启动后浏览器访问localhost:8081能看到 Web UI作业状态、日志、异常都在里面后面排查问题先看这里。SQL Client 是交互式会话也可以直接用-f参数跑脚本后面调度章节会用到。进入客户端后第一件事是切批模式SET execution.runtime-mode BATCH;CSV 是一次性快照不是持续产生的流用批模式才有「跑一遍出一个确定结果」的语义重跑也容易对账。如果不切默认流模式会把文件当无界流读聚合结果在流模式下默认是增量输出排序和窗口的语义都会和批模式不一样。数据分析和展示场景下批模式是首选实时化是后面的事。2.3 用 Flink SQL 清洗脏数据日期、千分位、空值与大小写先建源表。Flink 的 filesystem connector 负责读目录csv format 负责解析每一行CREATE TABLE album_raw ( album_id BIGINT, album_name STRING, artist STRING, genre STRING, release_date STRING, track_count BIGINT, duration_sec BIGINT, total_plays STRING, total_sales STRING, avg_rating DOUBLE, rating_cnt BIGINT ) WITH ( connector filesystem, path /home/user/music-album/input, format csv, csv.field-delimiter ,, csv.ignore-parse-errors true );这里故意把total_plays、total_sales定义为 STRING因为源数据里带千分位逗号直接按 BIGINT 解析会整行报错。csv.ignore-parse-errors是一把双刃剑它能跳过坏行保作业不挂但也会悄悄丢数据所以清洗完之后一定要做行数核对第 5 章会讲验证方法。接着建清洗视图这是整条链路的核心CREATE TEMPORARY VIEW album_clean AS SELECT album_id, TRIM(album_name) AS album_name, CONCAT(UPPER(SUBSTRING(TRIM(genre), 1, 1)), LOWER(SUBSTRING(TRIM(genre), 2))) AS genre, TO_TIMESTAMP(REGEXP_REPLACE(release_date, [./], -), yyyy-M-d) AS release_ts, COALESCE(track_count, 0) AS track_count, COALESCE(duration_sec, 0) AS duration_sec, CAST(REGEXP_REPLACE(TRIM(total_plays), ,, ) AS BIGINT) AS total_plays, COALESCE(CAST(REGEXP_REPLACE(TRIM(total_sales), ,, ) AS BIGINT), 0) AS total_sales, COALESCE(avg_rating, 0.0) AS avg_rating, rating_cnt FROM album_raw;这几条规则基本覆盖了 CSV 专辑数据的常规脏点。REGEXP_REPLACE(release_date, [./], -)把斜杠和点号统一成短横线再交给TO_TIMESTAMP(..., yyyy-M-d)解析一个格式参数吃下1959/08/17、1959.8.17、1959-08-17三种写法。genre 的首字母大写用CONCAT SUBSTRING实现比直接UPPER全大写更适合展示层读。千分位逗号先去后转空销量用COALESCE兜底成 0 而不是 NULL因为后面算占比时NULL 参与 SUM 会让整个结果变 NULL。注意CREATE TEMPORARY VIEW只活在当前 SQL 会话里脚本每次重跑都会重建这正符合批作业「可重复、可回滚」的预期。如果原始 CSV 大到几十 GB清洗结果可以考虑用INSERT INTO ... SELECT落到 parquet 目录避免每次跑批重新解析大文件百万行以内视图方案完全够用。3. 分析指标落到 SQL排行、趋势与相关性一次算清3.1 各流派 TOP 10 专辑ROW_NUMBER 排名写法「每个流派收听量最高的 10 张专辑」是榜单类页面的标准需求。Flink SQL 用窗口函数一次写完SELECT genre, album_name, total_plays FROM ( SELECT genre, album_name, total_plays, ROW_NUMBER() OVER (PARTITION BY genre ORDER BY total_plays DESC) AS rn FROM album_clean ) t WHERE rn 10 ORDER BY genre, rn;外层WHERE rn 10是实现分组 Top N 的固定姿势。rn 这个别名在同一个 SELECT 层里不能直接用必须包一层子查询这是 SQL 执行顺序决定的WHERE在窗口函数计算之前就已经处理完了。PARTITION BY genre按流派开窗和GROUP BY的差异在于它不折叠行每行保留自己的排名适合「每个分组取前几条」的场景。如果两张专辑播放量完全一致排名就有随机性前端刷新一次变一次。建议在ORDER BY里补第二排序键比如total_plays DESC, avg_rating DESC让并列情况也有确定顺序。榜单类页面最忌讳的就是结果不稳定。3.2 年代趋势与流派占比ROLLUP 一次出三档粒度前端要做「年代折线图 流派饼图 总计卡片」最省事的做法是用一个 SQL 把年、流派、年加流派三个粒度同时算出来SELECT YEAR(release_ts) AS release_year, genre, SUM(total_plays) AS plays, COUNT(DISTINCT album_id) AS album_cnt FROM album_clean GROUP BY ROLLUP (YEAR(release_ts), genre);ROLLUP 是 GROUPING SETS 的特例按括号里从左到右逐层收拢一次返回四类行年加流派的明细聚合、只有年、只有流派、总合计。前端拿到release_year为 NULL 的组就知道那是合计行饼图、折线图、总计卡片都能从同一份结果里取数不需要后端为每个图表写一条独立 SQL。如果两组数据量级相差悬殊比如 Pop 的播放量是 Classical 的几十倍饼图会被大块占满。想算占比可以用SUM(plays) / SUM(SUM(plays)) OVER ()但更常见的做法是 SQL 只输出分母和分子占比留给前端算。后端只给事实不给加工观点前端要什么粒度自己除这样改口径时不用前后端一起改。3.3 「叫好又叫座」判断销量与评分的相关性分析榜单看热度散点图看口碑。Flink SQL 内置聚合函数CORR可以直接算销量和评分的相关系数SELECT YEAR(release_ts) AS release_year, CORR(total_sales, avg_rating) AS rating_sales_corr, SUM(total_plays) / NULLIF(COUNT(DISTINCT album_id), 0) AS avg_plays_per_album FROM album_clean GROUP BY YEAR(release_ts) ORDER BY release_year;CORR(total_sales, avg_rating)输出 -1 到 1 的相关系数。音乐数据里它通常是弱正相关卖得好的专辑评分不一定最高但销量垫底的专辑评分大概率也不高。NULLIF(COUNT(...), 0)防止某一年一张专辑都没有时除零报错。如果想用散点图看单张专辑的位置把album_id、album_name也查出来GROUP BY改成album_id每个点就是一张专辑。相关系数这个指标单独看意义不大但放在「年度趋势」里能讲故事某一年相关系数突然变负说明那一年市场口味和评论口碑严重分裂这是数据分析展示里很有话题性的素材。3.4 把指标统一成视图给展示层留好查询口三张图的指标逻辑都有了但它们是三段独立 SQL。为了让展示层和落库层共用同一套口径我习惯把它们注册成统一结构的视图CREATE TEMPORARY VIEW top_by_genre AS SELECT top_by_genre AS indicator, genre AS dimension, total_plays AS value FROM ( ... 3.1 的 TOP 10 SQL ... ); CREATE TEMPORARY VIEW genre_trend AS SELECT genre_trend AS indicator, CONCAT(CAST(YEAR(release_ts) AS STRING), -, genre) AS dimension, SUM(total_plays) AS value FROM album_clean GROUP BY YEAR(release_ts), genre; CREATE TEMPORARY VIEW rating_corr AS SELECT rating_corr AS indicator, CAST(YEAR(release_ts) AS STRING) AS dimension, CORR(total_sales, avg_rating) AS value FROM album_clean GROUP BY YEAR(release_ts);统一视图的价值在口径也在结构。以后想改「TOP 10 变 TOP 20」只动top_by_genre这一处前端、落库、日报全部跟着变。三个视图输出都是indicator dimension value后面所有下游只消费这三个字段结构越简单传递过程中越不容易打架。第 4 章的落库就直接消费这三个视图。4. 结果落库与跑批调度Flink JDBC 连接器避坑与重跑策略4.1 JDBC Sink 建表与参数从指标视图到 MySQLFlink 官方 JDBC 连接器负责把结果写进 MySQL 这类关系库。目标表结构刻意简化落库 DDL 长这样CREATE TABLE result_sink ( indicator STRING, dimension STRING, value DOUBLE, stat_date STRING, PRIMARY KEY (indicator, dimension, stat_date) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/music_insight?useSSLfalsecharacterEncodingutf8serverTimezoneAsia/Shanghai, table-name result, username flink_writer, password change_me, sink.buffer-flush.max-rows 500, sink.buffer-flush.interval 5s, sink.max-retries 3 );PRIMARY KEY ... NOT ENFORCED的意思是 Flink 侧声明主键但不校验。带主键的 JDBC sink 在 MySQL 上会走 upsert 语义同一把 key 重复写入是更新而不是追加这是重跑不出脏数据的核心。JDBC Sink 的 4 个必调参数参数默认值建议sink.buffer-flush.max-rows100攒够行数再批量写500~1000 对 MySQL 更友好sink.buffer-flush.interval1s兜底刷出时间数据量小时设 5s 减少连接开销sink.max-retries3单次连接失败重试次数别滥用大数值sink.parallelism无批模式建议等于结果分区数避免写出乱序连接串里有三个参数基本是标配useSSLfalse省去本地环境证书开销characterEncodingutf8防中文乱码serverTimezoneAsia/Shanghai避免时间偏移八小时。MySQL 侧如果开批量写入连接串再加rewriteBatchedStatementstrue一次网络往返能塞进更多 INSERT吞吐差距肉眼可见。写入动作本身很直白三条 INSERT 对应三个视图INSERT INTO result_sink SELECT indicator, dimension, value, DATE_FORMAT(CURRENT_TIMESTAMP, yyyy-MM-dd) FROM top_by_genre; INSERT INTO result_sink SELECT indicator, dimension, value, DATE_FORMAT(CURRENT_TIMESTAMP, yyyy-MM-dd) FROM genre_trend; INSERT INTO result_sink SELECT indicator, dimension, value, DATE_FORMAT(CURRENT_TIMESTAMP, yyyy-MM-dd) FROM rating_corr;三条 INSERT 在 SQL Client 脚本里会依次触发三个 Flink 作业互不阻塞。哪个作业失败就在日志里标记哪条排查范围比一个大作业裹在一起清晰得多。4.2 用 cron 跑批SQL 脚本、日志与重跑兜底难度低的项目最常见做法就是单机调度。把前面所有 SQL 存成一个analysis.sql开头放批模式设置最后放三条 INSERT然后写 cron0 2 * * * cd /opt/music-analysis /opt/flink/bin/sql-client.sh -f analysis.sql logs/analysis.log 21sql-client.sh -f是非交互执行脚本里第一条语句出错就退出后面的不会继续。日志里抓ERROR就能做告警配合一个简单的grep ERROR logs/analysis.log巡检脚本就够用。cron 定在凌晨 2 点有个前提CSV 文件必须在 2 点前落到 input 目录且不能边写边读。我见过有人直接把文件传进 input 目录结果作业读到一半文件还在拷贝行数对不上。常见做法是先把文件传到临时目录再由一条mv命令移入 input 目录作业一跑就是一个干净的快照。如果指标要更实时把 CSV 换成 Kafka source执行模式切流JDBC Sink 几乎不用动想顺手把 MySQL 里的明细同步到 ClickHouse 做即席查询原理也是在同一套 Flink 里多配一个 ClickHouse Sink。批起步、流演进是这类项目最稳的路线难度低不等于不能长大。4.3 四条真实踩坑记录现象、原因与解决坑 1作业一提交就报 ClassNotFoundException: com.mysql.cj.jdbc.Driver现象JDBC Sink 初始化阶段直接失败TaskManager 日志里找不到 MySQL 驱动类。原因Flink 发行版默认不打包第三方驱动。flink-connector-jdbc虽然带了连接器但mysql-connector-j这个驱动 jar 要单独放进$FLINK_HOME/lib或打进作业 jar。解决用 SQL Client 就把mysql-connector-j-*.jar拷到lib目录后重启集群用 Java 工程就用maven-shade-plugin打 fat jar。这个坑几乎排在所有 Flink JDBC 连接器异常案例的第一位遇到先查依赖。坑 2运行中偶发 Communications link failure现象作业跑到一半连不上 MySQL重试几次后失败手动重跑又好了而且不是每次必现非常玄学。原因MySQL 的wait_timeout把空闲连接回收了而 Flink 连接池里的连接还在被复用或者连接串没配autoReconnect。解决连接串加autoReconnecttrueMySQL 侧调大wait_timeout同时把sink.max-retries设成 3 兜底。批作业大不了重跑但重跑必须保证幂等见下一条。坑 3同一指标跑两次MySQL 里行数翻倍现象昨天还正常的表今天一查COUNT(*)变成两倍值一模一样第一天跑批就翻车。原因作业失败后重跑目标表没有唯一键约束JDBC 连接器按追加模式把数据又插了一遍。解决目标表建(indicator, dimension, stat_date)唯一索引Flink 侧声明PRIMARY KEY ... NOT ENFORCED走 upsert或者跑批第一步先DELETE FROM result WHERE stat_date 当天再插入。我一般两种都做SQL 里幂等重跑前清当天分区双保险。坑 4专辑名写进 MySQL 全变问号时间整整差 8 小时现象中文全部变成?日期字段对不上。原因连接串没带characterEncodingutf8serverTimezone没指定驱动用了服务器默认时区和 Flink 本地时区错位时间就偏了 8 小时。解决连接串统一补characterEncodingutf8serverTimezoneAsia/ShanghaiJVM 加-Duser.timezoneAsia/Shanghai。CSV 文件本身也统一成 UTF-8避免源文件是 GBK 时读出来就是乱码清洗层再怎么写也救不回来。5. 展示与验证一套可复算的口径让图表经得起追问展示层最省事的架构是Flink 算完Spring Boot 只做查询接口。Spring Boot 整合 Flink 的项目里这是最常见分工——Flink 干重活应用服务只读 MySQL 里已经算好的 result 表前端拿到的 JSON 直接渲染图表GetMapping(/api/indicators/{type}) public ListIndicator list(PathVariable String type, RequestParam(defaultValue latest) String batch) { return mapper.selectByIndicator(type, resolveBatchId(batch)); }接口只按 indicator 类型查聚合已经在 Flink 里算完应用层不写一行统计代码。指标和前端的对应关系固定成一张表指标dimension 内容前端图表top_by_genre流派 专辑名横向柱状图genre_trend年份-流派折线图或堆叠面积图rating_corr年份散点图ECharts 侧只需要最朴素的 optionfetch(/api/indicators/top_by_genre) .then(r r.json()) .then(data myChart.setOption({ yAxis: { type: category, data: data.map(d d.dimension) }, xAxis: { type: value }, series: [{ type: bar, data: data.map(d d.value) }] }));写图表之前先把验证做完。我每次跑完批做三件事第一抽一个流派手算 SUM和 Flink 输出对不上就查清洗层第二SELECT COUNT(*)对比源 CSV 行数和清洗后行数差值必须能解释坏行丢了多少、空值过滤了多少都得有数第三用 Python 或 Spark 把同样逻辑跑一遍做双跑对照两边数字一致才敢让这张图见人。这套验证做下来别人质疑数据时你打印一条 SQL 就能讲清口径。最后补一句个人习惯指标 SQL 永远存成带日期版本的文件改口径就新增文件而不是改线上出问题 5 分钟就能回滚到上一个版本。希望帮到你。本文还有配套的精品资源点击获取