简介面向大数据入门学习者的 Apache Flink 实战资源以音乐专辑数据分析与展示为完整场景覆盖从数据读取、清洗聚合、窗口统计到可视化输出的基础流程难度较低适合初次接触 Flink 的读者快速上手。压缩包共 86 个文件大小约 2.21MB包括 50 个 class 编译文件、14 个 xml 工程配置、11 个 csv 测试数据、5 个 html 结果页面以及 scala、py 等源码与辅助脚本分别对应项目运行代码、工程配置、输入数据和展示结果。已有 561 人学习使用。附带参考代码、flinkProject 数据处理工程和 DrawPic 可视化代码可帮助理解 DataStream/DataSet API、窗口操作、状态管理等关键点并通过图表呈现分析结论标签虽显随性但项目门槛不高正好用来练习 Flink 编程与数据可视化。 很多人一直搞不清楚数据工程DE和数据科学DS的分工平时也总有人私信问我“想入门 Flink但不知道该拿什么项目练手有没有不太烧脑、但又能把实时分析链路跑通的例子”说实话Flink 入门的最大障碍不是 API 本身而是很多人一上来就对着电商双十一大屏那种高并发场景发怵。这次我拿一个“音乐专辑数据分析展示”项目做例子它是 FLink 里难得的低难度、全链路、能出成果的选题特别适合刚接触 Flink、或者正在准备数据分析类面试的人。这个项目本质上就是把音乐专辑的播放、收藏、销售行为数据做成实时流交给 Flink 按专辑、歌手、流派等维度做聚合统计再把结果落到 MySQL最后用可视化工具展示排行榜和趋势图。它不依赖复杂的机器学习算法也不要求你懂高深的流计算理论你用 Flink SQL 就能把整条链路打通。整个项目做完你既能理解“数据如何流动”“流批处理到底在解决什么问题”又能收获一个完整的实时数据分析应用。这篇文章会把我的设计和踩坑过程完整写出来从环境搭建到 SQL 编写再到常见异常排查按我实际操作的顺序走一遍。1. 项目整体设计与数据链路1.1 先搞清楚这个项目到底要算什么拿到“音乐专辑数据分析展示”这个题目第一件事不是写代码而是把分析目标定清楚。我做这个项目时把核心指标锁定在三类专辑热度排行统计每个专辑在固定时间窗口内的播放量、收藏量输出 TOP N。歌手作品趋势按歌手维度聚合看一段时间内哪些歌手的总播放量增长最快。流派分布按专辑所属的音乐流派流行、摇滚、民谣、电子等统计占比。数据字段的设计也要围绕这几个指标展开。我使用的原始数据是一条条的“专辑行为事件”模拟用户对专辑的操作字段如下字段名类型说明album_idSTRING专辑唯一编号album_nameSTRING专辑名称singerSTRING歌手genreSTRING音乐流派publish_yearINT发行年份play_countINT该事件对应的播放次数favorite_countINT该事件对应的收藏次数event_timeTIMESTAMP(3)事件发生时间选这些字段是有讲究的。时间字段event_time是流处理的核心Flink 的水位线和窗口全靠它驱动play_count、favorite_count是可累加的度量值其余album_id、singer、genre都是维度字段。一个明确的宽表结构让你后面写 SQL 聚合时不用反复去 join 维表大大降低了入门难度。1.2 为什么选 Flink 而不是 Spark Streaming这个项目叫“数据分析展示”很多人的第一反应是 Spark。我也被问过不少次“Spark 也能做实时为什么非得 Flink”从工程实现角度说Spark Streaming 本质上是把流切成一个个微批次它的实时性属于“准实时”延迟通常在秒级而且调优起来比较依赖批处理经验。Flink 是真正的流式计算引擎事件一来就能处理延迟低到了毫秒级。对于音乐专辑这种用户行为数据Flink 的窗口机制滚动窗口、滑动窗口、会话窗口写起来更自然而且在事件时间语义下能更好地处理乱序数据。但选择 Flink 还有一个更现实的原因你去看现在的招聘要求和开源社区的热度Flink 在实时数仓、数据集成方面的生态明显更活跃Flink CDC 和 Flink SQL 已经成了很多公司数据平台的基础设施。既然这个项目是为了练手和积累经验直接上 Flink 的性价比更高。1.3 整体架构长什么样整个数据链路是标准的“模拟采集 - 消息缓冲 - 实时计算 - 结果存储 - 可视化展示”五段式非常适合复现用 Python 脚本模拟产生音乐专辑行为事件发送到 Kafka 的album_behavior主题。Flink 通过 Kafka Source 消费事件数据。Flink SQL 执行流式计算按时间窗口聚合。聚合结果通过 JDBC Sink 写入 MySQL。可视化端连 MySQL 展示排行榜和趋势图。这个架构看起来组件不少其实每段都“平易近人”。Kafka 可以单机模式跑Flink 用本地集群或 Standalone 集群都行MySQL 是你肯定熟悉的可视化层我选的是 Superset也可以用 ECharts 写个简单页面。整个项目跑通以后你再看公司里那些复杂实时数仓架构会发现自己至少已经摸清了主干的每一段。2. 环境准备与数据模拟2.1 Flink 安装别踩版本坑Flink 的本地安装其实不复杂但版本坑非常隐蔽。我第一次装的时候直接用系统自带 JDK 跑结果 Flink 1.17 以后对 JDK 版本有要求启动时各种报UnsupportedClassVersionError。我的建议是使用 JDK 8 或 JDK 11我实测 JDK 8 最省心去 Flink 官网下载flink-1.18.0-bin-scala_2.12.tgz解压tar -zxvf flink-1.18.0-bin-scala_2.12.tgz cd flink-1.18.0 ./bin/start-cluster.sh启动后访问http://localhost:8081看到 Flink Web UI 就说明环境通了。这个 Web UI 是你排查作业是否正常运行的重要工具后面我在问题排查部分还会用到它。如果你在 Linux 服务器上安装记得给 Flink 目录和日志目录配好权限还要确认JAVA_HOME环境变量指向了正确的 JDK 路径。我见过不少人在start-cluster.sh时出现的诡异报错最后发现只是JAVA_HOME没写对。2.2 用 Python 快速模拟专辑行为数据既然是水平入门我们不需要真的去对接音乐 App 的线上数据写个模拟脚本就够了。我用 Python 生成 JSON 格式的事件每次随机选取专辑、歌手和流派随机生成播放量和收藏量然后推给 Kafka。这里需要注意如果你不想一开始就引入 Kafka也可以先用 Flink 的FileSource或DataGen连接器生成流式数据。不过考虑到 Kafka 是生产环境最常用的数据管道组件我还是建议直接上 Kafka一遍把链路跑通。Kafka 单机安装很简单下载解压后直接用自带的 ZooKeeper 启动即可。模拟脚本的关键代码逻辑如下import json import time import random from kafka import KafkaProducer albums [ {album_id: A001, album_name: 夜航西飞, singer: 陈粒, genre: 民谣, publish_year: 2024}, {album_id: A002, album_name: 丑奴儿, singer: 草东没有派对, genre: 摇滚, publish_year: 2016}, {album_id: A003, album_name: Leaves, singer: AGA, genre: 流行, publish_year: 2023}, ] producer KafkaProducer( bootstrap_serverslocalhost:9092, value_serializerlambda v: json.dumps(v).encode(utf-8) ) while True: album random.choice(albums) event { **album, play_count: random.randint(1, 100), favorite_count: random.randint(0, 10), event_time: time.strftime(%Y-%m-%d %H:%M:%S, time.localtime()) } producer.send(album_behavior, valueevent) time.sleep(0.5)这段脚本每秒生成两条事件已经能满足演示需求。真实业务里专辑行为数据的峰值可能是每秒几万条但处理逻辑是一样的只是把时间窗口和并行度调大。我建议你在做这个项目的时候把模拟数据理解成“用户行为日志的简化版”它和日志采集系统里真正跑的数据并无本质区别。3. 核心分析逻辑与 Flink SQL 实现3.1 用 Flink SQL 建表先把数据源接进来环境就绪、数据持续产生之后进入核心环节。我用的工具是 Flink SQL Client因为低难度项目最重要的是快速看到结果没必要一开始就写 Java 代码。Flink SQL 的体验和普通关系型数据库很接近但它明确区分Source表和Sink表。首先建立 Kafka Source 表CREATE TABLE album_source ( album_id STRING, album_name STRING, singer STRING, genre STRING, publish_year INT, play_count INT, favorite_count INT, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic album_behavior, properties.bootstrap.servers localhost:9092, properties.group.id album_analysis_group, format json, json.fail-on-missing-field false, scan.startup.mode latest-offset );这里最值得讲清楚的是水位线Watermark的设定。WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND表示允许事件时间最多乱序 5 秒也就是说5 秒之前的迟到数据会被丢弃5 秒之内的会按正确的事件时间归入所属窗口。刚接触 Flink 的人通常会纠结“为什么不能完全等数据到齐”等过你就明白了实时流是无限数据你永远等不到数据到齐只能靠水位线在准确性和延迟之间做取舍。3.2 滚动窗口聚合算出实时专辑热度在 Source 表之上可以直接写查询统计每分钟各专辑的总播放量和收藏量。我使用 10 秒的滚动窗口做演示方便在页面上看到数据变化CREATE TABLE album_count AS SELECT album_id, album_name, singer, genre, TUMBLE_START(event_time, INTERVAL 10 SECOND) AS window_start, TUMBLE_END(event_time, INTERVAL 10 SECOND) AS window_end, SUM(play_count) AS total_play, SUM(favorite_count) AS total_favorite FROM album_source GROUP BY TUMBLE(event_time, INTERVAL 10 SECOND), album_id, album_name, singer, genre;这段 SQL 是 Flink 流计算中最经典的滚动窗口聚合。注意GROUP BY后面除了窗口函数还必须把所有非聚合字段album_id、album_name、singer、genre都写进去这是 SQL 的语义规定不能偷懒。窗口字段window_start和window_end能帮你清楚地看到这批数据统计的是哪个时间段。为了观察歌手维度的整体趋势再写一条按歌手聚合的查询SELECT singer, TUMBLE_START(event_time, INTERVAL 10 SECOND) AS window_start, SUM(play_count) AS total_play FROM album_source GROUP BY TUMBLE(event_time, INTERVAL 10 SECOND), singer;做完这两个查询你对 Flink SQL 的窗口机制基本就上路了。实际上实时大屏的核心逻辑也就是这些再复杂的逻辑也无非是增加维度、增加窗口类型、增加明细关联。3.3 结果写出 MySQLJDBC 连接器参数别配错计算出来的结果不能一直存在 Flink 内存里需要写到一个稳定存储。我选择 MySQL原因很直接后面用 Superset 或者任何 BI 工具连 MySQL 都非常方便不需要额外搭建 Elasticsearch 或 Redis。在 Flink 中定义 MySQL Sink 表CREATE TABLE album_sink ( album_id STRING, album_name STRING, singer STRING, genre STRING, window_start TIMESTAMP(3), window_end TIMESTAMP(3), total_play BIGINT, total_favorite BIGINT, PRIMARY KEY (album_id, window_start) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://localhost:3306/album_analysis?useSSLfalseserverTimezoneAsia/Shanghai, table-name album_hot_rank, username root, password your_password, sink.buffer-flush.max-rows 100, sink.buffer-flush.interval 2s );然后把计算结果写入 Sink 表INSERT INTO album_sink SELECT album_id, album_name, singer, genre, TUMBLE_START(event_time, INTERVAL 10 SECOND) AS window_start, TUMBLE_END(event_time, INTERVAL 10 SECOND) AS window_end, SUM(play_count) AS total_play, SUM(favorite_count) AS total_favorite FROM album_source GROUP BY TUMBLE(event_time, INTERVAL 10 SECOND), album_id, album_name, singer, genre;这段看起来简单但 JDBC 连接器的坑特别多。我陆续踩过三个一是serverTimezone不配置连接 MySQL 8 时会报时区错误二是PRIMARY KEY如果不声明Sink 写入时容易产生重复数据三是sink.buffer-flush.max-rows和sink.buffer-flush.interval需要按业务场景调节缓冲太小容易频繁写库太大则结果延迟较高。我在用 Flink 1.18 时如果连接器报“No suitable driver found”通常是缺少 MySQL Connector/J 的 JAR 包把mysql-connector-java-8.0.30.jar放进 Flink 的lib目录重启即可。4. 可视化展示与数据联动4.1 可视化方案怎么选数据落库以后展示就变得自由了。我尝试过两种方式在这里直接给出对比方案优点缺点适合场景Superset配置快、仪表盘直接连 MySQL、无需前端基础部署略重、交互样式受限快速出分析大屏ECharts 自写页面灵活、图表好看、可定制交互需要写 HTML/JS需要个性化展示如果你只是想把项目“跑起来看效果”我推荐 Superset。Docker 一行命令启动数据源配置成 MySQL然后拖拽式创建图表非常快。如果你想顺便练一下数据可视化能力那就用 ECharts后端用 Flask 或 Spring Boot 暴露一个简单 API前端定时拉取 MySQL 数据渲染图表。4.2 大屏上放哪些图表我最终的可视化展示页放了三个核心图表专辑热度实时排行榜柱状图展示当前窗口 TOP 10 专辑的总播放量。歌手播放趋势折线图展示每个歌手多个窗口的播放量变化。流派占比饼图或环形图展示不同流派的播放量占比。对应的查询 SQL 也很简单。查实时排行榜SELECT album_name, singer, total_play FROM album_hot_rank WHERE window_start (SELECT MAX(window_start) FROM album_hot_rank) ORDER BY total_play DESC LIMIT 10;查歌手趋势SELECT singer, window_start, SUM(total_play) AS play_sum FROM album_hot_rank GROUP BY singer, window_start ORDER BY window_start;查流派占比SELECT genre, SUM(total_play) AS play_sum FROM album_hot_rank GROUP BY genre;这三条 SQL 覆盖了最常见的分析展示场景。有了实时排行榜、趋势图和占比图你就已经完成了一个“数据实时分析展示”项目该有的样子。更重要的是你会发现从 Flink 计算到 BI 展示整条链路是贯通的数据在实时流动图表也在实时刷新。5. 常见问题与排查技巧实录5.1 JDBC 连接器异常是最大拦路虎我在标题里的热搜词中看到“flink的jdbc连接器异常”被反复搜索说明这个问题真的很多人在踩。我梳理了几个典型的异常场景和解决办法异常信息可能原因解决办法No suitable driver found缺少 MySQL JDBC 驱动包把 mysql-connector-java JAR 放到 Flink lib 目录重启集群Connection refusedMySQL 地址或端口不对或 MySQL 不允许远程连接确认 URL、端口给 MySQL 开启远程访问权限Table not existsSink 表在 MySQL 中没有预先创建先在 MySQL 执行建表语句或配置 Flink 自动建表能力Data truncation数据类型长度不匹配检查 MySQL 字段长度比如 BIGINT 和 INT 对应关系这里重点说一下最后一种情况。Flink 的BIGINT对应 MySQL 的BIGINTTIMESTAMP(3)对应 MySQL 的DATETIME(3)如果你建表时把 MySQL 字段设置成了INT或TIMESTAMP写入时就会出现 Data truncation。我的经验是Sink 表字段类型尽量和 Flink 端保持一致不要依赖隐式转换。5.2 数据一直不更新先查水位线和时间字段这个项目运行以后你可能会遇到一个特别迷惑的现象Flink 作业没有报错但结果表迟迟不更新或者更新得很慢。排查方向主要是两个事件时间字段是否解析正确。如果 Kafka 数据里的event_time是字符串但 Flink 建表时指定为TIMESTAMP(3)解析失败会直接导致数据过滤掉。你可以用ts_proctime或者不用事件时间、改用ProcessingTime先跑通链路再看事件时间和水位线。消费位点是否正确。scan.startup.mode如果设成了latest-offset而 Flink 启动时 Kafka 里没有新数据那作业会一直等新消息。调试阶段建议先用earliest-offset确保消费到已有数据。我在做这个项目时有一次改了字段名但忘了改 JSON 里的 key结果 Flink 一直拿不到play_count字段聚合结果全是 0。检查的时候只盯着 SQL 有没有错完全没意识到是字段映射问题。后来在 Flink Web UI 里查看任务的吞吐量发现 Source 端每秒只有几条数据再用kafka-console-consumer查看消息原始内容才找到问题。5.3 关于 Flink CDC 的扩展玩法最后聊一个延伸点也回应一下热搜词里的“flink cdc”。你如果把模拟数据换成真实的业务库数据经常会遇到“数据在 MySQL 里更新希望 Flink 实时感知”的需求。Flink CDC 就是干这个的它可以直接监听 MySQL binlog把增删改事件转成流你再继续用 Flink SQL 处理。在本项目里可以做一个升级版用 Flink CDC 监听一张专辑信息维表把专辑的基础信息变更实时同步到 Kafka再和用户行为流做关联。这样就能实现“专辑改名了实时生效”“新增专辑自动参与统计”等效果。不过这不属于低难度范围了入门阶段可以先把主链路跑熟练再动手玩 CDC。5.4 低难度项目也要学会看 Flink Web UI无论你用的是本地集群还是独立集群Flink Web UI 都是你最核心的排查工具。它会显示作业的运行状态、算子耗时、反压情况等。我要求在项目里至少做到作业从Running变成Finished或FAILED时能通过 UI 找到失败点和日志入口。有一个小技巧使用 Flink SQL Client 跑INSERT INTO任务时控制台可能看不到自动生成的作业 ID但你在 Web UI 的 “Task Managers” 和 “Metrics” 里能看到负载。如果结果表一直没有数据优先查看各个 Task 的numRecordsIn和numRecordsOut哪一环数据量是明显下降的问题就出在哪一环。写在最后的实际操作体会整个项目做完我最想分享的一个体会是Flink 入门并没有想象中那么难但要“真正理解它的运行逻辑”比“能写出 SQL”重要得多。我见过很多人能把 Flink SQL 背得滚瓜烂熟但一旦作业出问题就完全不知道从哪里入手。原因就在于他们跳过了对数据链路的整体把握一味追求语法正确。而“音乐专辑数据分析展示”这个题目最大的价值不是教会你什么高深的调优技巧而是让你亲手把一条端到端的数据管道走通从模拟数据到 Kafka、再到 Flink 窗口计算、最终变成可视化图表。这个过程走完你对“数据工程”这个词的理解会完全不一样。如果后续你还想扩展这个项目我有三个建议一是把 Kafka 换成 Pulsar体会不同消息中间件的差异二是从 Flink SQL 换成 DataStream API写一段自定义 ProcessFunction加深对算子模型的理解三是把 MySQL 换成 ClickHouse看看即席查询和大屏展示的速度差异。每次只改一个组件通过对比才能理解每个组件在链路中承担的角色。本文还有配套的精品资源点击获取
