ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

基于Spark的网易云音乐数据分析:从数据清洗到机器学习与图计算

基于Spark的网易云音乐数据分析:从数据清洗到机器学习与图计算 简介面向计算机专业毕业设计及课程设计场景的Spark网易云音乐数据分析实战包围绕图计算、机器学习歌曲分类、评论词云、评论时间段统计等核心任务展开适合需要完成真实项目、准备毕设答辩或提升Spark工程能力的开发者。整个压缩包共484个文件主体由Java与Scala源码构成后端分析逻辑配合JSP动态页面及HTML/CSS/JS前端展示同时包含XML配置、SQL初始化脚本、CSV数据集、PNG/JPG图表截图等支撑材料整体大小约11.63MB目录按功能模块组织检索与复用都很方便。目前已有40人学习下载可作为毕业设计选题参考或课程项目进阶资料。资源内含可运行的项目源码、环境配置文件、预处理数据与可视化界面素材并附有详细说明文档覆盖从数据清洗、图计算、分类模型训练到评论词云与时间段统计的完整链路能让读者直观看到Spark在真实音乐数据场景中的落地方式尤其适合想深入理解分类预测与图计算结合的实践者。1. Spark在网易云音乐数据中的应用分析这道毕设题到底在考什么“Spark在网易云音乐数据中的应用分析”这个题目迷惑性在于四个分析点分开看都容易合起来却很容易翻车。很多人一开始把评论词云当作切入口跑出第一张图时很有成就感等到要做机器学习预测歌曲分类和图计算时才发现前面的数据根本没法复用——标签没设计、字段对不上、Spark任务跑一半内存就爆了。这道题真正考查的不是会不会调某个库而是能不能把“评论、歌曲、歌手”三类数据串成一条可复现的链路从接口取数、Spark清洗、按小时聚合到训练分类模型、构建歌曲关联图。它适合想做大数据毕设但又不想停留在WordCount的人也适合需要一份能讲清楚分析逻辑作品的求职者。下面这套落地顺序是我按实际交付这类项目的方式整理的照着走能把两周弯路省掉。2. 数据准备与Spark环境先把评论、歌曲和歌手字段整理成能算的表先泼盆冷水这个题目的工作量一半以上在取数和清洗。官方没有公开数据集常见做法是用歌曲ID循环请求评论接口再按歌单接口补歌曲元信息。规模上不需要全站2000首热门歌曲、一份歌单表就足够支撑四个分析模块。数据拿回来之后先别急着分析第一步是把三张表的结构定下来。2.1 数据来源与核心字段评论、歌曲、歌手怎么对齐网易云音乐的评论接口路径是R_SO_4_歌曲ID同一个接口里返回hotComments和comments两个数组。数组里每个元素都包含commentId、content、time、likedCount以及一个嵌套的user对象里面才是userId。歌曲元数据从歌单接口拿包含歌曲名、歌手名、专辑、时长和热度值。按“分析口径先于代码”的原则先建三张表表名关键字段模块用途commentscommentId、content、time毫秒、likedCount、userId词云、时间段、分类特征songssongId、name、artists、popularity、playlistId分类标签、图节点artistsartistId、name、followCount关联特征、图节点需要注意两个口径问题。第一time是毫秒级Unix时间戳不是字符串后面所有时间分析都必须先除以1000再转格式。第二user是嵌套结构在Spark里读取时要写成user.userId直接写userId会拿到一列null。这两个坑几乎每个做这个题的人都会踩后面避坑章还会展开。2.2 Spark环境与三个必调参数local模式同样需要调内存毕设阶段不需要一上来就搭Spark集群本地local[4]模式足够跑2000首歌的评论数据。集群搭建经验可以写在简历里但跑通分析逻辑最要紧的是先让单机不崩溃。下面这段SparkSession配置我会在项目里反复用from pyspark.sql import SparkSession spark (SparkSession.builder .appName(netease_music_analysis) .master(local[4]) .config(spark.driver.memory, 4g) .config(spark.serializer, org.apache.spark.serializer.KryoSerializer) .config(spark.sql.shuffle.partitions, 16) .getOrCreate()) spark.sparkContext.setCheckpointDir(file:///tmp/spark_checkpoint)三个必调参数的解释spark.driver.memory设成4g因为词云和分词阶段要把部分数据collect回驱动端内存太小直接OOMspark.sql.shuffle.partitions默认是200本地跑纯属浪费改成16能让小任务快很多Kryo序列化对后面图计算的边对象尤其友好Java默认序列化在节点数据多时很容易把内存撑爆。setCheckpointDir这行是给图计算做checkpoint用的不设的话connectedComponents这种迭代型算法会在长任务中反复回算同一批数据跑得慢不说还可能直接丢executor。2.3 用Spark读JSON评论文件explode和selectExpr的正确顺序评论接口返回的是一个嵌套JSON数组直接spark.read.json(data/comments/*.json)读进来之后顶层只有hotComments、comments、total这几个字段评论内容全在数组里。必须先explode把数组铺平再选字段raw spark.read.json(data/comments/*.json) raw.printSchema() comments (raw .select(F.explode(comments).alias(c)) .select( F.col(c.commentId).alias(cid), F.col(c.content).alias(text), F.col(c.time).alias(ts), F.col(c.likedCount).alias(likes), F.col(c.user.userId).alias(uid) ) .filter(text is not null and length(text) 0) .dropDuplicates([cid])) comments.write.mode(overwrite).parquet(data/clean_comments.parquet)explode的返回值是一个结构体所以第二次select里所有字段都要带c.前缀。过滤条件放在select之后而不是之前是为了在字段选完前保留原始行的可追溯性。最终结果落成Parquet后续每个模块都从这份干净表读数据而不是每次重新解析原始JSON。落盘这一步很重要它让你的中间结果可复用也方便答辩时被问到“数据从哪来”时快速回溯。3. 评论时间段与词云先用SQL和分词产出两张能答辩的图数据备好之后先做两个出图快的评论时间段和评论词云。它们不涉及训练适合用来验证清洗结果对不对。如果这两张图的结论和常识明显违背比如深夜0点评论量暴增那大概率是时间戳没除1000先把数据修好再继续往下走。3.1 把毫秒时间戳切成小时桶24小时活跃曲线的Spark SQL写法时间段分析的常见做法是直接按小时分组统计。注意ts是毫秒需要先除以1000再交给Spark的from_unixtime转成时间字符串最后用hour函数提取小时。不要图省事在Python里循环处理几十万行SparkgroupBy一次性算完再出图from pyspark.sql import functions as F hour_stat (comments .withColumn(datetime, F.to_timestamp(F.from_unixtime(F.col(ts) / 1000))) .withColumn(hour, F.hour(F.col(datetime))) .groupBy(hour) .agg( F.count(*).alias(comment_cnt), F.avg(likes).alias(avg_likes) ) .orderBy(hour)) hour_stat.show(24)comment_cnt是每个小时段的评论总量avg_likes是平均点赞数两者放在一起能看出“活跃高峰”和“高赞高峰”不一定重合——比如深夜评论少但单条赞多。如果想让结论更有层次再加一列星期几把工作日和周末拆开看weekday_hour (comments .withColumn(dow, F.dayofweek(F.to_timestamp(F.from_unixtime(F.col(ts) / 1000)))) .withColumn(hour, F.hour(F.to_timestamp(F.from_unixtime(F.col(ts) / 1000)))) .groupBy(dow, hour) .count() .orderBy(dow, hour))这段代码里to_timestamp被调用了两次实际项目里应该先用withColumn生成datetime列再复用。dayofweek返回1到71是周日分析时最好先做映射再聚合否则图画出来坐标轴全是数字。3.2 评论内容的中文分词不直接在Spark里调分词器的取舍词云的原料是text列。常见做法是先用jieba分词再统计词频。二十万条以下的中文评论直接把text列collect回驱动端用Python处理比在Spark里调JiebaSegmenter之类的第三方包省事得多——后者要额外编译jar包环境坑一个接一个。我一般这样处理text_list comments.select(text).toPandas()[text].dropna().tolist() import jieba from collections import Counter stopwords set(open(stopwords.txt, encodingutf-8).read().splitlines()) counter Counter() for line in text_list: counter.update([ w for w in jieba.cut(line) if len(w) 1 and w not in stopwords and not w.isspace() ]) top_words counter.most_common(200)stopwords.txt自己维护一份常用中文停用词表包含“我们”“你们”“就是”“真的”这类无意义词。过滤条件里len(w) 1很重要单字词在歌曲评论里大多是语气词混进词云会显得很脏。most_common(200)限定词云只展示前200个词词频差距太悬殊时可以取对数再做词云避免高频词一家独大。3.3 词云渲染WordCloud的字体与频率两个必调参数词频表出来后交给wordcloud库画图。这里有两个必调参数font_path和generate_from_frequencies。不指定中文字体画出来全是方块不传词频表而是直接传原始文本WordCloud会重新分词一次结果和你前面算好的词频不一致。from wordcloud import WordCloud import matplotlib.pyplot as plt wc WordCloud( font_path/usr/share/fonts/truetype/wqy/wqy-microhei.ttc, width1200, height800, background_colorwhite, max_words200 ) wc.generate_from_frequencies(dict(top_words)) plt.imshow(wc, interpolationbilinear) plt.axis(off) plt.savefig(comment_wordcloud.png, dpi150)font_path在不同系统位置不一样Windows常见的是C:/Windows/Fonts/simhei.ttfmacOS可以填/System/Library/Fonts/PingFang.ttcLinux服务器如果没有中文字体先apt install fonts-wqy-microhei。interpolationbilinear只是让图更平滑。导出图片后建议同时把top_words存成CSV答辩PPT里可以直接放生成的图附录里放词频表。4. 机器学习预测歌曲分类与图计算两个吃算力模块的完整落地时段和词云是“看得到结果”的分析到了歌曲分类和图计算才开始真正考验Spark的分布式能力。这两个模块的共同点是都需要构造训练样本、都需要控制数据规模。先做分类再做图计算因为分类产出的标签可以直接给图模块当节点属性用。4.1 歌曲分类的标签与特征不要自己手打5000首歌做有监督分类先解决标签从哪来。常见做法是借歌单名做弱标签把包含“欧美流行”“POP”的歌单里的歌标成pop包含“ACG”“二次元”的标成anime包含“纯音乐”“钢琴曲”的标成instrumental再用Spark的when表达式批量打标而不是手工一张张标。特征侧我常用五个维度歌曲热度popularity、评论总量、平均点赞、艺人关注数、发行年份。这些字段分散在comments和songs两张表先join再组装song_features (songs .join(comments.groupBy(song_id) .agg(F.count(*).alias(comment_total), F.avg(likes).alias(avg_likes)), onsong_id) .withColumn(category, F.when(F.col(playlist_name).rlike(流行|Pop|pop), pop) .when(F.col(playlist_name).rlike(纯音乐|钢琴|轻音乐), instrumental) .otherwise(other))) from pyspark.ml.feature import StringIndexer, VectorAssembler indexer StringIndexer(inputColcategory, outputCollabel) assembler VectorAssembler( inputCols[popularity, comment_total, avg_likes, artist_follows, release_year], outputColfeatures )StringIndexer会把pop、instrumental、other转成0、1、2的数字标签注意它按频次排序不是按字母。绝大多数评论集中在少数头部歌曲所以“other”类占比会很高后面的评估必须关注类别不平衡。VectorAssembler只能接受数值列如果artist_follows是字符串要先cast(double)否则运行时报类型错误。4.2 用随机森林训练歌曲分类模型选择与参数边界随机森林是这类表格分类问题的稳妥选择对特征尺度不敏感、抗过拟合、训练速度快而且能输出特征重要性。决策树容易过拟合逻辑回归对非线性关系欠拟合所以我在这个题目上一般直接上随机森林。核心代码如下train, test sf.randomSplit([0.8, 0.2], seed42) from pyspark.ml.classification import RandomForestClassifier from pyspark.ml.evaluation import MulticlassClassificationEvaluator rf RandomForestClassifier( featuresColfeatures, labelCollabel, numTrees50, maxDepth6, seed42 ) model rf.fit(train) pred model.transform(test) evaluator MulticlassClassificationEvaluator( labelCollabel, predictionColprediction, metricNamef1 ) print(F1 Score:, evaluator.evaluate(pred))numTrees设50是本地运行的平衡点再大收益有限但训练时间明显变长。maxDepth6控制树深度太浅欠拟合太深容易把训练集的噪声学进去。评估指标我不用accuracy而用f1因为如果某一类占比超过70%accuracy可能虚高。输出结果里的prediction列可以直接写回Parquet后续图计算的节点属性正好用上。4.3 图计算用“同一用户评论多首歌曲”构建关联图图计算的第一步是定义节点和边。节点取歌曲边表示两首歌被同一个人评论过——这是一种典型的共现关系不需要额外爬关注关系。对毕设体量来说这个定义足够支撑社区发现。先算共现矩阵# 每首歌被哪些用户评论过 user_song (comments .select(uid, song_id) .distinct()) # 同一用户评论过的歌曲两两成边只保留出现3次以上的强关联 cooc (user_song.alias(a) .join(user_song.alias(b), F.col(a.uid) F.col(b.uid)) .filter(F.col(a.song_id) F.col(b.song_id)) .groupBy(a.song_id, b.song_id) .count() .filter(count 3))这里用而不是!做过滤是为了让每对歌曲只保留一条边否则 (A,B) 和 (B,A) 会各算一遍边数翻倍且没有意义。阈值count 3是过滤弱关联的关键不设这个阈值热门歌曲之间会产生完全图图计算跑得慢且社区结构被噪声淹没。构图和社区发现用GraphFrame完成from graphframes import GraphFrame vertices songs.selectExpr(song_id as id, name, category) edges cooc.selectExpr(a.song_id as src, b.song_id as dst, count as weight) g GraphFrame(vertices, edges) cc g.connectedComponents() cc.orderBy(component).show(50)connectedComponents会把所有互相可达的歌曲归到同一个社区输出里多一列component。判断结果好坏有一个实用的行业习惯看同一个社区里的歌曲是不是同一风格。如果pop和instrumental被分到同一个社区说明边的定义或阈值有问题。另外图计算的规模最好控制在5000节点以内超过这个量级本机跑connectedComponents的时间会成倍增长。5. 毕设避坑指南六个让Spark分析卡壳的实际问题与排查每个做这个题目的人都会在相同的位置摔跤。下面这六条是我自己踩过、也看别人反复踩的坑按“现象→原因→解决”写清楚遇到同类问题时可以直接对照。5.1 评论字段全为null忘记把comments数组铺平现象spark.read.json读入后select(content)出来全是null但time字段却有值。原因评论接口返回的JSON里content嵌在comments数组内部顶层根本没有这个列。Spark读取嵌套数组时不会自动展开直接按顶层字段选自然拿到一堆null。解决先raw.printSchema()看清楚层级再用explode(comments)把数组拆成行然后从展开后的结构体里取字段。注意hotComments和comments是两个不同数组结构相同但含义不同分析时要么只用一个、要么用union把两者合并。5.2 词云输出全是方块字体缺失而不是数据问题现象词云生成成功程序没报错但图片里全是小方块。原因WordCloud默认字体不支持中文字形在Linux服务器上尤其常见因为系统没有安装任何中文字体。这是字体问题不是数据问题。解决WordCloud里显式传入font_path指向系统中文字体文件。运行前先确认路径存在用fc-list :langzh查看Linux系统已安装的中文字体。如果服务器实在没有字体直接安装fonts-wqy-microhei或者干脆在本地生成图片后再传到服务器。5.3 VectorAssembler报“Column does not exist”列名与空值的双重陷阱现象VectorAssembler.fit()报错说某个输入列不存在但你明明在DataFrame里见过这个列名。原因两个原因叠加。第一join之后列名重复Spark自动给其中一个加了后缀原来的列名被覆盖第二有些列全是nulltoPandas()时被丢弃导致后续select找不到列。解决join之前用withColumnRenamed把重复列名改掉VectorAssembler之前先对特征列做fillna(0)保证每个特征列都是完整数值。检查列是否完整用df.select(col).dropna().count()和总行数对比。5.4 图计算Executor OOMspark内存不够和checkpoint缺失现象connectedComponents跑到一半控制台报 Executor Lost / OOM重试还是崩。原因边数据量比预想大一个量级。同一用户评论N首歌会产生N*(N-1)/2对边一个深度用户就能贡献成百上千对。加上没设checkpoint目录图计算的长迭代阶段每次随机失败都要从头重建数据内存越来越紧张。解决构图前先按歌曲热度截断只保留TOP 500歌曲再算共现过滤条件从count 3往上加观察边数量级SparkSession里用setCheckpointDir设置检查点目录配合spark.graphx.pregel.checkpointInterval控制迭代检查频率。5.5 机器学习调参一跑几小时网格搜索和并行度的真实权衡现象用TrainValidationSplit做参数搜索配了5x5的参数网格跑了一个小时没出结果。原因TrainValidationSplit默认会把每个参数组合跑多折本地local[4]模式下并行度有限网格一大就是指数级耗时。初次调参时最常见的问题就是网格过密。解决第一轮用3x3小网格numTrees和maxDepth各取三个值手动切分train/test跑完记录F1第二轮在最优参数附近再加密。另一个行之有效的方式是关闭TrainValidationSplit的collectSubModels参数避免在驱动端堆积大量子模型既省内存又提速。5.6 时间统计普遍偏移12小时session时区没对齐现象24小时活跃曲线显示凌晨0点到4点的评论量最高曲线整体向后偏移。原因Spark的from_unixtime和hour依赖spark.sql.session.timeZone设置。我自己就遇到过JVM时区默认UTC北京时间凌晨的评论被转成了UTC时间峰值整体错位。这个坑特别隐蔽因为程序不报任何错误。解决SparkSession里显式设置spark.conf.set(spark.sql.session.timeZone, Asia/Shanghai)在读取时间戳之前完成配置。更稳妥的做法是统一用to_utc_timestamp把时间戳转成UTC再按小时聚合最后在结果解释时加回来8小时。6. 验证你的分析结果混淆矩阵、时段曲线与图社区的交叉印证模型和图表都跑出来之后答辩前留一天做结果验证。分类模块先看整体F1再看混淆矩阵哪两个类别容易互相混淆如果pop频繁被分到instrumental多半是特征里混入了评论量这类和风格无关的字段回4.1检查特征构造。时段曲线可以对照网易云音乐的运营规律工作日晚上8点到10点评论量最高周末的白天活跃期更早。如果曲线和这个常识偏差太远优先检查时区而不是去调整统计口径。图计算的社区结果则用“业务可解释性”来验证——同一个component里的歌曲风格是否集中能否在社区里找到周杰伦和林俊杰这种歌迷高度重叠的歌手组合。一个值得做的进阶是把图计算结果反哺给分类模型把每首歌所属的component编号当成新特征加入VectorAssembler再训练一次。社区信息携带了“同类听众”的群体信号往往能把F1再提升几个点。另一个进阶方向是词云按情感加权——评论点赞数归一化后作为词的权重再做词云这样被高赞评论使用的词会放大比单纯按词频更能反映用户真实情绪。我每次跑完这类项目都会把中间结果统一落成Parquet存一份字段清单放在项目根目录。答辩被追问“这个数据怎么来的”时能当场从原始JSON回溯到某首歌、某条评论、某个特征列而不是靠记忆回答。这个习惯帮我避免过好几次“明明算了却说不清”的尴尬。希望帮到你。本文还有配套的精品资源点击获取
返回列表