☰
基于Spark的网易云音乐数据分析:从数据采集到可视化看板全流程实践
2026/10/12 1:10:25 网站建设 项目流程

简介:这是一份基于Apache Spark的网易云音乐数据分析毕业设计源码包,面向大数据方向高年级学生、毕业设计开发者,以及想快速上手Spark分析流程的读者。项目覆盖数据采集、清洗、处理、分析和可视化全流程,并包含可运行的核心代码。压缩包共402个文件,约10.99MB,主要含123个Java与19个Scala源码文件、27个JSP页面、56个JS及丰富的CSS/HTML前端资源、42张效果图,另配properties/conf配置文件与SQL脚本,目录结构完整。资源附有项目说明文档、效果图展示,以及情感分析子模块,便于理解如何对音乐评论和用户行为做情绪挖掘;Flume、Log4j等配置示例展示了大数据环境中的数据接入与日志处理方式。已有4126人学习下载,适合作为毕设参考、实战练习或二次开发的起点。

1. 这类毕业设计真正卡人的不是 Spark,而是前面的数据准备

把“基于 Spark 网易云音乐数据分析”当成一个普通的 Spark 入门项目来做,十有八九会卡在第一步:没有现成的数据集。公开渠道能拿到的网易云音乐数据是零散的接口响应,字段命名混乱、时间戳秒毫秒混用、歌词和评论里塞满转义字符,这些问题任何一个都足以让新手在环境搭建之外多耗两周。而 Spark 本身恰恰是整个链路里最不稀奇的部分:读 JSON、过滤脏数据、做分组聚合、写回结果,这些都是 DataFrame API 的常规操作。

所以这篇博文把网易云音乐数据分析当成一个完整的数据工程小闭环来讲:从公开接口取数、设计落地格式、用 Spark 做清洗与聚合,最后落到可视化和验证方法。适合正在做毕业设计、想把“数据平台 + 业务分析”写成完整故事的人,也适合想借一个真实领域把 Spark 从跑通 demo 推到能交付状态的一线工程师。核心思路是先解决数据能不能用,再谈分析好不好看。

2. 基于 Spark 网易云音乐数据分析的链路设计:从接口到看板

2.1 抓数据:网易云音乐公开入口与请求整形

网易云音乐的网页端和客户端都依赖一组 HTTP 接口,歌曲详情、歌手信息、评论列表都有对应的开放路由。直接请求这些接口可以拿到 JSON 响应,不需要逆向客户端、也不需要模拟登录态,但有几个前置条件:请求头里必须带上常见的 User-Agent,Cookie 里至少要有NMTID这类匿名标识,否则部分接口会返回-460错误码。更重要的一个约束是频率:接口没有公开的限流文档,但高频请求会触发风控,建议每分钟控制在 30 到 60 个请求之间。

歌曲元数据推荐从歌单接口切入。一个歌单会返回完整的曲目列表,每个曲目都携带id、name、ar(歌手数组)、al(专辑信息)、dt(时长)等字段,省去了用关键词搜索再拼装数据的麻烦。下面的脚本用一个简单的uid和歌单id集合抓取原始 JSON,并把每个接口响应原样写到本地文件:

