简介:一份面向毕业设计场景的电影智能推荐系统完整实现,整合了Flask Web框架、Spark分布式计算、ALS协同过滤算法与MovieLens公开评分数据集。项目覆盖数据清洗、特征处理、推荐模型训练与网页交互展示等关键环节,适合计算机相关专业的学生用于毕业设计参考,也适合推荐系统初学者跟随项目代码理解完整实现流程。压缩包共50个文件,包含10个Python源码文件、10个HTML页面模板、4个CSV评分数据集,以及图片、配置和说明文档,整体仅6.59MB,结构清晰,便于本地部署与学习。目前已有66人学习下载。资源中提供了可运行的Flask服务、Spark处理脚本、ALS训练代码和项目说明,能够帮助读者快速搭建一个完整的电影推荐示例,并在现有代码基础上替换数据集或调整算法参数,进一步扩展为个性化推荐应用,具有较高的实践参考价值。
1. 为什么是Flask+Spark+ALS:毕业设计的选型逻辑
大部分推荐系统毕业设计翻车,不是算法太难,而是数据量稍微上来一点就崩。MovieLens 的 ml-25m 包含 2500 万条评分,用 Pandas 在单机上做矩阵分解,光是读 CSV 就能把内存吃满。这个项目把 Flask 作为用户请求入口,Spark 接管数据清洗与模型训练,ALS 负责矩阵分解,三个层级刚好对应 Web 开发、大数据处理、机器学习三块考核点。它适合想证明自己能把算法工程化的人,而不是只会在 notebook 里调库。整条链路从原始评分表到 Web 页面返回推荐结果,每一步都能复现。
2. MovieLens数据清洗与Spark DataFrame预处理
2.1 原始CSV与Schema设计
MovieLens 常见版本里,ml-latest-small 有 100836 条评分、9742 部电影、610 个用户,正好作为毕设的调试集。解压后核心文件是 ratings.csv 和 movies.csv。ratings.csv 的各列是 userId、movieId、rating、timestamp;movies.csv 是 movieId、title、genres。直接用 Spark 读取时,inferSchema 推断出的 timestamp 是 long,rating 是 double,这两个点在后面都要显式处理。
在 create_db.py 里我一般会单独写一个 build_schema() 函数,把 DataFrame 的列名统一成 user_id、movie_id、rating、rated_at。原因是:ALS 的 userCol、itemCol、ratingCol 需要精确匹配列名;如果原始列名大小写混用,后面调参改参数时很容易漏改。另一个原因是 MovieLens 的 userId 和 movieId 在 CSV 中是一致的整数,但评分可能有空值,应该在读入后立刻过滤,而不是等到训练时报错。顺便提醒,spark 的安装与使用不是零配置:pyspark 对应的 JVM 版本不一致时,启动就会抛 UnsupportedClassVersionError,这类问题在本地调试阶段最常见。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_unixtime spark = SparkSession.builder \ .appName("movie_reco_preprocess") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() ratings = spark.read.load( "data/ml-latest-small/ratings.csv", format="csv", header=True, inferSchema=True ) ratings = ratings.select( col("userId").cast("int").alias("user_id"), col("movieId").cast("int").alias("movie_id"), col("rating").cast("float").alias("rating"), from_unixtime(col("timestamp")).alias("rated_at") ).filter( col("user_id").isNotNull() & col("movie_id").isNotNull() & col("rating").isNotNull() ) ratings.show(5, truncate=False)这段代码把 userId 转成 user_id,用 from_unixtime 把 Unix 时间戳变成可读时间。cast 成 int 和 float 是为了避免 Spark 在 parquet 存储时出现类型不确定的问题。isNotNull 过滤会在后续 ALS.fit 时减少空值导致的脏数据。shuffle.partitions 设置为 4 是因为小数据集的 shuffle 分区不需要 200 个,太多小任务反而让调度时间变长。
2.2 与电影表Join、标签处理与缓存
ratings 表只有 id,最终展示推荐结果至少要带电影标题和年份。movies.csv 里 title 字段形如 “Forrest Gump (1994)”,可以在清洗时拆成年份,虽然 ALS 本身用不到文字特征,但 Web 页面展示时需要。
from pyspark.sql.functions import split movies = spark.read.load( "data/ml-latest-small/movies.csv", format="csv", header=True, inferSchema=True ).select( col("movieId").cast("int").alias("movie_id"), col("title"), col("genres") ) full_df = ratings.join(movies, on="movie_id", how="left") full_df = full_df.withColumn( "genre_list", split(col("genres"), "\\|") ) full_df.cache() print("有效评分总数:", full_df.count())这里的 join 使用 left,因为 MovieLens 里 rating 的 movieId 应该都能在 movies 表中找到,但线上数据不一定。left join 能保留评分记录,万一条电影信息缺失也不影响模型。genre_list 是从管道符分隔的 genres 拆出的数组,为后续做冷启动分析用。cache 在这里要说明一点:full_df 被 count 触发计算后,下一次迭代模型或做筛选时可以直接读内存中的缓存,避免重复解析 CSV。如果集群内存紧张,也可以先 filter 出训练需要的列,再 cache 一个更窄的 DataFrame。
2.3 时间序列划分与随机划分的差别
很多教程直接用 randomSplit 把评分数据分成两份,我建议在毕设里多做一个时间划分。原因是 ALS 评估的核心问题是“对用户未来行为的预测能力”,如果随机划分,同一用户的过去和未来评分可能同时出现在训练集和测试集,评测结果会偏乐观,导师很容易针对这点提问。
# 时间划分:提取 2023-01-01 之后的数据作为测试集 train = full_df.filter("rated_at < '2023-01-01'") test = full_df.filter("rated_at >= '2023-01-01'") print(f"train count: {train.count()}, test count: {test.count()}")如果使用的 ml-latest-small 时间范围是 1995 到 2018,那就需要把阈值调整到数据集 80% 分位处。可以先用full_df.selectExpr("percentile_approx(cast(rated_at as long), 0.8) as ts").collect()拿到阈值,再转成字符串。时间划分后,测试集必然包含训练集的尾部状态,同时也会引入新用户和新电影,这正好测试模型对冷启动的鲁棒性。如果最后 RMSE 偏高,先看是不是 test 里大量 userId 从未出现在 train 中。
| 字段 | 原始类型 | 清洗动作 | 训练中的作用 |
|---|---|---|---|
| userId | int | 重命名 user_id,过滤空 | ALS user 因子输入 |
| movieId | int | 重命名 movie_id,过滤空 | ALS item 因子输入 |
| rating | double | cast 为 float | 矩阵分解的监督值 |
| timestamp | long | from_unixtime 转为时间 | 时间序列切分 |
| genres | string | split 为 array | 冷启动和推荐解释 |
到了这一步,数据已经被规整成“用户-电影-评分-时间”四元组,接下来可以进 ALS 训练。需要注意的是 DataFrame 的分区数不用太大,否则每个分区太小,训练时很多时间浪费在任务调度上。
3. ALS交替最小二乘的训练参数与模型调优
3.1 矩阵分解的两个视角
ALS 全称是交替最小二乘法,核心是把稀疏的评分矩阵 R 拆成两个低维矩阵 U 和 V,分别表示用户特征和电影特征。R 中第 i 行第 j 列的评分约等于 U_i 和 V_j 的内积。这个分解没有解析解,所以采用交替迭代:固定 U,把 V 的每一列当成最小二乘问题求解;固定 V,再解 U。每轮只更新一个矩阵,Spark 可以把每个用户的向量分到不同 executor 上并行计算。
与 SVD 相比,ALS 有两个实际优势:一是能直接处理缺失值,MovieLens 中大量用户只看过几十部电影,矩阵稀疏度极高;二是正则化参数方便控制过拟合,rank 设置 10-20 通常就足够。很多旅游推荐、商品推荐的毕设也沿用这个思路,把评分替换成浏览时长或购买次数后,把 implicitPrefs 设为 true 即可,原理是同一套。
3.2 核心参数与常见取值范围
直接在 pyspark.ml.recommendation 里构造 ALS 时,用户需要关注下面几个参数。如果要做网格搜索,这些参数的取值范围也应该围绕这些区间展开。
| 参数 | 含义 | 我常用的取值 | 调大后的影响 |
|---|---|---|---|
| rank | 隐特征维度 | 10, 15, 20 | 拟合能力增强,但更容易过拟合,训练时间变长 |
| maxIter | 交替迭代轮数 | 10 | 过小欠拟合,过大后面几轮几乎无变化 |
| regParam | 正则化系数 | 0.05 | 太小会放大噪声,太大会把所有预测拉向均值 |
| alpha | 隐反馈置信度权重 | 40(仅 implicitPrefs=true) | 显式评分项目里直接忽略 |
| implicitPrefs | 是否使用隐式反馈 | false | 对于 MovieLens 评分数据必须是 false |
| coldStartStrategy | 冷启动处理 | drop | 不设置会导致 NaN 预测写入结果 |
特别说明 alpha 只在 implicitPrefs=True 时生效。原始 MovieLens 评分是显式行为,0.5 到 5 分,所以用显式模型。如果是隐式反馈项目,比如播放次数、点击次数,通常先把原始计数转换成 confidence = 1 + alpha * count,再传给模型。
3.3 训练、评估与模型落盘
训练用 2.3 节切出的 train 集。对测试集做 transform 后,recommendation 模块会生成 prediction 列。评估器用 RMSE 比较好解释,Map 类的评估还要考虑排序阈值,在毕设中通常作补充指标。
from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator from pyspark.ml.tuning import ParamGridBuilder, CrossValidator als = ALS( userCol="user_id", itemCol="movie_id", ratingCol="rating", rank=15, maxIter=10, regParam=0.05, coldStartStrategy="drop", implicitPrefs=False, seed=42 ) evaluator = RegressionEvaluator( metricName="rmse", labelCol="rating", predictionCol="prediction" ) pred = als.fit(train).transform(test) rmse = evaluator.evaluate(pred) print("test RMSE:", rmse) param_grid = (ParamGridBuilder() .addGrid(als.rank, [10, 15, 20]) .addGrid(als.regParam, [0.01, 0.05, 0.1]) .build()) cv = CrossValidator( estimator=als, estimatorParamMaps=param_grid, evaluator=evaluator, numFolds=3 ) cv_model = cv.fit(train) print("best rank:", cv_model.bestModel.rank, "best reg:", cv_model.bestModel.regParam) cv_model.bestModel.write().overwrite().save("models/als_model")上面的代码先跑一次单模型拿到基准 RMSE,再用交叉验证在 9 个参数组合里搜索。CrossValidator 内部会重复训练,所以在 ml-latest-small 上大约需要几分钟;如果直接把 train 换成 ml-25m,建议先把网格缩小到 rank 加 regParam 的 2×2 矩阵,否则一个 stage 要跑几十轮迭代,Spark UI 上能看到 task 数量成倍增加。write().save() 保存的是完整的 ALSModel 目录,里面包含 itemFactors、userFactors、item 与 user 映射的 parquet 数据,后续 Flask 加载时不需要重新训练。
提示:coldStartStrategy 默认是 nan,如果漏掉 drop,测试集里新用户的预测会变成 null,RMSE 计算直接产生 NaN。出现这种情况先去检查 model.transform(test) 结果里的 prediction 列是否有空值,而不是怀疑算法。模型保存之后,最好再加载一次打印 model.rank,确认写入与读取成功。
4. Flask推荐引擎:从模型到Web接口
4.1 项目结构与入口
项目根目录里的 run.py 是 Web 入口,create_db.py 负责初始化数据库和 Spark 环境,models.py 定义 SQLAlchemy 模型,views.py 放路由。recommend 目录通常放 ALS 初始化和模型加载相关代码,templates 与 static 是 Flask 默认模板和静态文件位置,config.py 存 Spark Session 和数据库连接串。类似的布局在 flask 开发里很常见,把数据库访问和算法引擎分开,防止单文件越写越长。
project/ ├── run.py ├── create_db.py ├── views.py ├── models.py ├── recommend/ │ ├── __init__.py │ └── loader.py ├── templates/ │ ├── base.html │ └── recommend.html ├── static/ └── config.py启动方式是python run.py。run.py 里基本是 create_app(),然后app.run(host="0.0.0.0", port=5000, debug=False)。生产环境不要开 debug,否则 Flask 会启用 reloader,和 SparkContext 的并发创建容易冲突。
4.2 SparkSession与ALS模型的延迟加载
ALS 模型被保存为一个目录,用户每次请求都去 load 一次绝不是好方案,底层的 parquet 文件读取和线性回归都耗时。这个项目里常见做法是用一个模块级单例,只在第一次请求时初始化 SparkSession 和 ALSModel,后面直接复用。注意 SparkSession 不是线程安全的,Flask 开发服务器默认是 single-threaded,但部署到 gunicorn 多 worker 后,要保证每个 worker 进程各持有一个单例。
# recommend/loader.py import threading from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALSModel _lock = threading.Lock() _spark = None _model = None def get_spark(): global _spark if _spark is None: with _lock: if _spark is None: _spark = SparkSession.builder \ .appName("movie_reco_web") \ .config("spark.driver.memory", "2g") \ .config("spark.ui.enabled", "false") \ .getOrCreate() return _spark def get_model(): global _model if _model is None: with _lock: if _model is None: _model = ALSModel.load("models/als_model") return _model双检锁保证了多线程环境下不会重复创建 SparkSession。spark.ui.enabled=false可以省去 Web 端口冲突的麻烦,不过如果要做集群排错,还是建议保留 4040 端口,方便查看 data locality 和 shuffle 数据量。
4.3 生成推荐列表的接口
对于“给某个用户推荐 10 部电影”的需求,ALS 原生提供 recommendForUserSubset。它接受一个只包含 user_id 列的 DataFrame,输出 recommendations 列,该列是由 struct 组成的数组,struct 里有 movie_id 和 rating。这里需要小心的坑:recommendations 的数组顺序是按预测评分从高到低排的,但 Spark 原生数组到 Python 的转换没有类型信息,要用 Row 的字段名逐个取出来。
# views.py from flask import Blueprint, render_template from pyspark.sql.functions import col from recommend.loader import get_model, get_spark bp = Blueprint("recommend", __name__) @bp.route("/recommend/<int:user_id>") def recommend_for_user(user_id): model = get_model() spark = get_spark() user_df = spark.createDataFrame([(user_id,)], ["user_id"]) recs = model.recommendForUserSubset(user_df, 10).collect() if not recs: return render_template("recommend.html", movies=[]) rec_movies = recs[0]["recommendations"] movie_ids = [row["movie_id"] for row in rec_movies] movies_df = spark.read.parquet("data/movies.parquet") \ .filter(col("movie_id").isin(movie_ids)) \ .collect() title_map = {row["movie_id"]: row["title"] for row in movies_df} result = [ {"movie_id": mid, "title": title_map.get(mid, "未知")} for mid in movie_ids ] return render_template("recommend.html", movies=result)collect 在参数维度固定时没有问题,因为推荐列表最多只有 N 条,不会把整个表拉回驱动端。用 isin(movie_ids) 读取电影信息,再按照原始推荐顺序做一次 map,可以保证页面展示的顺序与模型排序一致。很多新人直接返回recs[0]["recommendations"]给 Jinja2 模板,会导致模板里拿到一串结构体难以渲染,先转成字典数组会省去很多模板 Debug 工作。
4.4 用表单提交评分,形成反馈闭环
项目里的 forms.py 一般就是给用户提交评分用的。采用 Flask-WTF 定义 ScoreForm,字段包括 movie_id 和 rating,前端在 recommend.html 遍历 10 部电影,生成下拉选择框。提交后写入 MySQL 或 SQLite 的 ratings 表。常见做法是同时写一张 user_feedback 表,字段包含 user_id、movie_id、rating、created_at,模型训练日任务在夜间读取 feedback 表并重训。ALS 本身不支持在线增量更新,所以“提交即生效”的即时推荐只会出现在演示环节,真正工程化是把新数据累积后定时重训。
到这里推荐链路已经通畅:用户访问页面,接口调用模型,模型读取内存中的 userFactors 与 itemFactors,Spark 输出结果,Flask 转成 JSON 或直接渲染模板。
5. 推荐效果验证与内存排错技巧
5.1 离线指标之外的排序检查
RMSE 只能反映评分预测误差,不代表用户真的觉得推荐“准”。我会额外看每个用户推荐列表里的电影类型覆盖度,以及是否大量推荐续集。用 genre_list 做一次简单的聚合统计,可以看到某类用户被窄化到单一题材。检查代码可以复用 SparkSession,在模型输出后对推荐结果 join films,再按 genre 展开统计。
from pyspark.sql.functions import explode recs_df = model.recommendForAllUsers(10) recs_df = recs_df.withColumn("rec", explode("recommendations")) \ .select("user_id", "rec.movie_id", "rec.rating") recs_df.join(movies, "movie_id") \ .withColumn("genre", explode("genre_list")) \ .groupBy("genre").count() \ .orderBy("count", ascending=False) \ .show()如果发现某类电影占比过高,考虑调小 rank,或者把热门电影过滤掉再训练。另一种做法是给模型推荐结果做后置重排,例如对已经看过的电影直接去掉,再把同系列电影去重。这些规则写在 Flask 接口里就行,不必重训。
5.2 Spark内存与集群提交参数
毕设环境通常是一个虚拟机或单机 docker,容易出现 executor 内存不足。训练前用df.cache()能减少重复 IO,但缓存的还是 RDD 序列化后的对象,堆外内存也要留足。spark-submit 时典型配置如下:
spark-submit \ --master yarn \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 4 \ --conf spark.sql.shuffle.partitions=200 \ train_als.py参数含义:driver-memory 控制 collect 和模型保存时的内存;executor-memory 决定每个 worker 能加载的用户/电影因子块;executor-cores 大于 2 可能让 ALS 阶段的任务洗牌变慢;shuffle.partitions 在集群模式下要回到 200,不要沿用单机调试时的 4。如果做模型评估时 driver OOM,通常是 collect 了过多预测数据,先pred.select("user_id", "prediction").sample(0.1).collect()做抽样,而不是全量拉回。
5.3 模型加载后返回空列表的常见原因
ALSModel.load 成功但接口返回空,多半是 user_id 类型不匹配。保存模型时,userFactors 里 userId 是 int,而 Flask 路由接收的 user_id 默认是字符串,createDataFrame 生成 int 类型时判断不出来。解决方式是新 DataFrame 先 cast:
user_df = spark.createDataFrame([(int(user_id),)], ["user_id"])再有一种情况是推荐列表里全部是 NaN,原因是模型训练时设置了 coldStartStrategy 之外的策略,或者模型路径下 itemFactors 文件缺失。查看模型目录时应看到 itemFactors/part-r-.parquet、userFactors/part-r-.parquet 两个子目录。少一个,直接重新跑训练脚本,不要自己拼接文件。观察 Spark UI 里某个 stage 的 Shuffle Write 是否异常突增,也常能定位到数据倾斜,但 MovieLens 数据分布通常已经比较均匀,反复出现倾斜时优先检查 join 时 hashing key 是否为空。
本文还有配套的精品资源,点击获取