十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

基于Spark的网易云音乐离线数据分析全链路实践

基于Spark的网易云音乐离线数据分析全链路实践 简介一份基于Spark的网易云音乐数据分析毕业设计项目面向大数据相关专业学生与入门开发者解决从海量用户行为数据采集、清洗、分析到可视化展示的完整流程问题。项目覆盖Spark核心API、Spark SQL与MLlib、Hadoop分布式存储、MySQL结果落地、Azkaban任务调度以及Nest/ECharts可视化等环节同时涉及推荐系统建模思路适合用于课程设计、毕业设计或简历项目对照。压缩包共403个文件以java、scala源码为主体js、html/jsp构建前端展示xml/properties承担配置管理png/jpg保存图表结果sql/csv/dat存放数据与脚本整体约9.31MB。已有485人学习下载。内容包括可运行的Spark分析代码、可视化网页模板、Flume/Hadoop配置示例、Azkaban工作流描述、MySQL建表SQL及readme说明文档目录结构清晰便于按数据层、计算层、展示层快速定位与二次开发。1. 网易云音乐日志的全链路分析毕业设计里Spark到底在算什么打开这份「毕业设计基于Spark网易云音乐数据分析.zip」最先看到的不是一堆Scala源码而是fontawesome-webfont.woff2、amazeui.min.css、bootstrap.min.css这些前端静态资源外加flume-hdfs-ng.conf、log4j-es.conf、Azkaban_1.png、MySQL_1.png这类配置和架构图文件。这说明它不是一个“调一下Spark API出个结果”的Demo而是一条从日志采集、分布式存储、离线批处理到MySQL落库、Echarts可视化的完整数据管道。网易云音乐每天会产生大量用户行为日志比如播放、收藏、搜索、评论数据量一旦达到亿级单机SQL和Excel就扛不住了。这个毕业设计的核心价值是演示如何用Flume接日志、用HDFS做原始存储、用Spark SQL做清洗和指标计算再用Azkaban把整条链路定时串起来。适合正在做大数据方向毕设、想搞清楚“Spark在真实项目里到底算什么”的同学也适合工作后想快速回顾离线数仓全流程的工程师。它解决的最关键问题不是某个算法多难而是数据从哪来、怎么算、算完放哪、怎么展示。2. Flume采集与HDFS落地让Spark拿到干净的原始日志2.1 为什么选用Flume而不是Kafka项目中出现了flume-hdfs-ng.conf这个文件名暴露了技术选型Flume 1.xng版本负责日志采集sink目标是HDFS。很多人在设计这类系统时第一反应是上Kafka但Kafka的引入意味着还要维护Zookeeper集群、控制consumer offset、处理消息堆积监控对毕业设计这种规模的集群来说运维成本偏高。Flume的agent本身就是Java进程配置一个spooldir或taildir source就能监控日志目录sink直接写HDFS链路短、排查直观。常见做法是单节点部署一个Flume agentsource监控网易云音乐埋点日志的输出目录channel用memorysink到HDFS按天分目录存储。这套设计的前提是离线T1分析而不是实时推荐。如果是实时场景比如用户听歌后立刻更新推荐列表那确实应该用Kafka加Spark Structured Streaming。但毕业设计里的指标比如日活、热门歌曲TopN、24小时播放趋势延迟一天完全可接受。Flume落地HDFS的另一个好处是原始日志不丢后续重跑任务时可以从HDFS重新读取不需要回溯业务系统。2.2 flume-hdfs-ng.conf核心配置解读项目里的flume-hdfs-ng.conf应该是整个采集层的核心配置内容一般是agent三要素source、channel、sink。我基于这个场景还原了一份典型配置# flume-hdfs-ng.conf agent.sources songSrc agent.channels memChannel agent.sinks hdfsSink # 监控日志落地目录处理完改后缀避免重复读取 agent.sources.songSrc.type spooldir agent.sources.songSrc.spoolDir /data/music_logs agent.sources.songSrc.fileSuffix .DONE agent.sources.songSrc.deletePolicy never agent.sources.songSrc.ignorePattern ^.*\.DONE$ # memory channel容量根据单日日志峰值估算 agent.channels.memChannel.type memory agent.channels.memChannel.capacity 10000 agent.channels.memChannel.transactionCapacity 1000 # sink 写 HDFS按天分目录 agent.sinks.hdfsSink.type hdfs agent.sinks.hdfsSink.hdfs.path hdfs://namenode:8020/music/raw/%Y%m%d agent.sinks.hdfsSink.hdfs.fileType DataStream agent.sinks.hdfsSink.hdfs.writeFormat Text agent.sinks.hdfsSink.hdfs.rollInterval 600 agent.sinks.hdfsSink.hdfs.rollSize 134217728 agent.sinks.hdfsSink.hdfs.rollCount 0 agent.sinks.hdfsSink.hdfs.filePrefix music agent.sinks.hdfsSink.hdfs.fileSuffix .log配置里有几个参数要特别说明。spooldir适合目录内文件不再变动的场景Flume会为每个文件维护一个.flumespool元数据断点续传靠它实现fileSuffix .DONE是防止同一个文件被反复消费的常用手段处理完成后源文件被改名但不会被删除这比deletePolicy immediate安全出问题时还能找回原始日志。memory channel的capacity是channel中最多能放多少eventtransactionCapacity是每个事务最多取多少event前者必须大于后者否则启动时会直接报错。HDFS sink里的rollInterval 600表示每600秒强制落盘一个文件rollSize 134217728128MB是文件大小滚动阈值两者谁先触发都行。这里要注意如果rollCount 0就表示不按event条数滚动避免小文件刷爆NameNode。HDFS路径里用了%Y%m%d这个时间取自event header默认是Flume服务器本地时间。按天分目录的好处很直接Spark读取时可以精确到某一天的目录作为输入做分区裁剪不需要全表扫描。这里有一个建议如果后续要做小时级分析可以把path改成%Y%m%d/%H但这样会产生更多小文件Spark读取时反而变慢所以T1场景按天够用了。2.3 整条链路的组件职责边界从这份毕业设计涉及的文件来倒推完整链路各层职责大致是这样的层级组件职责关键产物采集层Flume监控埋点日志目录写入HDFS/music/raw/yyyyMMdd/xxx.log存储层HDFS保存原始日志和中间结果按天分区的原始数据计算层Spark清洗、去重、聚合、TopN指标结果DataFrame调度层Azkaban定时触发Spark作业和导出作业工作流DAG结果层MySQL存储计算结果供前端查询指标表、榜单表展示层Echarts/AmazeUI读取MySQL数据渲染图表折线图、柱状图、词云这个表格对应了压缩包里那些图片和配置文件背后的设计意图。Hadoop_1.png说明集群环境是Hadoop 2.x加Spark on YARN部署Azkaban_1.png说明不是手动spark-submit而是通过Azkaban调度MySQL_1.png说明结果数据最终是结构化存储。这套组合是2018到2022年间大数据毕业设计里最主流的架构现在看依然适合教学演示因为它把离线数仓的每个环节都覆盖到了但又没有引入Kafka、Iceberg这类对毕设来说过重的组件。3. Spark SQL清洗与指标计算去重、TopN与时序统计的SQL写法3.1 原始日志的典型字段结构与埋点格式网易云音乐的日志通常包含用户ID、歌曲ID、行为类型、时间戳、设备信息等字段。下面是一份常见的埋点日志schema实际项目中可能还有更多维度但对毕业设计来说这几个字段已经能支撑80%的指标计算。字段名类型示例值说明user_idString8374291匿名用户可能为空或为设备IDsong_idString27438123歌曲唯一标识actionStringplay / collect / download行为类型play为主tsLong1715846400000毫秒级Unix时间戳device_typeStringAndroid / iOS / Web设备端dtString20240520分区字段跟目录对应日志进入HDFS时往往是原始文本字段用制表符或逗号分隔甚至混入一些Nginx上报的异常行。Spark SQL的DataFrame API负责把这些文本解析成结构化表。第一步是read加载后按分隔符split再用to_timestamp把毫秒时间戳转成可读时间格式。这里有一个常见的坑不要用__HIVE_DEFAULT_PARTITION__这种默认分区名去join得先把dt字符串校验一遍否则脏数据会把整张表带偏。3.2 清洗逻辑去重、过滤、类型转换用PySpark写清洗逻辑最直观也方便在Jupyter里跑通后再改成Scala封装进JAR。核心代码如下from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_timestamp, from_unixtime, split spark SparkSession.builder \ .appName(music-etl) \ .enableHiveSupport() \ .getOrCreate() # 按天读取HDFS原始日志目录 raw_df spark.read.text(/music/raw/20240520) # 切分字段并转类型 parsed_df raw_df.withColumn(user_id, split(col(value), \t).getItem(0)) \ .withColumn(song_id, split(col(value), \t).getItem(1)) \ .withColumn(action, split(col(value), \t).getItem(2)) \ .withColumn(ts, split(col(value), \t).getItem(3).cast(long)) \ .withColumn(device_type, split(col(value), \t).getItem(4)) \ .withColumn(event_time, from_unixtime(col(ts) / 1000, yyyy-MM-dd HH:mm:ss)) # 过滤user_id为空、ts异常、action不在白名单 clean_df parsed_df.filter( col(user_id).isNotNull() (col(user_id) ! ) col(ts).isNotNull() col(action).isin(play, collect, download, search) ) # 去重同一个人同一秒对同一首歌的重复上报只保留一条 dedup_df clean_df.dropDuplicates([user_id, song_id, action, ts]) dedup_df.createOrReplaceTempView(music_event)这里的split(col(value), \t)是按Tab切分如果日志实际是逗号分隔就需要替换。from_unixtime(col(ts) / 1000)把毫秒时间戳转成秒再格式化注意如果不除以1000格式化出来的时间会是1970年附近。dropDuplicates的粒度是用户、歌曲、行为、时间戳四元组不是全字段去重因为device_type等字段可能上报不一致全字段去重会留下重复统计的隐患。清洗后创建临时视图music_event后面所有指标都基于它计算。这里的createOrReplaceTempView是Session级别的临时表作业结束就释放不会污染Hive元数据。如果需要跨SparkSession复用结果可以把清洗后的数据写成Parquet格式到HDFS的中间目录比如/music/cleaned/dt20240520这也是离线数仓里“清洗层”的标准做法。3.3 三个核心指标TopN榜单、24小时趋势、用户活跃度指标计算是Spark SQL最出彩的部分。直接上SQL-- 1. 播放量Top20歌曲榜单 SELECT song_id, COUNT(*) AS play_cnt FROM music_event WHERE action play GROUP BY song_id ORDER BY play_cnt DESC LIMIT 20; -- 2. 24小时播放量趋势 SELECT HOUR(event_time) AS hour, COUNT(*) AS play_cnt FROM music_event WHERE action play GROUP BY HOUR(event_time) ORDER BY hour; -- 3. 用户播放行为Top10用窗口函数替代全局排序 SELECT user_id, song_id, play_cnt FROM ( SELECT user_id, song_id, COUNT(*) AS play_cnt, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY COUNT(*) DESC) AS rn FROM music_event WHERE action play GROUP BY user_id, song_id ) t WHERE rn 10;HOUR(event_time)是从Spark 2.0开始内置的时间函数返回0到23的整数用于小时粒度聚合不需要额外UDF。第三个查询里的ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY COUNT(*) DESC)是标准SQL窗口函数其执行效率取决于group by之后的数据分布如果热门用户听歌量巨大这里会触发shuffle。实际执行时Spark会把整个查询翻译成物理计划窗口函数的排序操作默认用的全局排序器对内存压力较大跑不过去时优先调spark.sql.shuffle.partitions而不是盲目加executor内存。这三个指标基本就是毕业设计可视化页面的主图柱状图显示热门歌曲、折线图显示24小时趋势、表格或嵌套饼图显示用户偏好。项目中的nest_1.png、nest_2.png应该就是用嵌套饼图展示不同维度的占比情况。3.4 Spark运行参数如何影响这批SQL同样一段Spark SQL在不同参数配置下运行时长可能差好几倍。下面是这组指标最常用的一组spark-submit参数spark-submit \ --master yarn \ --deploy-mode client \ --executor-memory 4g \ --num-executors 4 \ --executor-cores 2 \ --conf spark.sql.shuffle.partitions200 \ --conf spark.sql.autoBroadcastJoinThreshold10485760 \ --class com.music.MusicAnalysis \ music-analysis.jarspark.sql.shuffle.partitions决定了group by或join时产生的reduce任务数默认200。如果输入数据只有几百MB200个分区意味着每个分区只处理几MB数据空转开销很大如果数据有几十GB200个分区又不够单个任务处理太久。常见做法是先看Spark UI里Shuffle Read和Shuffle Write的字节数再估算一个合适的分区数。spark.sql.autoBroadcastJoinThreshold表示小于10MB的表自动广播到每个executor避免shuffle join这个毕业设计里如果需要把歌曲维度表join进来补全歌名和歌手这个参数就能派上用场。4. 结果落库MySQL与Echarts展示Azkaban把整条链路自动跑起来4.1 DataFrame写出MySQL的两种模式与坑计算结果不能一直躺在Spark里最终要给前端查询。项目里出现了MySQL_1.png说明结果数据是落到MySQL的。写出代码用DataFrame的jdbc接口即可result_df.write \ .mode(overwrite) \ .format(jdbc) \ .option(url, jdbc:mysql://localhost:3306/music_analysis?useUnicodetruecharacterEncodingutf8useSSLfalse) \ .option(dbtable, top_song_rank) \ .option(user, root) \ .option(password, 123456) \ .option(batchsize, 1000) \ .option(truncate, true) \ .save()这里有两个坑值得展开。第一个是mode(overwrite)配合truncatetrueSpark默认overwrite是先drop表再重建如果表结构变了会出问题而truncate只是清空数据保留表结构更安全。第二个是URL里必须带useUnicodetruecharacterEncodingutf8否则中文歌名写入MySQL会变成问号。batchsize1000表示每批写入1000条太大会导致MySQL连接超时太小则写入慢。写完后推荐在MySQL里对song_id和dt建立联合索引因为前端查询基本都带WHERE dt ?条件没有索引的话TopN榜单查询在数据量大了之后会明显变慢。对于每日榜单这种结果表用overwrite是合理的因为每天只保留当天结果对于用户行为明细表则应该用append并加上dt字段做分区标识便于回溯。两种模式混用时要小心同一张表不要一会儿append一会儿overwrite否则容易把历史数据搞丢。4.2 Echarts可视化后端JSON怎么交给前端渲染项目里的bootstrap.min.css、amazeui.min.css、font-awesome.css这些资源说明前端采用了AmazeUI响应式框架加Echarts图表库。Echarts从MySQL拿数通常不是直连而是后端接口返回JSON再渲染。常见做法是后端写一个Servlet或Spring Boot接口查询top_song_rank表转成以下JSON结构{ code: 0, data: { songs: [晴天, 稻香, 七里香], playCounts: [98234, 87321, 80234] } }前端拿到数据后用setOption填充// echarts 柱状图每日播放Top10歌曲 $.getJSON(/api/top_song, { dt: 20240520 }, function (res) { if (res.code ! 0) return; var chart echarts.init(document.getElementById(topSongChart)); chart.setOption({ tooltip: { trigger: axis }, xAxis: { type: category, data: res.data.songs }, yAxis: { type: value }, series: [{ type: bar, data: res.data.playCounts, itemStyle: { color: #e74c3c } }] }); });这段代码里的xAxis.data和series.data分别对应榜单名称和播放量核心逻辑是数据从MySQL到JSON再到Echarts实例的映射。tooltip.trigger axis表示鼠标悬停在坐标轴附近时展示提示框柱状图适合这种模式。如果要做24小时趋势折线图只要把接口换成按小时聚合的数据源然后type: line即可。AmazeUI的作用是页面的栅格布局和响应式适配让图表在PC和手机上都能看这在毕设答辩时用平板展示是个加分项。4.3 Azkaban调度把ETL、指标计算、导出串成DAG项目里有Azkaban_1.png说明不是手动执行spark-submit而是用Azkaban定时触发。Azkaban提交的是一个zip包里面包含多个.job文件每个job定义一条命令job之间用dependencies声明依赖关系。典型的项目结构如下music_analysis/ ├── project.job ├── spark_etl.job ├── mysql_export.job └── chart_report.job# spark_etl.job typecommand commandspark-submit --master yarn --deploy-mode client --executor-memory 4g --num-executors 4 --executor-cores 2 --conf spark.sql.shuffle.partitions200 --class com.music.MusicAnalysis music-analysis.jar --date ${dt}# mysql_export.job typecommand dependenciesspark_etl command/opt/mysql/export_to_mysql.sh --date ${dt}Azkaban会按dependencies的依赖关系生成DAGspark_etl跑完才执行mysql_export。这里的${dt}是Azkaban的调度参数可以在创建Project时配置为20240520这种格式也可以用cron表达式每天自动替换。有一点需要注意deploy-mode client要求提交spark-submit的机器和Azkaban executor在同一台节点上而且要求该机器配置了Hadoop和Spark客户端如果集群是普通用户启动的还要确保Azkaban的executor用户有HDFS写入权限否则提交作业后会在YARN上报错。这个调度流程解决的核心问题是“重跑”。某天数据清洗逻辑改了一行代码需要重算过去一周的数据如果没有Azkaban你只能手动执行七次spark-submit有了它只需要把调度周期改一下或者手动触发带不同${dt}参数的重跑任务即可。5. 数据倾斜与Shuffle优化几十万歌单聚合时的排错顺序网易云音乐的歌曲热门程度呈典型的幂律分布少数头部歌曲贡献了绝大部分播放量。这会导致GROUP BY song_id时某一个或某几个Reducetask处理的数据量远超其他task也就是数据倾斜。具体表现是Spark UI里大部分task几十秒跑完剩下一两个task卡了十几分钟甚至报OOM。这是这类音乐数据分析项目里最容易踩、也最容易在答辩时被问到的性能问题。先讲排查顺序。打开Spark UI的Stages页面看Summary Metrics里的Duration列如果Max和Median差了一个数量级以上基本可以判定倾斜。再用鼠标点开那个最慢的task看Shuffle Read Size是不是远大于其他task。如果是那就要针对倾斜的key做处理。针对热门歌曲倾斜的经典方案是两阶段聚合加盐。所谓加盐就是给原来的key拼接一个随机数前缀把一个大key拆成多个小key完成局部聚合后再去掉前缀做全局聚合。代码示例如下from pyspark.sql.functions import rand, concat, lit, split, col, sum as fsum # 设置分区数让reduce端有足够并行度 spark.conf.set(spark.sql.shuffle.partitions, 300) # 第一阶段加随机盐拆key salted_df df.filter(col(action) play) \ .withColumn(salt, (rand() * 10).cast(int)) \ .withColumn(salted_song_id, concat(col(song_id), lit(_), col(salt))) # 局部聚合同一个盐内先聚合一次 partial_agg salted_df.groupBy(salted_song_id) \ .agg(fsum(lit(1)).alias(cnt)) # 第二阶段去掉盐对局部结果再聚合 final_result partial_agg.withColumn( song_id, split(col(salted_song_id), _).getItem(0) ) \ .groupBy(song_id) \ .agg(fsum(cnt).alias(play_cnt)) \ .orderBy(col(play_cnt).desc()) \ .limit(20)代码里rand() * 10生成0到9的随机整数拼在song_id后面让同一首歌的播放记录被打散到10个不同的key上reduce端就有10个task分担压力。这里的10是盐的粒度粒度太小效果不明显太大会产生过多中间文件常见做法是先试10看最长task耗时降了多少再逐步调到50。第二段代码里split(col(salted_song_id), _).getItem(0)是对字段切割取盐之前的原始song_id。注意fsum(lit(1))比count(*)在部分场景下更可控因为它的返回值类型是Long且不受null值干扰。这个方案有一个前提去盐后的二次聚合数据量已经被大幅压缩否则二次groupBy仍然可能倾斜。如果倾斜的key不多还有一种更轻的做法是加一个WHERE song_id NOT IN (热门key列表)的过滤把热门key单独计算再union回结果但这种方式需要维护一个动态的热门key清单在这个场景下不如加盐通用。再补充一个和Shuffle直接相关的内存参数。之前提到spark.sql.shuffle.partitions但它只影响shuffle输出的分区数不直接影响单个task的内存上限。真正和OOM相关的是spark.executor.memoryOverhead它表示每个executor在堆外内存之外额外预留的内存默认只有executor-memory的10%。如果你用--executor-memory 4goverhead只有约400MB当Shuffle spill严重时这个值不够用建议显式设置为1g到2g。判断依据还是Spark UI里Executor页面的GC时间如果GC TIME超过总运行时间的10%说明堆内存也紧张这时候应该降低executor并发task数而不是盲目加内存。最有效的调优操作其实是复现一次完整流程记录每次调参前后Spark UI里Shuffle Spill和Task Duration两列的变化。先把spark.sql.shuffle.partitions从默认200翻三倍观察Shuffle Read和Spill是否下降再把spark.executor.memoryOverhead调高到1g看GC耗时是否回落。这一步对任何基于Spark做用户行为数据分析的项目都适用也是这份毕业设计里最值得单独拿出来讲的一页PPT。本文还有配套的精品资源点击获取
返回列表