又到一年毕业季我后台收到最多的私信就是“猫眼电影数据分析与票房预测系统怎么做”选这个题目的同学确实不少Python爬虫、Hadoop、Spark、Django、协同过滤、随机森林一个题目把大数据和机器学习的主流技术栈全占了光是列在简历上就很有分量。但真动手做起来网上资料东一块西一块好多教程还停留在“装好环境就跑个wordcount”的程度离一个能展示、能答辩、能跑通全流程的系统差得远。这篇就把我完整做完这个项目的思路和代码细节全部摊开来讲从架构设计、数据采集、Spark清洗、特征工程、协同过滤推荐、随机森林票房预测到Django后端接口和可视化展示一条线串下来。无论你是刚把Python语法过完的小白还是已经能写爬虫但没碰过Spark的同学按这个路线走都能把这个系统从零搭出来。文章会比较长建议先收藏再慢慢看。1. 项目整体设计与技术选型1.1 系统到底要做什么先捋清楚需求。猫眼电影数据分析与票房预测系统拆开来看就四个核心功能第一能自动采集猫眼电影数据包括电影基本信息、评分、评论、票房这些第二把采集到的数据做清洗和分析比如票房走势、评分分布、类型占比、地区票房对比第三给用户做电影推荐这一步用的是协同过滤第四基于历史数据预测一部新电影的票房区间这一步用随机森林回归模型。完整的数据链路是这样Python爬虫从猫眼采集数据原始数据存到Hadoop HDFS做分布式存储再用Spark做清洗和特征工程清洗后的结果落到MySQL供Django后端查询。协同过滤推荐和随机森林票房预测这两个算法模型一个在Spark MLlib里训练一个在Python sklearn里训练最后通过Django REST Framework提供接口前端用ECharts把分析结果可视化展示出来。很多同学一上来就问“能不能不用Hadoop和Spark直接用Pandas处理”技术上当然可以但那样这个题目就失去了一半的意义。毕业设计不只是把功能做出来更要展示你对大数据技术栈的理解。Hadoop和Spark的价值在于处理海量数据的分布式能力面试官问你“为什么用Spark而不用Pandas”时你要能答上来“因为当评论数据达到千万级别单机内存根本扛不住Spark基于RDD的懒加载机制和分布式计算可以在多节点并行处理”。1.2 技术栈选型的思路对比先说我踩过的坑。一开始我把所有数据都塞进MySQL一张评论表存了上百万条之后统计类SQL查询直接卡到十几秒这还没算爬虫不断写入的压力。后来才把原始日志和中间结果挪到HDFSMySQL只保留分析后的结果表和业务表查询立刻快了一个数量级。这套技术栈的选择逻辑很简单每层解决每层的问题层级技术选型解决什么问题替代方案对比数据采集Python requests BeautifulSoup从猫眼爬取电影与评论数据Scrapy更重但更规范毕设用requests足够分布式存储Hadoop HDFS海量原始数据存储本地文件过于简单体现不了大数据场景数据处理Spark SQL DataFrame清洗、聚合、特征提取单机Pandas在大数据量下内存不够业务数据库MySQL存储分析结果和业务数据SQLite太简单MySQL最通用推荐算法Spark MLlib ALS协同过滤基于用户评论评分做电影推荐Python实现ItemCF也OK但MLlib更贴合大数据场景预测算法Python sklearn 随机森林回归票房区间预测XGBoost效果略好但随机森林更容易解释答辩好讲Web后端Django REST Framework提供接口给前端Flask轻但需要自己拼凑组件Django自带Admin和ORM开发效率高前端展示ECharts HTML/JS数据可视化大屏纯静态页面即可不需要Vue这类框架选这套方案最大的好处是每个技术点都有明确的“存在理由”。答辩的时候老师问任何一个组件你都能讲清楚它在这个系统里干了什么、为什么非它不可这就比“我看教程里用了所以我用了”高一个层次。1.3 系统整体架构和各模块职责把系统按层次划分的话从上到下是应用展示层、业务逻辑层、算法模型层和数据存储层。应用展示层就是浏览器里的可视化大屏包含票房排行、评分分布、类型占比、推荐结果、预测结果等模块。业务逻辑层是Django写的REST API负责把数据从MySQL读出来转成JSON返回给前端。算法模型层两个模块协同过滤走Spark MLlib的ALS算法票房预测走sklearn的随机森林回归。数据存储层是HDFS加MySQL的组合HDFS放原始评论和日志数据MySQL放电影表、票房表、分析结果表。这个分层架构的核心思想是解耦。爬虫只管采集Spark只管处理Django只管提供服务前端只管展示。每一层都可以独立测试答辩时演示也很有条理——先看爬虫运行再看Spark任务日志最后打开网页看结果。2. 数据采集与预处理爬虫是第一步也是最重要的一步2.1 爬虫目标分析需要哪些数据字段做数据分析系统数据是地基。我定义的爬虫目标包括两个数据源电影列表页和电影评论页。电影列表页拿到的字段有电影名、主演、导演、上映日期、类型、评分、评论数、票房部分电影列表页有这些是电影基本信息。电影评论页拿到的是评论用户ID、用户名、评分、评论内容、评论时间这些是协同过滤推荐算法的输入数据。最终这两类数据分别存成两张表movie_info和user_comment。评分这里有个坑猫眼的评分是整数1到5星加小数点后一位的分数比如4.5分、3.0分但评论用户的个人打分是整数1到10分爬取的时候要注意区分推荐算法用的是评论用户打分统计展示用的是猫眼评分。2.2 爬虫编写的关键细节与反爬策略猫眼的反爬不算特别变态但有几个点必须处理。首先是User-Agent不用默认的Python-requests标识要伪装成浏览器。其次是请求频率必须加随机延时我一般用time.sleep(random.uniform(1, 3))太快了IP会被封。最后是字体反爬猫眼的某些数字是动态字体渲染的爬下来可能是乱码。这个我在评论区数据上遇到过处理方案是下载字体文件解析映射关系或者直接换一种数据来源——比如从页面JSON接口里提取原始分数数据而不是解析渲染后的HTML文本。下面是我爬取电影信息的核心代码框架请求加了解析每一步都有注释import requests import time import random import pandas as pd from bs4 import BeautifulSoup HEADERS { User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36, Accept-Language: zh-CN,zh;q0.9, Referer: https://maoyan.com/ } def fetch_movie_list(offset0): 抓取猫眼电影列表页offset控制分页 url fhttps://maoyan.com/films?offset{offset} resp requests.get(url, headersHEADERS, timeout10) soup BeautifulSoup(resp.text, html.parser) movies [] for item in soup.select(.movie-item): title_tag item.select_one(.movie-title) score_tag item.select_one(.movie-score) if not title_tag: continue movies.append({ title: title_tag.text.strip(), score: score_tag.text.strip() if score_tag else None }) return movies def run_crawler(): all_movies [] for offset in range(0, 100, 30): # 控制采集页数 try: movies fetch_movie_list(offset) all_movies.extend(movies) print(f已抓取 offset{offset}, 累计 {len(all_movies)} 条) except Exception as e: print(f抓取失败 offset{offset}: {e}) time.sleep(random.uniform(1, 3)) # 随机延时防反爬 df pd.DataFrame(all_movies) df.to_csv(maoyan_movies.csv, indexFalse, encodingutf-8-sig)有一点我必须强调爬到的数据一定要先落成CSV或JSON存一份然后再导入数据库。这样就算后面清洗逻辑写错了原始数据还在不用重新爬一遍。我一开始就是直接爬一条写一条到MySQL后来发现有个字段解析错了改完代码只能重爬白白浪费了大半天时间。2.3 数据清洗规则与入库设计原始数据不可避免会有缺失、重复、格式不一致的问题。爬虫阶段拿到的数据不能直接用我整理了一套清洗规则第一去重。相同title加相同上映日期的记录只保留一条用DataFrame的drop_duplicates(subset[title, release_date])解决。第二缺失值处理。评分为空的记录如果是统计用就保留但标记为None如果是训练模型用就直接剔除票房字段有空值用同类电影的平均值填充。第三类型字段规范化。猫眼的类型可能是“剧情/爱情”这种斜杠分隔格式要拆成主类型和副类型方便后面按类型统计。第四上映日期统一成datetime类型因为后面按档期分析要按月份、季度分组。数据入库我用的是pandas的to_sql方法一条命令就能把DataFrame写进MySQL但要注意先建好表和索引from sqlalchemy import create_engine engine create_engine(mysqlpymysql://root:passwordlocalhost:3306/maoyan_db?charsetutf8mb4) df.to_sql(movie_info, engine, if_existsreplace, indexFalse, chunksize1000)评论数据同理但建议按movie_id做分区或者至少建索引不然后续做协同过滤时要扫描全表慢到怀疑人生。3. Spark大数据处理与特征工程把数据变成模型能吃的样子3.1 为什么处理层必须上Spark这部分是我在答辩时讲得最多的地方。用户评论数据爬到几十万条之后用Pandas做groupby聚合已经有点吃力了更别说后面还要算用户-物品评分矩阵。评分矩阵在Python里就是几千乘几千的二维数组内存还好说但数据清洗和特征提取的全量计算量是真的大。Spark的RDD和DataFrame把这套计算分布到集群多个节点上并行执行处理千万级数据也就是分钟级的事。Spark对新手最友好的地方是DataFrame API跟Pandas长得非常像会Pandas的基本能无缝迁移。但背后逻辑完全不同——Pandas是一次性把数据load进内存操作Spark是懒加载遇到action操作才真正触发计算中间过程全是DAG调度。这个区别我建议你务必理解因为这是面试高频题。3.2 Spark作业的编写与提交我用的是Spark SQL做清洗和聚合统计核心代码逻辑分成三步从HDFS读原始评论数据、清洗过滤、按业务维度聚合。from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg, to_date, year, month, desc spark SparkSession.builder \ .appName(MaoyanDataAnalysis) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() # 1. 从HDFS读取评论数据 df spark.read.csv(hdfs://localhost:9000/maoyan/comment_data.csv, headerTrue, inferSchemaTrue) # 2. 清洗过滤空评论、去掉评分缺失的 df_clean df.filter(col(content).isNotNull() col(score).isNotNull()) # 3. 按电影统计评论数与平均评分 movie_stats df_clean.groupBy(movie_id, movie_title) \ .agg(count(content).alias(comment_count), avg(score).alias(avg_score)) \ .orderBy(desc(comment_count)) movie_stats.show(20)跑完之后结果数据我用同样的方式写回MySQL。这里要注意一个坑Spark写MySQL时如果表里有大量数据默认的写模式会非常慢建议用partitionBy按日期或电影ID分区写入然后设置batchsize批量插入。3.3 特征工程构建可用的机器学习特征集推荐算法和票房预测模型的输入都不是原始数据而是处理后的特征。票房预测模型我用了下面这些特征特征字段类型构造方式导演历史票房均值数值型从movie_info里按导演聚合历史票房取平均主演热度数值型演员参与电影的总票房/平均评分电影类型类别型多标签one-hot编码上映档期类别型春节档、暑期档、国庆档、平日档首日评论量数值型评论时间分组统计猫眼评分数值型爬取的基础字段预告片播放量数值型这个字段需要扩展爬虫可选特征不是越多越好关键是每个特征都要有业务含义。我一开始把“电影ID”也当特征丢进模型结果模型完全过拟合预测结果失去参考价值。后来才意识到ID是标识符不是特征移除之后模型效果才恢复正常。这个教训写在这里希望你别踩同样的坑。4. 算法模型协同过滤推荐与随机森林票房预测4.1 协同过滤Spark MLlib的ALS实现协同过滤分两种基于用户的UserCF和基于物品的ItemCF。用户数量远大于电影数量而且用户评分数据稀疏UserCF计算用户相似度矩阵会非常庞大。所以我选的是基于物品的ItemCF——计算电影之间的相似度然后根据用户看过的电影推荐相似电影。这个思路就好理解你给一个看过《流浪地球》的用户推荐《星际穿越》因为两部电影被大量相同用户打过高分所以在向量空间里它们距离很近。Spark MLlib里实现ItemCF最方便的方式其实是ALS交替最小二乘法。ALS本质上做的是矩阵分解把用户-物品评分矩阵分解成两个低维矩阵的乘积这样既能降维又能补全未评分的部分。训练代码很简单from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 评分数据user_id, movie_id, rating ratings spark.read.csv(hdfs://localhost:9000/maoyan/ratings.csv, headerTrue, inferSchemaTrue) # 拆分训练集和测试集 train, test ratings.randomSplit([0.8, 0.2], seed42) als ALS( maxIter10, regParam0.1, userColuser_id, itemColmovie_id, ratingColrating, coldStartStrategydrop ) model als.fit(train) predictions model.transform(test) evaluator RegressionEvaluator(metricNamermse, labelColrating, predictionColprediction) rmse evaluator.evaluate(predictions) print(fRMSE {rmse})踩过的坑ALS的regParam正则化系数非常敏感默认是0.1但如果评分数据特别稀疏这个值偏小容易过拟合RMSE会居高不下。我后来调到0.5才稳定下来。另一个坑是coldStartStrategy必须设置成drop否则评分矩阵里有新用户或新电影时预测结果会返回NaN后续评估直接报错。4.2 随机森林从原理到调参实践随机森林是集成学习里最容易讲清楚的一个模型。它训练多棵决策树每棵树用训练集的随机子集Bootstrap采样和随机特征子集来训练最后回归任务取所有树的预测均值。为什么用随机森林做票房预测而不是线性回归因为票房跟特征之间的关系高度非线性比如“档期”的影响——国庆档一部普通电影可能比平日档的好电影票房还高这种关系线性模型很难拟合树模型天然能捕捉。我在sklearn里的实现代码如下from sklearn.ensemble import RandomForestRegressor from sklearn.model_selection import train_test_split, GridSearchCV from sklearn.metrics import mean_absolute_error, r2_score import pandas as pd # 加载特征数据 features pd.read_csv(movie_features.csv) X features.drop(columns[movie_id, box_office]) y features[box_office] # 训练集测试集划分 X_train, X_test, y_train, y_test train_test_split( X, y, test_size0.2, random_state42 ) # 随机森林回归模型 rf RandomForestRegressor( n_estimators200, # 树的数量 max_depth12, # 最大深度 min_samples_leaf4, # 叶子节点最小样本数 max_featuressqrt, # 每次分割随机选特征数 random_state42, n_jobs-1 # 使用全部CPU核心 ) rf.fit(X_train, y_train) # 评估 y_pred rf.predict(X_test) print(fMAE {mean_absolute_error(y_test, y_pred)}) print(fR2 {r2_score(y_test, y_pred)}) # 特征重要性排序 importance pd.DataFrame({ feature: X.columns, importance: rf.feature_importances_ }).sort_values(importance, ascendingFalse) print(importance)随机森林的超参数没有标准答案要结合数据量调。我给的这套参数组合来自我实际跑出来的经验值数据量在几百部电影时n_estimators200已经足够再大收益很小反而增加耗时max_depth设12是因为特征维度本身不多太深容易过拟合训练集。很多新手一上来就n_estimators1000以为树越多越好实际只是白白增加训练时间。4.3 模型评估与结果解读模型训练完不能只看准确性还要看业务上有没有意义。票房预测的RMSE如果是一亿对于小成本电影来说误差就非常大对于大制作来说可能还能接受。所以我在评估时把预测结果分成了三个区间低成本1亿、中成本1亿到5亿、高成本5亿分别计算准确率。这样做的好处是答辩时能讲出更细的分析结论而不仅仅是“模型R2达到了0.85”这种干巴巴的指标。特征重要性输出是答辩时另一个加分项。我跑出来的结果里“导演历史票房均值”和“主演热度”贡献了超过50%的重要性“首日评论量”排第三“评分”反而靠后。这个结论跟行业认知是一致的——票房更依赖主创阵容和宣发热度评分对票房的解释力有限。5. Django后端与可视化展示把数据变成能看的页面5.1 Django项目结构与API设计后端我用的是Django REST Framework。Django的好处是全套组件都有ORM、Admin后台、认证系统、模板引擎全内置了写业务代码不用到处找第三方库。项目结构我建议这样组织maoyan_backend/ ├── manage.py ├── config/ # 项目配置 │ ├── settings.py │ ├── urls.py │ └── wsgi.py ├── api/ # API接口应用 │ ├── models.py # 数据模型 │ ├── views.py # 视图函数 │ ├── serializers.py # 序列化器 │ └── urls.py ├── analysis/ # 数据分析应用 │ ├── spark_jobs.py # Spark任务调用 │ └── recommend.py # 推荐算法封装 └── static/ # 前端静态文件API接口设计是按前端页面需求来的不需要设计得很复杂够用就行接口路径方法功能/api/movies/GET电影列表分页查询/api/movies/stats/GET电影统计评分分布、类型占比/api/movies/boxoffice/trend/GET月度票房趋势/api/recommend/?user_idxxxGET给指定用户推荐电影/api/predict/boxoffice/POST传入电影特征返回票房预测/api/comments/GET评论数据查看5.2 推荐结果接口的代码实现推荐接口的封装逻辑比较简单训练好的ALS模型保存下来Django启动时加载模型收到请求后调用模型的recommendForUserSubset方法拿到推荐列表再查MySQL补全电影信息返回JSON。# recommend.py from pyspark.ml.recommendation import ALSModel from pyspark.sql import SparkSession from django.http import JsonResponse from .models import Movie spark SparkSession.builder \ .master(local[2]) \ .appName(RecommendService) \ .getOrCreate() # 启动时加载HS保存的模型 als_model ALSModel.load(hdfs://localhost:9000/models/als_model) def recommend_for_user(request): user_id request.GET.get(user_id) if not user_id: return JsonResponse({error: missing user_id}, status400) # 生成推荐 users spark.createDataFrame([(int(user_id),)], [user_id]) recs als_model.recommendForUserSubset(users, 10) recs_list recs.collect()[0].recommendations # 解析推荐电影ID列表 movie_ids [row.movie_id for row in recs_list] # 查MySQL补全电影信息 movies Movie.objects.filter(movie_id__inmovie_ids) result [{movie_id: m.movie_id, title: m.title, score: m.score} for m in movies] return JsonResponse({user_id: user_id, recommendations: result})这里有个性能问题要说明每次请求都创建SparkSession是不行的非常慢且消耗资源。正确做法是Django启动时创建全局唯一实例单例模式后面所有请求复用它。我当时也是踩了这个坑第一次请求等了十几秒后面才反应过来SparkSession初始化开销大。5.3 ECharts可视化大屏搭建前端我用的是原生HTML加ECharts没有用Vue或者React。理由很简单毕设阶段不追求工程复杂度只要把数据展示清楚就行。ECharts做数据可视化大屏的优势是组件丰富、文档全、社区资源多从折线图、柱状图到饼图、地图都开箱即用。我页面里放了六个图表模块顶部是系统标题和数据总量摘要卡片左侧是票房Top10柱状图和类型占比玫瑰图中间是月度票房趋势折线图和评分分布直方图右侧是推荐结果卡片和票房预测输入表单。整个布局用CSS Grid实现1920x1080分辨率下可以全屏展示答辩时视觉效果很好。前端获取数据就是fetch调用Django接口然后setOption更新图表。以票房趋势图为例async function loadBoxOfficeTrend() { const response await fetch(/api/movies/boxoffice/trend/); const data await response.json(); const chart echarts.init(document.getElementById(trendChart)); chart.setOption({ title: { text: 月度票房趋势, left: center }, tooltip: { trigger: axis }, xAxis: { type: category, data: data.months }, yAxis: { type: value, name: 票房亿 }, series: [{ name: 票房, type: line, data: data.box_offices, smooth: true, // 曲线平滑 areaStyle: { opacity: 0.3 } }] }); } loadBoxOfficeTrend();一个细节Django开发服务器默认只开一个线程前端同时请求多个接口时如果某个接口因为Spark计算慢而阻塞其他接口会排队等待。建议在settings.py里配置一下线程支持或者把耗时的算法接口放到Celery任务队列里。毕设场景下用django的runserver加上多线程选项基本够用了。6. 常见问题排查与避坑指南6.1 环境与版本兼容问题速查表这一节是血泪史。我整个项目大概有三分之一的时间花在解决环境问题上这里直接把最常踩的坑列成表格希望你能少走这些弯路。问题现象根本原因解决办法Spark启动报Java版本错误Hadoop 2.x和Spark 3.x对Java 8/11要求不同统一JDK1.8不要用默认的Java 17HDFS NameNode无法启动未格式化或多次格式化导致存储状态异常先停掉所有进程删除tmp目录下的dfs数据重新执行hdfs namenode -formatSpark写MySQL中文乱码MySQL连接串没指定utf8mb4编码在JDBC URL加characterEncodingutf8useUnicodetrueDjango跨域请求被拦截前端页面和后端API不同端口安装django-cors-headers并在settings中注册MiddlewareALS预测结果全是NaN没设置coldStartStrategydrop模型训练时显式配置该参数爬虫被封IP请求频率太高无代理池降低请求频率加随机延时有条件就用代理IPRandomForest训练内存溢出数据集未标准化但树太深调小max_depth增加min_samples_leaf前端图表不显示数据接口返回格式不是ECharts预期的结构打开浏览器开发者工具查看Network里响应JSON结构6.2 Spark与Hadoop整合的深层排查案例挑一个最典型的展开说Hadoop和Zookeeper整合的问题。我用的是Spark Standalone模式不需要Hadoop YARN做资源调度。但Hadoop本身需要运行在分布式文件系统模式下HDFS要求集群中节点状态一致。有一次我重启了虚拟机HDFS一直报“There appears to be a gap: expected block but received X”查了半天发现是数据节点的心跳超时导致NameNode认为节点挂了。排查思路是这样的先看NameNode日志发现报的是block replica缺失然后用hdfs fsck命令检查文件完整性发现确实有块丢失最后定位到是Windows虚拟机重启后网卡IP变了dataNode注册的地址和NameNode记录的地址对不上。解决办法是给每台虚拟机配静态IP并且重新格式化NameNode后把所有数据重新上传。这个案例在答辩时讲出来非常有说服力证明你确实自己排查过问题不是照着教程敲的。6.3 调试与验证的实用技巧最后聊几个提高效率的调试方法。第一先用小数据量验证全流程。别一上来就跑全量数据否则出了问题很难定位是哪个环节错了。我在写Spark清洗逻辑时先拿100条数据放到本地文件测试确认统计结果对了再扔到HDFS上跑全量。第二日志分级是关键。可以给每个模块加上不同级别的日志爬虫用DEBUG级别记录抓取详情Spark作业用INFO级别记录执行阶段Django接口用WARNING级别记录异常。这样排查问题时能看到关键节点的执行状态而不是全靠print输出。第三模型保存和加载要提前设计好。ALS模型训练好后用save方法持久化到HDFS随机森林用joblib.dump保存到本地。这样Django服务启动时直接加载不用重新训练接口响应速度会快很多。最后再分享一个实用经验整套系统做完回头看最核心的收益不是代码本身而是把“数据采集-数据清洗-数据分析-算法建模-Web展示”这条完整链路跑通的经验。很多人写毕设时容易陷入一个误区要么只关注爬虫和数据分析结果页面很简陋要么只关注Django页面展示结果底层算法都是假的。实际上评阅老师更看重的是系统完整性和逻辑自洽性哪怕每个环节做得不深但只要链路是通的、每个模块都有真实的数据在流转就已经达到毕业设计的标准了。如果你也想在这个题目上做些扩展后续可以考虑两个方向一是把票房预测模型从回归改成分类预测“爆款/中等/扑街”三档准确率会直观很多二是给推荐系统加上简单的Web交互反馈用户点“不喜欢”后实时更新推荐结果。这两个方向我在项目里做了初步验证效果都还不错但篇幅有限就不展开了。先按上面的路线把基础版本做出来跑通再来聊进阶玩法。