import requests import json import time import os def fetch_playlist(playlist_id, cookie="NMTID=xxx; MUSIC_U=xxx", sleep_sec=1.5): url = f"https://music.163.com/api/v6/playlist/detail?id={playlist_id}" headers = { "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36", "Cookie": cookie, "Referer": "https://music.163.com/", } resp = requests.get(url, headers=headers, timeout=10) if resp.status_code != 200: return None return resp.json() os.makedirs("raw_playlist", exist_ok=True) for pid in [3778678, 3779629, 2884035]: data = fetch_playlist(pid) if data: with open(f"raw_playlist/{pid}.json", "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False) time.sleep(2)

这段代码做了三件必要的事:把playlist_id拼进接口路径、把Cookie写在请求头里、把每次请求的响应完整落盘。参数上最重要的不是超时时间,而是sleep_sec——不同歌单的请求间隔至少要留出 1 秒以上,连续快速请求非常容易触发临时封禁。这里刻意不做解析,因为原始响应里还嵌套着推荐语、标签、创建者信息,后续清洗时统一处理比抓取时逐个抽取更高效。

落盘格式上,建议一行为一个 JSON 对象,而不是把整个歌单存成一个 JSON 数组。行式 JSON 能让 Spark 的read.json直接以multiLine=false的方式读入,未来如果有多个歌单增量抓取,追加写入也不会破坏格式。评论数据同理,把每一个评论作为一行单独落盘,字段名保持和接口返回一致。

2.2 原始数据的落地格式:先存 JSON 还是先入表

毕设场景里经常出现“抓完数据先导进 MySQL”的冲动,从结果看这是给自己加工作量。网易云音乐的接口返回是嵌套 JSON,歌手不是字符串而是对象数组,专辑信息也是对象,硬塞进关系表意味着抓取阶段就要做 3 到 4 张表的拆解设计;反过来,如果只做数据分析,很多嵌套字段根本用不到。

常见做法是按“原始层 → 清洗层 → 分析层”分层落地,原始层全部存行式 JSON。目录结构按采集日期分区,例如raw_playlist/2025-06-01/3778678.json,这样后续增量抓取只需要添加新目录,不需要改表结构。分析层的数据则从 JSON 里挑出核心字段转成 Parquet 格式,Parquet 的列式存储配合 Spark 做聚合时扫描的数据量远小于 JSON。

这个阶段还有一个容易忽略的问题:接口返回的字段名不是稳定的。歌单详情里的评论数曾经叫commentCount,有些旧的缓存响应里则叫comment_count。JSON 落地方案的好处在这里体现出来——字段名变化不会导致入库失败,清洗层统一做字段映射即可,不用回改采集脚本。

2.3 ETL 入湖:清洗、格式标准化与 ID 统一

拿到原始 JSON 后,Spark 的活才正式开始。清洗层要处理四类典型问题:接口返回了null节点导致整个对象解析失败、时间戳单位不统一(歌曲时长dt是毫秒,评论时间有时返回毫秒有时返回秒)、歌手以 “群星/Various Artists” 形式出现无法直接 join、以及同一首歌在不同歌单里重复出现。

下面这段 PySpark 代码把行式 JSON 读进来,完成去重、类型转换和脏数据过滤,最后落成按日期分区的 Parquet 表:

from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date, from_unixtime, row_number from pyspark.sql.window import Window spark = SparkSession.builder.appName("ncm_etl").getOrCreate() df = spark.read.json("raw_playlist/*.json") # 把嵌套的歌曲列表炸开,每首歌一行 songs = df.select( col("playlist.id").alias("playlist_id"), explode(col("playlist.tracks")).alias("track"), ) # 提取核心字段,时间戳统一转为 date 类型 clean = songs.select( col("track.id").cast("long").alias("song_id"), col("track.name").alias("song_name"), col("track.ar.name").alias("artist"), (col("track.dt") / 1000).cast("int").alias("duration_sec"), col("track.al.name").alias("album"), to_date(from_unixtime(col("track.publishTime") / 1000)).alias("publish_date"), ) # 按歌曲 ID 去重,保留第一次出现的那条 window = Window.partitionBy("song_id").orderBy("playlist_id") dedup = clean.withColumn("rn", row_number().over(window)).filter(col("rn") == 1).drop("rn") dedup = dedup.filter(col("duration_sec") > 0).filter(col("song_id").isNotNull()) dedup.write.mode("overwrite").partitionBy("publish_date").parquet("clean_songs")

这段代码的关键逻辑是explode与row_number的配合。explode负责把歌单里的tracks数组拆成多行,这一步做完才能用歌曲字段做聚合;row_number开窗去重则避免同一首歌被多个歌单重复计数。参数上要注意from_unixtime的输入必须是秒,所以publishTime先除以 1000;dt也要先除以 1000 再cast("int"),否则时长会变成一个很大的毫秒整数。

清洗前后的数据量对比如下:

指标清洗前清洗后
记录数10000+ 条曲目(含重复)8500 条唯一歌曲
duration_sec 为 0 的脏数据约 3%0
artist 为空约 1.5%已过滤或标记为 unknown
时间戳单位毫秒/秒混杂统一为秒

做完这一步,后面的聚合查询不用再关心字段单位、null 和重复值,分析层代码可以写得非常短。

3. Spark 集群搭建与内存调优:本地任务如何跑到集群上

3.1 spark 集群搭建:最小成本方案选 standalone 还是 YARN

毕设场景的数据量通常落在 GB 级以内,单机 Spark 完全能跑完,但答辩时“我搭了一个集群”和“我用单机跑了一下”是完全不同的两个故事。spark 集群搭建不是非要机房级的机器配置,最常见的做法是本机装一个 Spark 发行版,再用 Docker 起一个 standalone 集群作为演示环境。

# docker-compose.yml 片段 services: spark-master: image: bitnami/spark:3.5 ports: - "8080:8080" - "7077:7077" spark-worker: image: bitnami/spark:3.5 depends_on: - spark-master environment: - SPARK_MASTER_URL=spark://spark-master:7077

这个方案的意义不在于替代生产集群,而是让你在开发机上能复现分布式执行的调度行为。standalone 模式下spark-submit会把任务提交到 master 的7077端口,worker 节点从镜像里读取任务执行,驱动仍然在本地。如果你只有一台笔记本,需要注意给 Docker 分配至少 4GB 内存,否则 worker 进程会因内存不足反复重启。

实际做数据分析时,并不需要每次都在集群里跑。开发阶段的代码可以直接用local[*]模式运行,逻辑验证通过后再用--master spark://localhost:7077提交到集群,这样能把“逻辑错误”和“分布式配置问题”分开排查。

3.2 Spark on YARN 提交:是不是只需要一个 Spark 客户端

这是搜索热度很高的问题,答案比大多数人想的简单:只要提交节点能通过网络访问 YARN 的 ResourceManager,就只需要在这一个节点上安装 Spark。提交时 Spark 会把自己打好的 JAR 包和依赖一起上传到集群,由 YARN 在 NodeManager 上为 ApplicationMaster 和 Executor 分配容器。

./bin/spark-submit \ --master yarn \ --deploy-mode cluster \ --name ncm_analysis \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 3 \ --queue default \ ncm_analysis.py

这段命令里最值得理解的不是参数数量,而是deploy-mode。cluster模式下 Driver 运行在 YARN 容器里,提交命令可以直接退出;client模式下 Driver 跑在提交节点的 JVM 里,日志直接打在终端,调试时更直观。毕设演示用client模式更方便看日志,但需要保证提交节点和集群之间的网络稳定。executor-memory的设置要和 YARN 的yarn.nodemanager.resource.memory-mb匹配,给 Executor 4g 时容器实际占用的内存会略大于这个值,因为要加上 overhead,默认是0.1 * executor-memory。

数据在 YARN 模式下经 shuffle 写盘,spark.shuffle.service.enabled=true时中间结果由 NodeManager 持有,这样 Executor 释放后 shuffle 文件不会立刻丢失。值得注意的一个细节是,如果你在本地起过 standalone 集群又提交到 YARN,一定要把SPARK_HOME/conf/spark-defaults.conf里残留的spark.master配置注释掉,否则提交命令会优先读配置而忽略命令行参数。

3.3 Spark 内存模型与 OOM:90% 的原因出在三个地方

Spark 内存模型是“期末答辩前最容易暴露底气”的话题。当前版本的 Spark 把 Executor 内存分成 Reserved、User Memory、Execution 与 Storage 三个区域,其中 Execution 与 Storage 可以互相借用。spark.memory.fraction=0.6表示堆内总内存的 60% 分给这两块,spark.memory.storageFraction=0.5表示 Storage 初始占这块区域的 50%。

实际跑网易云音乐数据分析时,最容易触发 OOM 的场景有三个:一是groupBy做全量聚合时,某个热点 key 对应的数据量远大于其他 key,单个 Executor 要处理的数据超过 partition 限制;二是collect()全量拉取结果,把集群内存灌回 Driver;三是cache()之后又触发大量 shuffle,Storage 内存被借用后缓存被清空。应对办法是不用加内存,而是调整并行度或改变代码写法。

--conf spark.sql.shuffle.partitions=200 --conf spark.memory.fraction=0.7 --conf spark.executor.extraJavaOptions=-XX:+UseG1GC

spark.sql.shuffle.partitions默认是 200,但毕设数据量也就几万行,200 个分区意味着大量空任务,反而让每个任务的开销变大。数据量小的时候把分区数降到 50 以下,shuffle 效率和稳定性都会明显改善。G1GC 选项则是经验值,大堆内存下 CMS 的 Full GC 更容易造成 Executor 失联。

还有一个容易踩的坑是collect的误用。从 DataFrame 转成 Python 列表来遍历很方便,但 Driver 端内存上限通常只有 2g,一次拉了 20 万行数据就可能 OOM。改成用df.write.parquet落盘后再读取,或者用foreachPartition按分区处理,都不会把数据集中到单点。

4. 分析建模:网易云音乐数据在 Spark 里怎么算

4.1 常规聚合:歌手作品量与评论量分布

清洗后的clean_songs表已经可以做业务分析了。比较有代表性的聚合维度有三个:歌手维度、年份维度、歌曲时长分布。下面的代码统计每个歌手的作品数量,并按作品数降序排列:

artist_stats = dedup.groupBy("artist").agg( count("song_id").alias("song_cnt"), avg("duration_sec").alias("avg_duration"), ).orderBy(col("song_cnt").desc()) artist_stats.show(20)

groupBy之后 Spark 会自动触发一次全量 shuffle,把相同艺术家的记录分到同一个 Executor 上。这里要注意artist字段虽然清洗过,但“周杰伦”和“周杰伦 / 浪花兄弟”不是同一个值,做严格的分组统计时建议先用split(artist, " / ")[0]提取第一位歌手。如果做的是评论分析,则建议以歌曲 ID 为粒度先 join 评论表,再聚合到歌手,而不是直接对歌手分组,否则一首歌的多条评论会被重复计算。

聚合结果如果只需要 Top 20,用limit(20)在排序之后截取。这个顺序不要反:如果先limit(20)再排序,拿到的只是分区内的前 20 条,排序结果不准确。

4.2 时间维度:发布趋势与窗口函数

网易云音乐数据的publish_date字段在 ETL 时已经转成了标准日期,时间类分析可以直接对年月做粒度聚合:

monthly = dedup.withColumn("month", date_format(col("publish_date"), "yyyy-MM")) \ .groupBy("month").agg(count("song_id").alias("release_cnt")) \ .orderBy("month") trend = monthly.withColumn( "yoy", (col("release_cnt") - lag("release_cnt", 12).over(Window.orderBy("month"))) / lag("release_cnt", 12).over(Window.orderBy("month")) )

这里使用了lag窗口函数计算同比,本质是同一列在不同行的对比。窗口函数与groupBy的区别在于它不会把多行合并成一行,因此可以在保留整体聚合结果的同时增加新的计算列。需要注意窗口内部的orderBy与 DataFrame 的orderBy完全独立,只影响窗口内的排序。

做时间序列时还要注意分区边界问题。如果publish_date存在 NULL,date_format会直接把它变成null并分到null组,建议在聚合时用filter(col("publish_date").isNotNull())提前过滤,避免图表上出现一个异常的“null”柱子。

4.3 推荐近似:歌曲相似度计算的简单实现

网易云音乐数据分析如果只做排行榜,答辩时很容易被追问“分析完了然后呢”。一个低成本的增量是做一个基于歌曲特征的相似度计算,从歌曲名称、专辑名和歌手名中抽取关键词,把每首歌转成向量后计算余弦相似度:

from pyspark.ml.feature import Tokenizer, HashingTF, IDF, Normalizer from pyspark.ml.linalg import Vectors tokenizer = Tokenizer(inputCol="search_text", outputCol="words") words_df = tokenizer.transform(dedup.select( col("song_name"), col("artist"), concat_ws(" ", col("song_name"), col("artist")).alias("search_text") )) htf = HashingTF(inputCol="words", outputCol="rawFeatures", numFeatures=1000) idf = IDF(inputCol="rawFeatures", outputCol="features").fit(words_df) tfidf_df = idf.transform(htf.transform(words_df)) normed = Normalizer(inputCol="features", outputCol="norm_features").transform(tfidf_df)

这个流程是标准的 TF-IDF,HashingTF把词语哈希到固定维度的稀疏向量,IDF降低常见词的权重,最后用Normalizer把向量归一化,使得内积即余弦相似度。文本字段大小写、英文缩写和空白符号要先清洗,否则同样的词会被哈希成不同的特征。相似度计算可以缩小到“同歌手下的歌”或“同一年代的歌”来减少笛卡尔积规模,全量两两比较在数据量大时会轻易产生上亿条中间结果,这在毕设阶段的单机环境完全跑不动。

这套方案只是学术演示意义上的推荐,不能对标生产推荐系统,但作为毕业设计“从数据到应用”的收尾足够。关键是要在文档里说清楚特征来源和相似度的局限:没有用户行为数据参与,纯内容相似很难反映真实听感。

5. 可视化看板:从结果 Parquet 到可交互报表

5.1 PyECharts 生成静态 HTML

Spark 计算完的结果通常落成 Parquet 或者 CSV,可视化阶段不需要再让 Python 进程连着 Spark 跑,而是直接读取结果文件。PyECharts 是最顺手的一层封装,输出 HTML 不需要起 Web 服务,本地浏览器打开就能展示,答辩演示零额外依赖。

import pandas as pd from pyecharts.charts import Bar from pyecharts import options as opts df = pd.read_parquet("result/top_artists.parquet") bar = ( Bar() .add_xaxis(df["artist"].head(10).tolist()) .add_yaxis("作品数量", df["song_cnt"].head(10).tolist()) .set_global_opts(title_opts=opts.TitleOpts(title="歌手作品量 Top 10")) ) bar.render("top_artists.html")

这里用read_parquet读回结果,Bar().add_yaxis的yaxis数据必须是 Python list,所以取了head(10).tolist()。如果要做更复杂的图表联动,可以用 PyECharts 的Grid把趋势折线图和歌手柱状图放在同一个页面里,比生成多张独立图片更直观。需要避开的一个问题是把全量数据塞进图表,几万条柱子的 HTML 文件会超过 50MB,浏览器渲染直接卡死。常规做法是聚合到 Top 50 或按月份降采样后再可视化。

5.2 Streamlit 搭一个可筛选的交互看板

静态 HTML 的局限是没有筛选交互,而 Streamlit 只用少量代码就能把看板变成可操作页面。它的好处是不用写前端逻辑,st.selectbox和st.slider会直接映射成页面控件,选择不同值后重新过滤数据并更新图表。

import streamlit as st import pandas as pd df = pd.read_parquet("clean_songs") artists = df["artist"].unique().tolist() selected_artist = st.selectbox("选择歌手", artists) range_data = df[(df["artist"] == selected_artist)] st.line_chart(range_data.groupby("publish_date").size())

st.line_chart接收的是 Pandas DataFrame,一行代码就能输出趋势图;st.selectbox的默认值取artists的第一个元素,所以列表为空时要单独处理。Streamlit 的脚本是从上到下执行,数据量过大时每次交互都会重跑一次全部逻辑,因此最耗时的过滤步骤要加@st.cache_data装饰器做缓存。

交付物格式使用场景
Top 榜单 HTML静态页面答辩演示初版,零环境依赖
Streamlit 看板Web 页面现场演示筛选、下钻、趋势切换
图表截图 PNG图片插入论文和答辩 PPT,避免现场演示崩

Streamlit 看板跑起来只需要streamlit run app.py,但要注意它会默认占用 8501 端口,如果集群上的端口被防火墙拦住了,本地开发时直接用--server.address localhost限定监听地址即可。分析的结果文件建议放在单独的output/目录,可视化脚本只读不写,这样改图表样式时不会碰坏 Spark 算出来的数据。

6. 三条硬指标:判断 Spark 任务是否真的正确

6.1 Spark UI 的 DAG 与 Executor 界面看什么

Spark UI 不只是看进度条用的。提交任务后打开http://localhost:4040,先看 Executors 标签页里的Shuffle Read/Write总量——如果 shuffle write 达到 10 倍于源数据量,说明代码里有大量重复的宽依赖,可以检查是否多做了几次无意义的 groupBy 或 join。再看 SQL 标签页里的物理计划,注意有没有Exchange节点的数量异常偏多。正常情况下一次 groupBy 只有一个 Exchange,两个以上就要警惕笛卡尔积或复用了未持久化的中间结果。

6.2 数据质量校验:主键完整率与空值率

交付分析结果前,至少跑一次全量校验脚本。用 Spark SQL 算三个数字:song_id的重复率是否为 0、artist的空值占比是否低于 1%、publish_date是否都在合理时间范围内。这三个指标能挡住 ETL 阶段的大部分隐性错误。校验脚本单独放一个文件,和主分析代码分开,每次数据更新后跑一遍即可。

6.3 与关系型数据库结果对拍:最简单的验收技巧

拿同一份清洗后的数据导出成 CSV,在 MySQL 或 SQLite 里用 SQL 重新算一遍 Top 歌手的作品数,再和 Spark 的结果对比。两个独立系统算出来的结果如果不一致,按字段逐级排查;如果一致,说明聚合逻辑本身没有歧义。这个方法成本很低,却能非常有效地回应“分布式算出来的结果可信吗”这类答辩提问,远比“Spark 是成熟的框架所以结果没错”有说服力。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询