
简介这份毕业设计项目以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之后的海量处理。但要注意采集到的数据是嵌套JSONSpark直接读会有Schema推断问题所以采集阶段就要把字段结构设计好。我的目标表字段如下字段名类型来源说明song_idLong接口内嵌歌曲唯一标识song_nameString接口内嵌歌曲名artist_nameString接口内嵌歌手名多歌手用斜杠拼接album_nameString接口内嵌专辑名comment_countLong接口内嵌评论总数hot_comment_contentString接口内嵌热评内容可能为空play_countLong接口内嵌播放次数部分接口有favorite_countLong接口内嵌收藏数collect_timeTimestamp采集端生成数据采集时间用于后续分区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(offset0, limit100): url https://music.163.com/api/v3/playlist/detail params {id: 3778678, offset: offset, limit: limit} resp requests.get(url, paramsparams, headersheaders, timeout10) data resp.json() return data.get(playlist, {}).get(tracks, []) rows [] for i in range(0, 500, 100): tracks fetch_toplist(offseti, limit100) 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, truncateFalse)参数说明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], seed42) als ALS( maxIter10, regParam0.01, userColuserId, itemColsongId, ratingColrating, coldStartStrategydrop ) model als.fit(train) predictions model.transform(test) evaluator RegressionEvaluator(metricNamermse, labelColrating, predictionColprediction) rmse evaluator.evaluate(predictions) print(fRMSE {rmse:.4f})参数解释maxIter10是ALS迭代次数毕设级别够用regParam0.01是正则化系数用来防止过拟合coldStartStrategydrop非常关键——测试集里可能出现训练集没见过的songId如果不设置这个参数预测结果会包含空值后面评估RMSE时直接报错。randomSplit([0.8, 0.2], seed42)注意固定随机种子这能保证每次运行结果一致方便你写进论文的实验设置里复现。到这里分析链路已经从清洗走通了聚合、榜单和推荐。但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.partitions200 \ --conf spark.memory.fraction0.6 \ --conf spark.memory.storageFraction0.5 \ your_job.py参数说明spark.sql.shuffle.partitions200是Shuffle默认分区数适合中等数据量spark.memory.fraction0.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.memory4G、spark.sql.shuffle.partitions200、spark.memory.fraction0.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()rowsPerSecond100是模拟每秒100条播放事件withWatermark(timestamp, 30 seconds)表示允许最多30秒的延迟数据超过窗口的数据会被丢弃。window(timestamp, 10 seconds, 5 seconds)是滑动窗口窗口长度10秒、滑动间隔5秒所以相邻窗口有一半重叠。这里outputMode(complete)要求每次触发输出全量聚合结果配合console输出就能在终端看到实时榜单变化。跑完实时模块我再整理一个验证清单防止第二天答辩时环境变量变了导致跑不出结果验证项命令/操作预期结果Spark版本spark-submit --version与代码兼容2.4或3.xHDFS目录hdfs dfs -ls /data/netease_cleanParquet文件存在非空离线TopN重新运行榜单分析脚本输出30天内歌手Top10实时模块运行Streaming脚本终端每5秒刷新榜单内存参数生效spark-submit --conf spark.memory.fraction0.6无OOM日志这组验证做完你的毕设从数据采集到实时计算就是一条完整的链路了。我印象最深的一次教训是第一次跑Streaming时忘了设置spark.sql.streaming.schemaInference为true结果读Kafka时Schema全是二进制折腾了两小时才发现是官方文档里一句小字。从那以后我每次搭Streaming任务都强制先打印半小时的explain和schema再开始调窗口参数。做Spark毕设就是这样一个过程——每个坑踩完代码的工程感就厚一层而不再只是个跑WordCount的Demo。希望这篇拆解能帮你把数据链路自己搭起来少走我走过的弯路。本文还有配套的精品资源点击获取