
简介基于Spark平台的TMDB电影数据分析毕设项目提供完整源代码、ECharts可视化页面与配套说明文档适合大数据及相关专业学生用于毕业设计参考、课程设计或是Spark离线分析实战入门。项目围绕TMDB电影数据集展开清洗、聚合与可视化呈现覆盖从数据接入到前端展示的完整链路可作为理解Spark Core/DataFrame编程模型的直观案例。资源共45个文件包含Scala与Python编写的核心分析代码、17个HTML可视化结果页、15个JSON数据文件和6个JavaScript脚本整体约2.55MB目录划分清晰便于按模块阅读。代码全部通过运行验证项目答辩评审平均分高达96分完成度与稳定性有保障。读者可重点参考其数据读取方式、RDD/DataFrame算子调用、统计结果导出至ECharts渲染的完整流程并可在现有代码基础上调整分析维度或二次开发。目前已有244人学习下载适合希望借助完整案例加速上手Spark开发的在校学生与自学者。1. 一份TMDB毕设源码包我为什么拆了它重写分析链路上个月帮人复现一个打包成Zip的Spark TMDB电影数据分析项目目录结构很典型TDMBMain.scala是分析入口Web_Echarts单独放可视化页面README放在最后。数据来自TMDB公开电影信息预算、票房、评分、类型、语言、制作公司这些字段都能拿来练Spark DataFrame、窗口函数和聚合统计。这类项目特别适合做毕设或课程设计因为它是“数据接入→清洗→统计分析→可视化”的完整闭环改几个统计口径就能当新作业交。但我拆完发现不少工程问题SQL字符串写死在代码里、结果无脑coalesce(1)、ECharts读的是手工复制粘贴后的Json。我按自己平时做数据平台的习惯把链路重排了一遍下面讲的是能从本地跑到YARN集群的完整方案。2. TMDB数据建模与Spark DataFrame选型2.1 TMDB字段结构先看懂CSV再动代码TMDB数据集通常以CSV或NDJSON格式发布仓库里那份是CSV。拿到手先别急着开Spark用任意文本编辑器看前50行确认三件事有没有header、字段里是否带逗号和换行、日期和数字列有没有缺失值。常见字段和它们的使用场景如下字段名类型分析价值常见脏数据budgetInteger与票房对比算ROI0或空值占比高revenueLong票房统计核心缺失记为0需要过滤genresString类型统计需JSON解析空字符串或格式不统一original_languageString分语言对比大小写不一致release_dateString按年份聚合存在多种日期格式vote_averageDouble评分画像投票人数过少时虚高vote_countInteger评分可信度权重0值常见runtimeInteger时长分布0表示缺失应过滤字段解析顺序直接决定后面清洗成本。原项目用spark默认类型推断本来能跑但遇到release_date缺值那一行整列会被推断成string后面的year()函数当场报错。这个问题在真实数据集里几乎必然出现所以建议显式定义schema不要依赖推断。数据源方面TMDB官方CSV在国内访问不稳定我一般从国内镜像站或HuggingFace离线数据集仓库拿到同构文件再跑字段顺序先和本地对齐省得改代码。2.2 为什么选DataFrame而不是RDD这个项目选DataFrame有三个理由。一是数据量在几GB级别DataFrame的Tungsten二进制存储和Catalyst查询优化能把执行计划压缩掉一大截RDD的Java对象序列化开销在这里是纯浪费二是后续所有统计都走SQL语义groupBy、join、窗口函数直接表达不需要手写reduceByKey三是写起来短毕设答辩时能讲清楚“我把筛选条件下推到文件读取阶段”这种优化点。RDD只在一种场景下更有优势底层自定义分区或者手写复杂UDF。本项目的genres解析用from_json加explode就能解决不需要下探到RDD层面。2.3 显式Schema定义与CSV加载参数确认字段结构后先定义schema再读入DataFrameimport org.apache.spark.sql.SparkSession import org.apache.spark.sql.types._ val spark SparkSession.builder() .appName(TMDB Movie Analysis) .master(local[*]) .getOrCreate() val schema StructType(Array( StructField(budget, IntegerType, true), StructField(genres, StringType, true), StructField(original_language, StringType, true), StructField(popularity, DoubleType, true), StructField(release_date, StringType, true), StructField(revenue, LongType, true), StructField(runtime, IntegerType, true), StructField(title, StringType, true), StructField(vote_average, DoubleType, true), StructField(vote_count, IntegerType, true) )) val df spark.read .option(header, true) .option(multiLine, true) .option(escape, \) .schema(schema) .csv(tmdb_movies.csv)这里三个option各有用途header必须打开multiLine允许字段值内含换行符escape处理引号嵌套TMDB的overview字段里大量出现逗号和换行不打开这两个选项会出现列错位。revenue用LongType是因为部分电影票房超过Int上限。注意multiLine选项会影响整列解析性能CSV的转义和回车会让解析器做更多状态判断确认数据里确实有换行再开启不要无脑照抄。2.4 清洗三件套空值、零值、日期格式清洗阶段最容易出的问题是把“缺数据”和“零数据”混为一谈。budget和revenue为0不代表真实值是0很多早期电影根本没收录预算做对比分析前必须过滤这类记录import org.apache.spark.sql.functions._ val cleaned df .filter(col(vote_count) 10) .filter(col(budget) 0 col(revenue) 0) .filter(col(runtime) 0) .withColumn(year, year(to_date(col(release_date), yyyy-MM-dd))) .filter(col(year).isNotNull col(year) 1900) .na.drop(Seq(title, vote_average))筛选条件顺序有讲究vote_count过滤最先执行它在分析中是可信度门槛提前把噪声数据裁掉能减少后续shuffle数据量日期解析失败会返回null用isNotNull兜底而不是一开始就drop能保留那些日期格式异常但其他字段完整的记录方便排查。na.drop只作用于关键字段不会把genres为空但其他字段正常的行误删。3. TDMBMain.scala分析链路清洗、聚合与Top-N统计3.1 年度评分与票房画像TDMBMain这个类名保留了原项目笔误正式命名应该是TMDBAnalysisMain但为了不打乱工程结构我沿用了原名。类里只放一个main入口不要写长SQL字符串用DataFrame DSL链路更利于定位错误。先做年度画像val yearlyStats cleaned .groupBy(year) .agg( count(*).alias(movie_count), round(avg(vote_average), 2).alias(avg_rating), round(avg(revenue), 0).alias(avg_revenue) ) .orderBy(year)groupBy之后Spark会按year做一次hash聚合如果文件特别大且年份分布不均匀会有一个task成为长尾。如果某一年份数据量占比过大光加executor数量没用需要按第5章的方法做加盐或开AQE这里先按常规聚合写。3.2 按语言分组取评分前十窗口函数比groupBy更合适需求“每种语言取评分前十”用groupBy做不出来必须用窗口函数import org.apache.spark.sql.expressions.Window val top10ByLang cleaned .withColumn(rn, row_number().over( Window.partitionBy(original_language) .orderBy(col(vote_average).desc, col(vote_count).desc) )) .filter(col(rn) 10) .drop(rn)窗口函数这里有两个细节容易搞错。第一orderBy用vote_average降序做主排序但评分并列时row_number是随机分配必须加vote_count做次级排序否则同一部电影在不同次运行里可能排进或排出前十。第二row_number、rank、dense_rank三者的区别要说清楚row_number严格排序不重号rank并列会留空位dense_rank并列不留空位。如果毕设里用rank前十会变成十二条记录答辩时大概率被追问。3.3 类型统计from_json与explode拆解genres字段TMDB的genres是JSON字符串比如[{id: 28, name: Action}, {id: 12, name: Adventure}]。要把数组打平成多行用explode配合from_jsonimport org.apache.spark.sql.functions._ val genreSchema ArrayType(StructType(Seq( StructField(id, IntegerType), StructField(name, StringType) ))) val genreStats cleaned .withColumn(genre, explode(from_json(col(genres), genreSchema))) .groupBy(genre.name) .agg( count(*).alias(movie_cnt), round(avg(vote_average), 2).alias(avg_rating), round(sum(revenue), 0).alias(total_revenue) ) .orderBy(col(movie_cnt).desc)from_json的第二个参数必须与JSON字符串真实结构一致这里genres数组的元素是对象所以要声明ArrayType(StructType(...))如果写成ArrayType(StringType)整列会解析成null。排错时先在spark-shell里单独select一个from_json结果观察返回再接explode能省很多时间。explode遇到空数组自动产零行所以genres为空字符串的记录在这步被自然过滤不需要额外写filter。3.4 结果输出不要只写coalesce(1)原项目把所有结果都coalesce(1)写出理由是前端省事。但这是集群作业里最不该养成的习惯coalesce(1)把所有分区数据拉到单个task写文件数据量一大必然OOM。我的做法分场景处理场景推荐写法原因结果集小于1万行coalesce(1).write.csv文件数少本地和前端都好读结果集超过10万行partitionBy加parquet避免单task写出OOM利于下次过滤前端需要Jsoncollect后手动拼Json可控制数值精度与字段顺序代码按这个表来// 小结果集年份统计这类几KB的数据coalesce(1)没问题 yearlyStats .coalesce(1) .write .mode(overwrite) .option(header, true) .csv(output/yearly_stats) // 大结果集保持分区写出后续合并 genreStats .write .mode(overwrite) .partitionBy(genre_name) .parquet(output/genre_stats.parquet)判断标准是结果集行数几千行直接合并没有问题超过十万行就必须分区写或走Parquet。Parquet列式存储对按类型过滤的查询明显快于CSV还天然携带schema给后续Web端解析也更友好。4. Web_Echarts可视化对接从Spark Json到前端图表4.1 结果Json化的两种方案仓库里的Web_Echarts目录只放HTML和Js没有后端接口。我建议把Spark结果导出成单个固定名的Json文件省掉接口层毕设演示最稳import scala.collection.mutable.ListBuffer val rows yearlyStats.collect() val buf new ListBuffer[String] buf [ rows.zipWithIndex.foreach { case (r, i) val comma if (i rows.length - 1) else , buf s{year:${r.getAs[Int](year)},avg_revenue:${r.getAs[Double](avg_revenue)}}$comma } buf ] spark.sparkContext.parallelize(buf.toSeq, 1) .saveAsTextFile(web/data/yearly_stats.json)这里用collect加saveAsTextFile而不是coalesce(1)写csv是因为前端拿到Json可以直接消费省去解析CSV的依赖。collect的约束是结果集必须能塞进driver内存几千行统计没问题如果要导全量明细这个方法会撑爆driver。数据量大的场景改用df.write.json输出到HDFS再用hadoop fs -getmerge拉回本地。Spark写Json默认每个分区产出一个part-00000文件前端fetch并不知道分区文件名因此需要合并。getmerge是最快的办法hadoop fs -getmerge hdfs:///output/yearly_stats_json/ web/data/yearly_stats.json4.2 ECharts怎么吃这个Json前端部分按ECharts 5的写法在Web_Echarts目录里加一个页面核心代码如下script srchttps://cdn.jsdelivr.net/npm/echarts5/dist/echarts.min.js/script div idchart stylewidth: 100%; height: 480px;/div script fetch(data/yearly_stats.json) .then(res res.json()) .then(rows { const chart echarts.init(document.getElementById(chart)); chart.setOption({ title: { text: TMDB电影年度平均票房 }, tooltip: { trigger: axis }, xAxis: { type: category, data: rows.map(r r.year) }, yAxis: { type: value, name: 平均票房 }, series: [{ type: bar, data: rows.map(r r.avg_revenue) }] }); }) .catch(err console.error(Json load failed:, err)); /scriptfetch用的是相对路径所以HTML文件不能以file://直接双击打开浏览器会拦截跨源请求。我一般用python3 -m http.server 8080在当前目录起一个静态服务然后访问localhost:8080。series.type换成line就是走势图换成pie之前需要先把数据按降序排好这些都不用改Spark代码。4.3 字段口径对应与两个坑字段映射关系如下Spark输出字段Json键ECharts配置项yearyearxAxis.dataavg_revenueavg_revenueseries.dataavg_ratingavg_rating第二个图表series.datamovie_countmovie_counttooltip自定义内容两个坑值得单独说。第一个是精度。Spark的Long写到Json后超过2^53会丢精度JavaScript的Number安全整数上限是9007199254740991。单部电影revenue达不到但sum(revenue)聚合会轻松越过这个值。稳妥做法是在Spark侧把数值cast成Decimal(20,0)或者直接格式化成字符串再导出。第二个是空值。Json里出现null时ECharts柱状图会有一段空缺。处理方式是在Spark侧用coalesce(col(avg_revenue), lit(0))补零或者在前端filter掉null行。我更倾向Spark侧兜底前端代码能保持简单。5. YARN提交、内存调优与Spark任务验证方法5.1 从本地到集群spark-submit参数不要照抄本地用local[*]跑通只是第一步放到集群上的spark-submit参数直接决定任务成败。我常用的模板spark-submit \ --class TDMBMain \ --master yarn \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 8 \ spark-tmdb-assembly.jar \ hdfs:///data/tmdb_movies.csv hdfs:///output/tmdb_analysis参数作用调优建议--driver-memorydriver端可用内存collect前先估算返回行数--executor-memory每个executor堆内存不能超过单机物理内存上限--executor-cores单个executor核数2到3为宜越大GC越严重--num-executorsexecutor总数按队列配额上限取不是越多越好executor-memory和executor-cores的比例是关键单个executor核数超过3时JVM GC会成为瓶颈遇到Container killed去YARN UI看是内存超限还是虚拟内存超限虚拟内存问题通常要调yarn.nodemanager.vmem-check-enabled但这是集群管理员权限自己搭的测试集群才建议动。另外在IDEA里跑local模式时task运行在executor线程而不是driver线程IDE经常提示“当前不会命中断点”调试Spark任务应该靠日志加collect而不是依赖断点。5.2 Spark内存与OutOfMemory的真相很多Spark作业死在OOM但OOM分两种处理方式完全相反driver端OOM是collect回来的数据太大解决方法是分页collect或先聚合再collectexecutor端OOM是单个分区数据量超过executor可用内存解决方法是增大分区数而不是增大executor内存。这个项目里如果报executor OOM优先加大输入文件读取并行度而不是调executor-memoryspark.sql.adaptive.enabled(true) spark.sql.adaptive.coalescePartitions.enabled(true)开AQE后Spark会在shuffle结束动态合并小分区对TMDB这种年份分布不均匀的数据集很有效。验证分区是否合并去Spark UI的SQL页签看“Shuffle Query Stage”的分区数变化即可。5.3 验证Top-N口径对拍命令统计写完我习惯用一个对拍方法确认结果没被写偏。拿“每语言评分前十”举例直接拉全量数据在本地过滤对比行数# 用spark-shell验证 val a spark.read.parquet(output/genre_stats.parquet).count() val b spark.read.csv(hdfs:///data/tmdb_movies.csv).filter(...).count() println(s$a vs $b)如果对不上优先查清洗环节八成是TopN结果里混进了vote_count小于阈值的记录。验证通过后把结果文件用hadoop fs -getmerge拉下来替换Web_Echarts的data目录整个过程不用改前端代码。Spark UI里看每个stage的shuffle read大小能反推哪一步聚合消耗最大比如所有记录按original_language分区时单个task拉的数据超过500MB就要考虑加盐拆key。把中间结果缓存成Parquet后二次迭代分析能直接复用文件整条链路跑下来比原项目省一半时间参数和口径也都留在了代码里换数据源时只需改schema和路径。本文还有配套的精品资源点击获取