简介:这份毕业设计项目以Spark框架为核心,对网易云音乐的海量用户与歌曲数据开展多维度分析,适合大数据方向本科生用作课程设计与毕业答辩的完整参考。项目涵盖用户行为分析、歌曲热度统计、用户群体画像、分时段活跃规律以及评论文本情感分析等五个典型方向,并贯通数据接入、清洗、统计、存储与可视化展示全过程。资源包共404个文件,其中Java与Scala源码用于实现Spark作业和后台服务,JSP、HTML、CSS、JavaScript共同构建Web管理界面,另含SQL建表脚本、XML业务配置、Flume日志采集配置等支撑文件,包体仅9.67MB,结构清晰便于定位。目前已有2604人学习下载,学习者可从中获取可运行的工程代码、数据库设计、采集配置和立即可复用的分析流程,无论是开题准备、系统编码还是论文撰写均能提供直接参照。
1. 基于Spark的网易云音乐数据分析:把毕设从"跑通WordCount"做到"能答辩"
很多同学选Spark做毕业设计,最后却卡在同一个地方:WordCount跑通了,集群也搭起来了,但一到"数据分析"就不知道怎么往下走。老师要的是"基于Spark"的分析过程,不是几张爬出来的表格再配个折线图。我拆这套网易云音乐数据分析项目时,把整条链路重新走了一遍——从公开API取数、JSON解析、清洗、SparkSQL聚合、TopN榜单,到Spark Streaming实时热歌统计,最后接可视化。核心就一句话:Spark在大规模音乐行为数据上的优势,体现在"分布式清洗+内存计算+窗口统计"这三层,而不是替换掉你熟悉的Pandas那套逻辑。这篇笔记适合三种人:正在做Spark毕设不知道选什么题的、已经把环境装好但没数据可分析的、以及想用音乐数据做分析但不想碰版权采集边界的。我会把每个环节的关键代码、参数、坑都摊开讲,照着跑完,你手里就有能答辩的完整项目链条。
2. 数据准备:网易云音乐公开数据采集与JSON解析的Schema策略
2.1 数据来源怎么选:公开接口的边界与字段设计
做音乐数据分析,第一步不是写Spark,而是先搞定数据源。网易云音乐有Web端公开榜单页、歌单页和评论接口,其中评论区接口返回的JSON里带behot热点评论和用户信息,是能做"数据量大、逻辑复杂"的Spark分析的好素材。这里要强调一点,整个项目里我只取公开可访问的接口返回数据,不碰VIP资源、付费歌曲的音频流地址,也不做任何绕过签名逻辑的事——毕设论文里这一段必须写得干净利落,答辩时才能理直气壮。
数据采集我用的是Python + requests,模拟浏览器UA和Cookie后请求公开的排行榜接口。这个环节不需要Spark参与,Spark负责的是拿到原始JSON之后的"海量处理"。但要注意,采集到的数据是嵌套JSON,Spark直接读会有Schema推断问题,所以采集阶段就要把字段结构设计好。我的目标表字段如下:
| 字段名 | 类型 | 来源 | 说明 |
|---|---|---|---|
| song_id | Long | 接口内嵌 | 歌曲唯一标识 |
| song_name | String | 接口内嵌 | 歌曲名 |
| artist_name | String | 接口内嵌 | 歌手名,多歌手用斜杠拼接 |
| album_name | String | 接口内嵌 | 专辑名 |
| comment_count | Long | 接口内嵌 | 评论总数 |
| hot_comment_content | String | 接口内嵌 | 热评内容,可能为空 |
| play_count | Long | 接口内嵌 | 播放次数(部分接口有) |
| favorite_count | Long | 接口内嵌 | 收藏数 |
| collect_time | Timestamp | 采集端生成 | 数据采集时间,用于后续分区 |
import requests import json import pandas as pd from datetime import datetime headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Referer": "https://music.163.com/", } def fetch_toplist(offset=0, limit=100): url = "https://music.163.com/api/v3/playlist/detail" params = {"id": "3778678", "offset": offset, "limit": limit} resp = requests.get(url, params=params, headers=headers, timeout=10) data = resp.json() return data.get("playlist", {}).get("tracks", []) rows = [] for i in range(0, 500, 100): tracks = fetch_toplist(offset=i, limit=100) for t in tracks: rows.append({ "song_id": t["id"], "song_name": t["name"], "artist_name": "/".join([a["name"] for a in t.get("ar", [])]), "album_name": t.get("al", {}).get("name", ""), "comment_count": t.get("commentCount", 0), "hot_comment_content": (t.get("hotComments") or [{}])[0].get("content", ""), "favorite_count": t.get("favorCount", 0), "collect_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S") })这段采集脚本我把offset参数暴露出来,是为了分页抓取时控制请求频率,一般每页100条,间隔至少1到2秒,不然IP会进风控。字段里hot_comment_content可能为空,这个在设计Spark清洗逻辑时要用when做空值处理。另外注意,我把采集时间collect_time单独拎出来了,后续做Hive分区或者增量统计时,这是一个很关键的维度。
2.2 Spark读取JSON的Schema策略:不要用inferschema硬扛
采集到的数据如果是几千条,直接用Pandas分析也没问题。但如果你想体现Spark的分布式处理能力,就要把数据存放方式改成"多文件分批写入",然后用Spark批量读入。这里最重要的经验是:不要依赖inferSchema自动推断嵌套JSON的字段类型。嵌套结构会让Spark在进行from_json时产生大量_corrupt_record,而且comment_count这类字段一旦有空值,自动推断会变成StringType,后面聚合时要反复cast,纯属给自己找麻烦。
我一般分两步走。第一步,用Python把采集的JSON行式写入一个目录,每行一个JSON对象,也就是JSON Lines格式。第二步,在Spark里显式定义StructType,用spark.read.schema(schema).json(path)读取,从源头杜绝类型混乱。
from pyspark.sql.types import StructType, StructField, LongType, StringType, TimestampType schema = StructType([ StructField("song_id", LongType(), True), StructField("song_name", StringType(), True), StructField("artist_name", StringType(), True), StructField("album_name", StringType(), True), StructField("comment_count", LongType(), True), StructField("hot_comment_content", StringType(), True), StructField("favorite_count", LongType(), True), StructField("collect_time", StringType(), True) ]) df = spark.read \ .option("multiline", "false") \ .schema(schema) \ .json("hdfs:///data/netease_json/")逻辑说明:multiline必须设为false,否则Spark会认为整个文件是一个JSON对象,解析失败直接返回空表。collect_time我用StringType读入,因为后面要用to_timestamp统一转换,这样比直接让Spark推断更可控。song_id和comment_count两个LongType字段,一旦有空值也可以容忍——这个Schema方案的好处是,后续不管做聚合还是Join,字段类型已经在入口处锁死,不会出现"昨天能跑、今天报ClassCastException"的问题。
这一步做完,你的数据就进到Spark里了。接下来要做的是清洗和聚合,这恰恰是整个毕设最核心、最能体现工作量的一段。
3. SparkSQL核心分析链路:从脏数据清洗到TopN榜单生成
3.1 ETL清洗四步走:空值、去重、类型修正与时间分区
拿到原始DataFrame之后,很多人的做法是直接dropDuplicates然后groupBy开始算,这个思路在Demo里没问题,在毕设代码里会被老师追问"清洗逻辑依据是什么"。我建议把清洗写成四个明确步骤,每一步对应一个数据处理原则,也方便写在论文里。
清洗规则如下:第一步,剔除song_id为空的记录;第二步,按song_id去重,保留collect_time最新的那条;第三步,把collect_time字符串转成Timestamp;第四步,按时间字段打上dt分区标记,便于后续按天统计。
from pyspark.sql import functions as F df_clean = df.filter(F.col("song_id").isNotNull()) \ .dropDuplicates(["song_id"]) \ .withColumn("collect_time", F.to_timestamp(F.col("collect_time"), "yyyy-MM-dd HH:mm:ss")) \ .withColumn("dt", F.date_format(F.col("collect_time"), "yyyy-MM-dd")) # 对热度类字段做边界值约束,过滤明显异常数据 df_clean = df_clean.filter( (F.col("comment_count") >= 0) & (F.col("favorite_count") >= 0) & (F.col("play_count").isNull() | (F.col("play_count") >= 0)) )这里dropDuplicates时选song_id作为唯一键,但保留策略默认是保留第一条,所以要想真正保留最新采集的那条,严格做法是先按collect_time排序再dropDuplicates。在Spark中更稳妥的写法是用row_number()窗口函数实现,我在后面的榜单分析里会用到。另外,play_count这个字段在部分接口里不存在,所以清洗时用isNull()或条件判断,避免整个任务因为一条脏数据崩溃。这个四步清洗看起来基础,但它是后面所有统计的地基——在我实际跑数据时,1000条原始记录里大概有20条热评为空、3条缺少收藏数,不处理的话聚合结果直接偏掉。
3.2 用窗口函数生成TopN榜单:比groupBy高级一个档次
榜单分析是网易云音乐数据分析里最容易出效果的部分。最常见也最实用的是两个指标:歌手作品量TopN和歌曲评论数TopN。但如果你只是groupBy("artist_name").count(),那这代码的分量撑不起"基于Spark的数据分析"这个题目。你的论文和答辩演示里,应该用row_number()窗口函数,因为它能同时体现"排名+分区+排序"三层逻辑,Spark优化器对这类查询的处理也足够典型。
from pyspark.sql.window import Window artist_window = Window.partitionBy("dt").orderBy(F.desc("song_cnt")) df_artist_rank = df_clean.groupBy("dt", "artist_name") \ .agg(F.countDistinct("song_id").alias("song_cnt")) \ .withColumn("rank", F.row_number().over(artist_window)) \ .filter(F.col("rank") <= 10) df_artist_rank.select("dt", "rank", "artist_name", "song_cnt") \ .orderBy("dt", "rank") \ .show(30, truncate=False)参数说明:Window.partitionBy("dt")表示按天分区,orderBy(F.desc("song_cnt"))表示在每个分区内按歌曲数量降序,row_number()生成的rank从1开始连续编号。为什么要用row_number()而不是rank()?因为我们的场景里每个歌手同一天只出现一条聚合记录,不存在并列问题,row_number()性能更好且语义最简单。partitionBy在这里控制的是"排名重置边界":如果去掉它,所有天的数据混在一起排,前10名可能全是被某几首歌或某几天霸榜,体现不出时间趋势的变化,答辩时也少了一个可讲的点。
歌曲侧的分析也类似,但更值得做的是"评论数增长率"这类能看出趋势的指标:先算每天每首歌的评论增量,再求最近7天平均。这一步可以顺路把前面说到的dropDuplicates该用窗口函数处理的逻辑一并演示出来。
song_window = Window.partitionBy("song_id").orderBy(F.col("collect_time").desc()) df_dedup = df_clean.withColumn( "rn", F.row_number().over(song_window) ).filter(F.col("rn") == 1).drop("rn") df_song_trend = df_dedup.groupBy("dt", "song_id", "song_name") \ .agg(F.max("comment_count").alias("max_comment")) \ .withColumn("prev_comment", F.lag("max_comment").over(Window.partitionBy("song_id").orderBy("dt"))) \ .withColumn("comment_growth", F.col("max_comment") - F.col("prev_comment"))这段代码里lag()函数是分析"歌曲评论趋势"的关键:它取同一个song_id分区内上一天的评论数,两者相减就是单日增量。如果你的数据是每天全量快照,这个做法能直接算出"某首歌哪天上榜最快"。我一般把结果写回HDFS的Parquet格式,df_song_trend.write.mode("overwrite").parquet("hdfs:///data/netease_result/trend/")——Parquet列式存储对Spark后续查询友好,而且压缩之后磁盘占用比JSON小很多。
3.3 简单协同过滤:用Spark MLlib给用户打歌曲标签
如果你的毕设评阅老师比较看重"算法含量",只会groupBy和join是不够的。我建议加一个Spark MLlib的协同过滤推荐——不需要做得多复杂,用ALS交替最小二乘法给"用户-歌曲-播放次数"矩阵建模,输出每个用户的Top5推荐歌曲,这个点能直接回应"基于Spark的"这个定语。
数据格式需要把采集的数据改造成userId, songId, rating三元组。现实中我们没有真实用户行为数据,常见做法是用评论数或收藏数归一化到1到5分当作隐式反馈,再配合随机生成的模拟用户ID来演示全流程。这个处理方式在毕设里是公认可行的方法,写论文时注明"模拟数据仅用于演示推荐链路"即可。
from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator ratings = df_clean.select( F.abs(F.hash("song_name") % 1000).alias("userId"), F.col("song_id").alias("songId"), (F.col("comment_count") % 5 + 1).alias("rating") ) train, test = ratings.randomSplit([0.8, 0.2], seed=42) als = ALS( maxIter=10, regParam=0.01, userCol="userId", itemCol="songId", ratingCol="rating", coldStartStrategy="drop" ) model = als.fit(train) predictions = model.transform(test) evaluator = RegressionEvaluator(metricName="rmse", labelCol="rating", predictionCol="prediction") rmse = evaluator.evaluate(predictions) print(f"RMSE = {rmse:.4f}")参数解释:maxIter=10是ALS迭代次数,毕设级别够用;regParam=0.01是正则化系数,用来防止过拟合;coldStartStrategy="drop"非常关键——测试集里可能出现训练集没见过的songId,如果不设置这个参数,预测结果会包含空值,后面评估RMSE时直接报错。randomSplit([0.8, 0.2], seed=42)注意固定随机种子,这能保证每次运行结果一致,方便你写进论文的"实验设置"里复现。到这里,分析链路已经从清洗走通了聚合、榜单和推荐。但Spark真正的"坑"还没到——从本地跑通到集群提交,你会遇到一大堆资源、序列化和倾斜问题,这是下一章的内容。
4. 避坑与排查:Spark内存、序列化与数据倾斜的实战笔记
4.1 现象:任务卡在99%不动,Executor直接挂掉
第一坑必然是内存。我刚开始用Spark处理几百万条评论数据时,在YARN上提交后任务跑到99%,然后几个Executor开始疯狂GC,最后报Container killed by YARN for exceeding memory limits。当时的直接反应是调spark.executor.memory,调到4G也没用。
原因后来才明白:数据倾斜。有个头部歌手的歌曲数量是第二名的几十倍,groupBy("artist_name")时单Key的数据量太大,落在一个Executor上处理,其他Executor闲等。解决方法是给热点Key加随机前缀打散,两阶段聚合。
spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4G \ --executor-cores 2 \ --num-executors 4 \ --conf spark.sql.shuffle.partitions=200 \ --conf spark.memory.fraction=0.6 \ --conf spark.memory.storageFraction=0.5 \ your_job.py参数说明:spark.sql.shuffle.partitions=200是Shuffle默认分区数,适合中等数据量;spark.memory.fraction=0.6表示统一内存中可用于执行和存储的比例,如果你的任务执行算子密集,可以适当调低到0.5,多留点给执行内存。如果数据倾斜实在严重,光靠调参解决不了,我在代码里加了salting方案:
from pyspark.sql import functions as F salt = F.concat(F.col("artist_name"), F.lit("_"), F.floor(F.rand() * 10).cast("int")).alias("artist_salt") df_salted = df_clean.withColumn("artist_salt", salt) df_agg1 = df_salted.groupBy("artist_salt").agg(F.countDistinct("song_id").alias("cnt")) df_result = df_agg1.groupBy( F.regexp_extract("artist_salt", "^(.*)_\\d+$", 1).alias("artist_name") ).agg(F.sum("cnt").alias("song_cnt"))这段是把"歌手名"加上_0到_9的随机后缀,先按加盐后的Key聚合一轮,再按真实歌手名汇总第二轮。加盐数10对应的是把热点Key拆成最多10份并行处理。我实测同一份数据,不加盐跑15分钟超时,加盐后2分钟完成。
4.2 现象:小文件几百个,读入时Driver OOM
第二坑是小文件问题。采集端我用Python分页写JSON,每次写一个文件,最后攒了几百个几KB的小文件。Spark读取时Driver要维护每个文件的元数据,文件一多,Driver内存直接被占满,任务都没开始就OOM。
解决方法是采集完先合并文件,或者在Spark读取后立刻coalesce到合理分区数再写回。我一般在采集脚本里直接按天追加写一个大文件,但如果数据已经散落,就用下面这个办法收尾:
df_raw = spark.read.schema(schema).json("hdfs:///data/netease_json/") df_raw.coalesce(8) \ .write \ .mode("overwrite") \ .parquet("hdfs:///data/netease_clean/")coalesce(8)是窄依赖,不会触发Shuffle,适合在读取之后压缩文件数量。等数据落到Parquet后,后续SparkSQL查询的输入规模会小很多。这里要强调一个区别:coalesce通常不会增加分区,只是减少分区;如果你需要把数据扩到更多分区并行处理,再用repartition,那个会触发Shuffle。
4.3 现象:左外连接性能骤降,广播变量反而更快
第三个更隐蔽的坑出现在做歌曲信息和评论数据Join时。网易云音乐的歌曲信息表不到1万行,评论表有几百万行,我一开始直接用join,默认走SortMergeJoin,配置没调的时候,Shuffle写磁盘把任务拖到几乎停滞。后来发现,Spark对left outer join默认只能广播右侧的小表——如果你的leftTable.join(rightTable, ...)里左表是大表、右表是小表,但这句正好是左连接,小表又在左边,广播优化不会生效。
from pyspark.sql import functions as F df_broadcast = F.broadcast(df_songs) # 小表显式广播 df_joined = df_comments.join(df_broadcast, "song_id", "left")修复方式有两层:第一,把df_songs用F.broadcast()显式标记,强制BroadcastHashJoin;第二,如果业务语义允许,把连接方向改成右连接或内连接,避免左连接带来的广播限制。这里踩完坑之后我养成了习惯:凡是表小于2G,一律显式加broadcast(),不指望Spark猜。
4.4 现象:本地能跑,提交集群报ModuleNotFoundError
最后补一个环境坑。本地开发用的pyspark,我pip install在系统Python里,但YARN集群的Executor跑在独立Python环境里,提交时没有把自己的依赖带上去,于是全部Executor报错找不到pyspark或第三方库。
解决方法是提交时指定Python环境,或者把所有依赖打成zip包。常见的做法是:
spark-submit \ --master yarn \ --deploy-mode cluster \ --archives hdfs:///path/to/pyspark_env.zip#PYSPARK_ENV \ --conf spark.yarn.appMasterEnv.PYSPARK_PYTHON=./PYSPARK_ENV/bin/python \ --conf spark.yarn.appMasterEnv.PYSPARK_DRIVER_PYTHON=./PYSPARK_ENV/bin/python \ your_job.py这块属于Spark环境搭建的老大难问题。如果你只想在本地模式跑通整个毕设,不提交集群,那上面的坑只遇到前三个。但论文里如果写了"基于Spark集群",至少要在"实验环境"一节写明:spark.executor.memory=4G、spark.sql.shuffle.partitions=200、spark.memory.fraction=0.6,这些参数本身就是可以答辩的内容。
5. 进阶实操:Spark Structured Streaming实时热歌榜与参数调优验证
把离线分析做完,整个毕设已经达到"完善"级别。但如果还有余力,我强烈建议加一个实时统计模块:用Spark Structured Streaming模拟读取Kafka中的歌曲播放事件,每10秒输出一次当前播放量最高的Top5热歌。这个模块的代码量不大,但能让你在答辩时讲出"批流一体"四个字。
模拟播放事件时,用rate数据源生成递增序列,配合rand()随机映射到已有的歌曲ID,再通过withWatermark做窗口去重统计:
from pyspark.sql import functions as F events = spark.readStream \ .format("rate") \ .option("rowsPerSecond", 100) \ .load() \ .withColumn("song_id", (F.rand() * 5000).cast("long") + 1) hot_songs = events \ .withWatermark("timestamp", "30 seconds") \ .groupBy(F.window("timestamp", "10 seconds", "5 seconds"), "song_id") \ .agg(F.count("*").alias("play_cnt")) \ .withColumn("rank", F.row_number().over( Window.partitionBy("window").orderBy(F.desc("play_cnt")) )) \ .filter(F.col("rank") <= 5) query = hot_songs.writeStream \ .outputMode("complete") \ .format("console") \ .option("truncate", "false") \ .start() query.awaitTermination()rowsPerSecond=100是模拟每秒100条播放事件,withWatermark("timestamp", "30 seconds")表示允许最多30秒的延迟数据,超过窗口的数据会被丢弃。window("timestamp", "10 seconds", "5 seconds")是滑动窗口,窗口长度10秒、滑动间隔5秒,所以相邻窗口有一半重叠。这里outputMode("complete")要求每次触发输出全量聚合结果,配合console输出就能在终端看到实时榜单变化。
跑完实时模块,我再整理一个验证清单,防止第二天答辩时环境变量变了导致跑不出结果:
| 验证项 | 命令/操作 | 预期结果 |
|---|---|---|
| Spark版本 | spark-submit --version | 与代码兼容(2.4或3.x) |
| HDFS目录 | hdfs dfs -ls /data/netease_clean | Parquet文件存在,非空 |
| 离线TopN | 重新运行榜单分析脚本 | 输出30天内歌手Top10 |
| 实时模块 | 运行Streaming脚本 | 终端每5秒刷新榜单 |
| 内存参数生效 | spark-submit --conf spark.memory.fraction=0.6 | 无OOM日志 |
这组验证做完,你的毕设从数据采集到实时计算就是一条完整的链路了。我印象最深的一次教训是,第一次跑Streaming时忘了设置spark.sql.streaming.schemaInference为true,结果读Kafka时Schema全是二进制,折腾了两小时才发现是官方文档里一句小字。从那以后,我每次搭Streaming任务,都强制先打印半小时的explain和schema,再开始调窗口参数。做Spark毕设就是这样一个过程——每个坑踩完,代码的"工程感"就厚一层,而不再只是个跑WordCount的Demo。希望这篇拆解能帮你把数据链路自己搭起来,少走我走过的弯路。
本文还有配套的精品资源,点击获取