简介:一份基于Spark的网易云音乐数据分析毕业设计项目,面向大数据相关专业学生与入门开发者,解决从海量用户行为数据采集、清洗、分析到可视化展示的完整流程问题。项目覆盖Spark核心API、Spark SQL与MLlib、Hadoop分布式存储、MySQL结果落地、Azkaban任务调度以及Nest/ECharts可视化等环节,同时涉及推荐系统建模思路,适合用于课程设计、毕业设计或简历项目对照。压缩包共403个文件,以java、scala源码为主体,js、html/jsp构建前端展示,xml/properties承担配置管理,png/jpg保存图表结果,sql/csv/dat存放数据与脚本,整体约9.31MB。已有485人学习下载。内容包括可运行的Spark分析代码、可视化网页模板、Flume/Hadoop配置示例、Azkaban工作流描述、MySQL建表SQL及readme说明文档,目录结构清晰,便于按数据层、计算层、展示层快速定位与二次开发。
1. 网易云音乐日志的全链路分析:毕业设计里Spark到底在算什么
打开这份「毕业设计基于Spark网易云音乐数据分析.zip」,最先看到的不是一堆Scala源码,而是fontawesome-webfont.woff2、amazeui.min.css、bootstrap.min.css这些前端静态资源,外加flume-hdfs-ng.conf、log4j-es.conf、Azkaban_1.png、MySQL_1.png这类配置和架构图文件。这说明它不是一个“调一下Spark API出个结果”的Demo,而是一条从日志采集、分布式存储、离线批处理到MySQL落库、Echarts可视化的完整数据管道。网易云音乐每天会产生大量用户行为日志,比如播放、收藏、搜索、评论,数据量一旦达到亿级,单机SQL和Excel就扛不住了。这个毕业设计的核心价值,是演示如何用Flume接日志、用HDFS做原始存储、用Spark SQL做清洗和指标计算,再用Azkaban把整条链路定时串起来。适合正在做大数据方向毕设、想搞清楚“Spark在真实项目里到底算什么”的同学,也适合工作后想快速回顾离线数仓全流程的工程师。它解决的最关键问题不是某个算法多难,而是数据从哪来、怎么算、算完放哪、怎么展示。
2. Flume采集与HDFS落地:让Spark拿到干净的原始日志
2.1 为什么选用Flume而不是Kafka
项目中出现了flume-hdfs-ng.conf,这个文件名暴露了技术选型:Flume 1.x(ng版本)负责日志采集,sink目标是HDFS。很多人在设计这类系统时第一反应是上Kafka,但Kafka的引入意味着还要维护Zookeeper集群、控制consumer offset、处理消息堆积监控,对毕业设计这种规模的集群来说运维成本偏高。Flume的agent本身就是Java进程,配置一个spooldir或taildir source就能监控日志目录,sink直接写HDFS,链路短、排查直观。常见做法是单节点部署一个Flume agent,source监控网易云音乐埋点日志的输出目录,channel用memory,sink到HDFS按天分目录存储。
这套设计的前提是离线T+1分析,而不是实时推荐。如果是实时场景,比如用户听歌后立刻更新推荐列表,那确实应该用Kafka加Spark Structured Streaming。但毕业设计里的指标,比如日活、热门歌曲TopN、24小时播放趋势,延迟一天完全可接受。Flume落地HDFS的另一个好处是原始日志不丢,后续重跑任务时可以从HDFS重新读取,不需要回溯业务系统。
2.2 flume-hdfs-ng.conf核心配置解读
项目里的flume-hdfs-ng.conf应该是整个采集层的核心,配置内容一般是agent三要素:source、channel、sink。我基于这个场景还原了一份典型配置:
# flume-hdfs-ng.conf agent.sources = songSrc agent.channels = memChannel agent.sinks = hdfsSink # 监控日志落地目录,处理完改后缀,避免重复读取 agent.sources.songSrc.type = spooldir agent.sources.songSrc.spoolDir = /data/music_logs agent.sources.songSrc.fileSuffix = .DONE agent.sources.songSrc.deletePolicy = never agent.sources.songSrc.ignorePattern = ^.*\.DONE$ # memory channel,容量根据单日日志峰值估算 agent.channels.memChannel.type = memory agent.channels.memChannel.capacity = 10000 agent.channels.memChannel.transactionCapacity = 1000 # sink 写 HDFS,按天分目录 agent.sinks.hdfsSink.type = hdfs agent.sinks.hdfsSink.hdfs.path = hdfs://namenode:8020/music/raw/%Y%m%d agent.sinks.hdfsSink.hdfs.fileType = DataStream agent.sinks.hdfsSink.hdfs.writeFormat = Text agent.sinks.hdfsSink.hdfs.rollInterval = 600 agent.sinks.hdfsSink.hdfs.rollSize = 134217728 agent.sinks.hdfsSink.hdfs.rollCount = 0 agent.sinks.hdfsSink.hdfs.filePrefix = music agent.sinks.hdfsSink.hdfs.fileSuffix = .log配置里有几个参数要特别说明。spooldir适合目录内文件不再变动的场景,Flume会为每个文件维护一个.flumespool元数据,断点续传靠它实现;fileSuffix = .DONE是防止同一个文件被反复消费的常用手段,处理完成后源文件被改名但不会被删除,这比deletePolicy = immediate安全,出问题时还能找回原始日志。memory channel的capacity是channel中最多能放多少event,transactionCapacity是每个事务最多取多少event,前者必须大于后者,否则启动时会直接报错。HDFS sink里的rollInterval = 600表示每600秒强制落盘一个文件,rollSize = 134217728(128MB)是文件大小滚动阈值,两者谁先触发都行。这里要注意如果rollCount = 0,就表示不按event条数滚动,避免小文件刷爆NameNode。
HDFS路径里用了%Y%m%d,这个时间取自event header,默认是Flume服务器本地时间。按天分目录的好处很直接:Spark读取时可以精确到某一天的目录作为输入,做分区裁剪,不需要全表扫描。这里有一个建议:如果后续要做小时级分析,可以把path改成%Y%m%d/%H,但这样会产生更多小文件,Spark读取时反而变慢,所以T+1场景按天够用了。
2.3 整条链路的组件职责边界
从这份毕业设计涉及的文件来倒推完整链路,各层职责大致是这样的:
| 层级 | 组件 | 职责 | 关键产物 |
|---|---|---|---|
| 采集层 | Flume | 监控埋点日志目录,写入HDFS | /music/raw/yyyyMMdd/xxx.log |
| 存储层 | HDFS | 保存原始日志和中间结果 | 按天分区的原始数据 |
| 计算层 | Spark | 清洗、去重、聚合、TopN | 指标结果DataFrame |
| 调度层 | Azkaban | 定时触发Spark作业和导出作业 | 工作流DAG |
| 结果层 | MySQL | 存储计算结果供前端查询 | 指标表、榜单表 |
| 展示层 | Echarts/AmazeUI | 读取MySQL数据渲染图表 | 折线图、柱状图、词云 |
这个表格对应了压缩包里那些图片和配置文件背后的设计意图。Hadoop_1.png说明集群环境是Hadoop 2.x加Spark on YARN部署,Azkaban_1.png说明不是手动spark-submit,而是通过Azkaban调度,MySQL_1.png说明结果数据最终是结构化存储。这套组合是2018到2022年间大数据毕业设计里最主流的架构,现在看依然适合教学演示,因为它把离线数仓的每个环节都覆盖到了,但又没有引入Kafka、Iceberg这类对毕设来说过重的组件。
3. Spark SQL清洗与指标计算:去重、TopN与时序统计的SQL写法
3.1 原始日志的典型字段结构与埋点格式
网易云音乐的日志通常包含用户ID、歌曲ID、行为类型、时间戳、设备信息等字段。下面是一份常见的埋点日志schema,实际项目中可能还有更多维度,但对毕业设计来说,这几个字段已经能支撑80%的指标计算。
| 字段名 | 类型 | 示例值 | 说明 |
|---|---|---|---|
| user_id | String | 8374291 | 匿名用户可能为空或为设备ID |
| song_id | String | 27438123 | 歌曲唯一标识 |
| action | String | play / collect / download | 行为类型,play为主 |
| ts | Long | 1715846400000 | 毫秒级Unix时间戳 |
| device_type | String | Android / iOS / Web | 设备端 |
| dt | String | 20240520 | 分区字段,跟目录对应 |
日志进入HDFS时往往是原始文本,字段用制表符或逗号分隔,甚至混入一些Nginx上报的异常行。Spark SQL的DataFrame API负责把这些文本解析成结构化表。第一步是read加载后按分隔符split,再用to_timestamp把毫秒时间戳转成可读时间格式。这里有一个常见的坑:不要用__HIVE_DEFAULT_PARTITION__这种默认分区名去join,得先把dt字符串校验一遍,否则脏数据会把整张表带偏。
3.2 清洗逻辑:去重、过滤、类型转换
用PySpark写清洗逻辑最直观,也方便在Jupyter里跑通后再改成Scala封装进JAR。核心代码如下:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_timestamp, from_unixtime, split spark = SparkSession.builder \ .appName("music-etl") \ .enableHiveSupport() \ .getOrCreate() # 按天读取HDFS原始日志目录 raw_df = spark.read.text("/music/raw/20240520") # 切分字段并转类型 parsed_df = raw_df.withColumn("user_id", split(col("value"), "\t").getItem(0)) \ .withColumn("song_id", split(col("value"), "\t").getItem(1)) \ .withColumn("action", split(col("value"), "\t").getItem(2)) \ .withColumn("ts", split(col("value"), "\t").getItem(3).cast("long")) \ .withColumn("device_type", split(col("value"), "\t").getItem(4)) \ .withColumn("event_time", from_unixtime(col("ts") / 1000, "yyyy-MM-dd HH:mm:ss")) # 过滤:user_id为空、ts异常、action不在白名单 clean_df = parsed_df.filter( col("user_id").isNotNull() & (col("user_id") != "") & col("ts").isNotNull() & col("action").isin("play", "collect", "download", "search") ) # 去重:同一个人同一秒对同一首歌的重复上报只保留一条 dedup_df = clean_df.dropDuplicates(["user_id", "song_id", "action", "ts"]) dedup_df.createOrReplaceTempView("music_event")这里的split(col("value"), "\t")是按Tab切分,如果日志实际是逗号分隔就需要替换。from_unixtime(col("ts") / 1000)把毫秒时间戳转成秒再格式化,注意如果不除以1000,格式化出来的时间会是1970年附近。dropDuplicates的粒度是用户、歌曲、行为、时间戳四元组,不是全字段去重,因为device_type等字段可能上报不一致,全字段去重会留下重复统计的隐患。清洗后创建临时视图music_event,后面所有指标都基于它计算。
这里的createOrReplaceTempView是Session级别的临时表,作业结束就释放,不会污染Hive元数据。如果需要跨SparkSession复用结果,可以把清洗后的数据写成Parquet格式到HDFS的中间目录,比如/music/cleaned/dt=20240520,这也是离线数仓里“清洗层”的标准做法。
3.3 三个核心指标:TopN榜单、24小时趋势、用户活跃度
指标计算是Spark SQL最出彩的部分。直接上SQL:
-- 1. 播放量Top20歌曲榜单 SELECT song_id, COUNT(*) AS play_cnt FROM music_event WHERE action = 'play' GROUP BY song_id ORDER BY play_cnt DESC LIMIT 20; -- 2. 24小时播放量趋势 SELECT HOUR(event_time) AS hour, COUNT(*) AS play_cnt FROM music_event WHERE action = 'play' GROUP BY HOUR(event_time) ORDER BY hour; -- 3. 用户播放行为Top10(用窗口函数替代全局排序) SELECT user_id, song_id, play_cnt FROM ( SELECT user_id, song_id, COUNT(*) AS play_cnt, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY COUNT(*) DESC) AS rn FROM music_event WHERE action = 'play' GROUP BY user_id, song_id ) t WHERE rn <= 10;HOUR(event_time)是从Spark 2.0开始内置的时间函数,返回0到23的整数,用于小时粒度聚合不需要额外UDF。第三个查询里的ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY COUNT(*) DESC)是标准SQL窗口函数,其执行效率取决于group by之后的数据分布,如果热门用户听歌量巨大,这里会触发shuffle。实际执行时Spark会把整个查询翻译成物理计划,窗口函数的排序操作默认用的全局排序器,对内存压力较大,跑不过去时优先调spark.sql.shuffle.partitions,而不是盲目加executor内存。
这三个指标基本就是毕业设计可视化页面的主图:柱状图显示热门歌曲、折线图显示24小时趋势、表格或嵌套饼图显示用户偏好。项目中的nest_1.png、nest_2.png应该就是用嵌套饼图展示不同维度的占比情况。
3.4 Spark运行参数如何影响这批SQL
同样一段Spark SQL,在不同参数配置下运行时长可能差好几倍。下面是这组指标最常用的一组spark-submit参数:
spark-submit \ --master yarn \ --deploy-mode client \ --executor-memory 4g \ --num-executors 4 \ --executor-cores 2 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.sql.autoBroadcastJoinThreshold=10485760 \ --class com.music.MusicAnalysis \ music-analysis.jarspark.sql.shuffle.partitions决定了group by或join时产生的reduce任务数,默认200。如果输入数据只有几百MB,200个分区意味着每个分区只处理几MB数据,空转开销很大;如果数据有几十GB,200个分区又不够,单个任务处理太久。常见做法是先看Spark UI里Shuffle Read和Shuffle Write的字节数,再估算一个合适的分区数。spark.sql.autoBroadcastJoinThreshold表示小于10MB的表自动广播到每个executor,避免shuffle join,这个毕业设计里如果需要把歌曲维度表join进来补全歌名和歌手,这个参数就能派上用场。
4. 结果落库MySQL与Echarts展示:Azkaban把整条链路自动跑起来
4.1 DataFrame写出MySQL的两种模式与坑
计算结果不能一直躺在Spark里,最终要给前端查询。项目里出现了MySQL_1.png,说明结果数据是落到MySQL的。写出代码用DataFrame的jdbc接口即可:
result_df.write \ .mode("overwrite") \ .format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/music_analysis?useUnicode=true&characterEncoding=utf8&useSSL=false") \ .option("dbtable", "top_song_rank") \ .option("user", "root") \ .option("password", "123456") \ .option("batchsize", "1000") \ .option("truncate", "true") \ .save()这里有两个坑值得展开。第一个是mode("overwrite")配合truncate=true,Spark默认overwrite是先drop表再重建,如果表结构变了会出问题,而truncate只是清空数据保留表结构,更安全。第二个是URL里必须带useUnicode=true&characterEncoding=utf8,否则中文歌名写入MySQL会变成问号。batchsize=1000表示每批写入1000条,太大会导致MySQL连接超时,太小则写入慢。写完后推荐在MySQL里对song_id和dt建立联合索引,因为前端查询基本都带WHERE dt = ?条件,没有索引的话,TopN榜单查询在数据量大了之后会明显变慢。
对于每日榜单这种结果表,用overwrite是合理的,因为每天只保留当天结果;对于用户行为明细表,则应该用append,并加上dt字段做分区标识,便于回溯。两种模式混用时要小心,同一张表不要一会儿append一会儿overwrite,否则容易把历史数据搞丢。
4.2 Echarts可视化:后端JSON怎么交给前端渲染
项目里的bootstrap.min.css、amazeui.min.css、font-awesome.css这些资源,说明前端采用了AmazeUI响应式框架加Echarts图表库。Echarts从MySQL拿数通常不是直连,而是后端接口返回JSON再渲染。常见做法是后端写一个Servlet或Spring Boot接口,查询top_song_rank表,转成以下JSON结构:
{ "code": 0, "data": { "songs": ["晴天", "稻香", "七里香"], "playCounts": [98234, 87321, 80234] } }前端拿到数据后用setOption填充:
// echarts 柱状图:每日播放Top10歌曲 $.getJSON("/api/top_song", { dt: "20240520" }, function (res) { if (res.code !== 0) return; var chart = echarts.init(document.getElementById("topSongChart")); chart.setOption({ tooltip: { trigger: "axis" }, xAxis: { type: "category", data: res.data.songs }, yAxis: { type: "value" }, series: [{ type: "bar", data: res.data.playCounts, itemStyle: { color: "#e74c3c" } }] }); });这段代码里的xAxis.data和series.data分别对应榜单名称和播放量,核心逻辑是数据从MySQL到JSON再到Echarts实例的映射。tooltip.trigger = 'axis'表示鼠标悬停在坐标轴附近时展示提示框,柱状图适合这种模式。如果要做24小时趋势折线图,只要把接口换成按小时聚合的数据源,然后type: "line"即可。AmazeUI的作用是页面的栅格布局和响应式适配,让图表在PC和手机上都能看,这在毕设答辩时用平板展示是个加分项。
4.3 Azkaban调度:把ETL、指标计算、导出串成DAG
项目里有Azkaban_1.png,说明不是手动执行spark-submit,而是用Azkaban定时触发。Azkaban提交的是一个zip包,里面包含多个.job文件,每个job定义一条命令,job之间用dependencies声明依赖关系。典型的项目结构如下:
music_analysis/ ├── project.job ├── spark_etl.job ├── mysql_export.job └── chart_report.job# spark_etl.job type=command command=spark-submit --master yarn --deploy-mode client --executor-memory 4g --num-executors 4 --executor-cores 2 --conf spark.sql.shuffle.partitions=200 --class com.music.MusicAnalysis music-analysis.jar --date ${dt}# mysql_export.job type=command dependencies=spark_etl command=/opt/mysql/export_to_mysql.sh --date ${dt}Azkaban会按dependencies的依赖关系生成DAG,spark_etl跑完才执行mysql_export。这里的${dt}是Azkaban的调度参数,可以在创建Project时配置为20240520这种格式,也可以用cron表达式每天自动替换。有一点需要注意:deploy-mode client要求提交spark-submit的机器和Azkaban executor在同一台节点上,而且要求该机器配置了Hadoop和Spark客户端;如果集群是普通用户启动的,还要确保Azkaban的executor用户有HDFS写入权限,否则提交作业后会在YARN上报错。
这个调度流程解决的核心问题是“重跑”。某天数据清洗逻辑改了一行代码,需要重算过去一周的数据,如果没有Azkaban,你只能手动执行七次spark-submit,有了它只需要把调度周期改一下或者手动触发带不同${dt}参数的重跑任务即可。
5. 数据倾斜与Shuffle优化:几十万歌单聚合时的排错顺序
网易云音乐的歌曲热门程度呈典型的幂律分布,少数头部歌曲贡献了绝大部分播放量。这会导致GROUP BY song_id时,某一个或某几个Reducetask处理的数据量远超其他task,也就是数据倾斜。具体表现是Spark UI里大部分task几十秒跑完,剩下一两个task卡了十几分钟,甚至报OOM。这是这类音乐数据分析项目里最容易踩、也最容易在答辩时被问到的性能问题。
先讲排查顺序。打开Spark UI的Stages页面,看Summary Metrics里的Duration列,如果Max和Median差了一个数量级以上,基本可以判定倾斜。再用鼠标点开那个最慢的task,看Shuffle Read Size是不是远大于其他task。如果是,那就要针对倾斜的key做处理。
针对热门歌曲倾斜的经典方案是两阶段聚合加盐。所谓加盐,就是给原来的key拼接一个随机数前缀,把一个大key拆成多个小key,完成局部聚合后再去掉前缀做全局聚合。代码示例如下:
from pyspark.sql.functions import rand, concat, lit, split, col, sum as fsum # 设置分区数,让reduce端有足够并行度 spark.conf.set("spark.sql.shuffle.partitions", "300") # 第一阶段:加随机盐,拆key salted_df = df.filter(col("action") == "play") \ .withColumn("salt", (rand() * 10).cast("int")) \ .withColumn("salted_song_id", concat(col("song_id"), lit("_"), col("salt"))) # 局部聚合:同一个盐内先聚合一次 partial_agg = salted_df.groupBy("salted_song_id") \ .agg(fsum(lit(1)).alias("cnt")) # 第二阶段:去掉盐,对局部结果再聚合 final_result = partial_agg.withColumn( "song_id", split(col("salted_song_id"), "_").getItem(0) ) \ .groupBy("song_id") \ .agg(fsum("cnt").alias("play_cnt")) \ .orderBy(col("play_cnt").desc()) \ .limit(20)代码里rand() * 10生成0到9的随机整数,拼在song_id后面,让同一首歌的播放记录被打散到10个不同的key上,reduce端就有10个task分担压力。这里的10是盐的粒度,粒度太小效果不明显,太大会产生过多中间文件,常见做法是先试10,看最长task耗时降了多少,再逐步调到50。第二段代码里split(col("salted_song_id"), "_").getItem(0)是对字段切割,取盐之前的原始song_id。注意fsum(lit(1))比count("*")在部分场景下更可控,因为它的返回值类型是Long且不受null值干扰。
这个方案有一个前提:去盐后的二次聚合数据量已经被大幅压缩,否则二次groupBy仍然可能倾斜。如果倾斜的key不多,还有一种更轻的做法是加一个WHERE song_id NOT IN (热门key列表)的过滤,把热门key单独计算再union回结果,但这种方式需要维护一个动态的热门key清单,在这个场景下不如加盐通用。
再补充一个和Shuffle直接相关的内存参数。之前提到spark.sql.shuffle.partitions,但它只影响shuffle输出的分区数,不直接影响单个task的内存上限。真正和OOM相关的是spark.executor.memoryOverhead,它表示每个executor在堆外内存之外额外预留的内存,默认只有executor-memory的10%。如果你用--executor-memory 4g,overhead只有约400MB,当Shuffle spill严重时这个值不够用,建议显式设置为1g到2g。判断依据还是Spark UI里Executor页面的GC时间,如果GC TIME超过总运行时间的10%,说明堆内存也紧张,这时候应该降低executor并发task数,而不是盲目加内存。
最有效的调优操作其实是复现一次完整流程,记录每次调参前后Spark UI里Shuffle Spill和Task Duration两列的变化。先把spark.sql.shuffle.partitions从默认200翻三倍,观察Shuffle Read和Spill是否下降;再把spark.executor.memoryOverhead调高到1g,看GC耗时是否回落。这一步对任何基于Spark做用户行为数据分析的项目都适用,也是这份毕业设计里最值得单独拿出来讲的一页PPT。
本文还有配套的精品资源,点击获取