尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

基于Spark的网易云音乐数据分析:从数据采集到可视化看板全流程实践

基于Spark的网易云音乐数据分析:从数据采集到可视化看板全流程实践 简介这是一份基于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-AgentCookie 里至少要有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, cookieNMTIDxxx; MUSIC_Uxxx, sleep_sec1.5): url fhttps://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, headersheaders, timeout10) if resp.status_code ! 200: return None return resp.json() os.makedirs(raw_playlist, exist_okTrue) for pid in [3778678, 3779629, 2884035]: data fetch_playlist(pid) if data: with open(fraw_playlist/{pid}.json, w, encodingutf-8) as f: json.dump(data, f, ensure_asciiFalse) time.sleep(2)这段代码做了三件必要的事把playlist_id拼进接口路径、把Cookie写在请求头里、把每次请求的响应完整落盘。参数上最重要的不是超时时间而是sleep_sec——不同歌单的请求间隔至少要留出 1 秒以上连续快速请求非常容易触发临时封禁。这里刻意不做解析因为原始响应里还嵌套着推荐语、标签、创建者信息后续清洗时统一处理比抓取时逐个抽取更高效。落盘格式上建议一行为一个 JSON 对象而不是把整个歌单存成一个 JSON 数组。行式 JSON 能让 Spark 的read.json直接以multiLinefalse的方式读入未来如果有多个歌单增量抓取追加写入也不会破坏格式。评论数据同理把每一个评论作为一行单独落盘字段名保持和接口返回一致。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先除以 1000dt也要先除以 1000 再cast(int)否则时长会变成一个很大的毫秒整数。清洗前后的数据量对比如下指标清洗前清洗后记录数10000 条曲目含重复8500 条唯一歌曲duration_sec 为 0 的脏数据约 3%0artist 为空约 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_URLspark://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.enabledtrue时中间结果由 NodeManager 持有这样 Executor 释放后 shuffle 文件不会立刻丢失。值得注意的一个细节是如果你在本地起过 standalone 集群又提交到 YARN一定要把SPARK_HOME/conf/spark-defaults.conf里残留的spark.master配置注释掉否则提交命令会优先读配置而忽略命令行参数。3.3 Spark 内存模型与 OOM90% 的原因出在三个地方Spark 内存模型是“期末答辩前最容易暴露底气”的话题。当前版本的 Spark 把 Executor 内存分成 Reserved、User Memory、Execution 与 Storage 三个区域其中 Execution 与 Storage 可以互相借用。spark.memory.fraction0.6表示堆内总内存的 60% 分给这两块spark.memory.storageFraction0.5表示 Storage 初始占这块区域的 50%。实际跑网易云音乐数据分析时最容易触发 OOM 的场景有三个一是groupBy做全量聚合时某个热点 key 对应的数据量远大于其他 key单个 Executor 要处理的数据超过 partition 限制二是collect()全量拉取结果把集群内存灌回 Driver三是cache()之后又触发大量 shuffleStorage 内存被借用后缓存被清空。应对办法是不用加内存而是调整并行度或改变代码写法。--conf spark.sql.shuffle.partitions200 --conf spark.memory.fraction0.7 --conf spark.executor.extraJavaOptions-XX:UseG1GCspark.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存在 NULLdate_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(inputColsearch_text, outputColwords) words_df tokenizer.transform(dedup.select( col(song_name), col(artist), concat_ws( , col(song_name), col(artist)).alias(search_text) )) htf HashingTF(inputColwords, outputColrawFeatures, numFeatures1000) idf IDF(inputColrawFeatures, outputColfeatures).fit(words_df) tfidf_df idf.transform(htf.transform(words_df)) normed Normalizer(inputColfeatures, outputColnorm_features).transform(tfidf_df)这个流程是标准的 TF-IDFHashingTF把词语哈希到固定维度的稀疏向量IDF降低常见词的权重最后用Normalizer把向量归一化使得内积即余弦相似度。文本字段大小写、英文缩写和空白符号要先清洗否则同样的词会被哈希成不同的特征。相似度计算可以缩小到“同歌手下的歌”或“同一年代的歌”来减少笛卡尔积规模全量两两比较在数据量大时会轻易产生上亿条中间结果这在毕设阶段的单机环境完全跑不动。这套方案只是学术演示意义上的推荐不能对标生产推荐系统但作为毕业设计“从数据到应用”的收尾足够。关键是要在文档里说清楚特征来源和相似度的局限没有用户行为数据参与纯内容相似很难反映真实听感。5. 可视化看板从结果 Parquet 到可交互报表5.1 PyECharts 生成静态 HTMLSpark 计算完的结果通常落成 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_optsopts.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 是成熟的框架所以结果没错”有说服力。本文还有配套的精品资源点击获取
返回列表