简介一份面向毕业设计场景的电影智能推荐系统完整实现整合了Flask Web框架、Spark分布式计算、ALS协同过滤算法与MovieLens公开评分数据集。项目覆盖数据清洗、特征处理、推荐模型训练与网页交互展示等关键环节适合计算机相关专业的学生用于毕业设计参考也适合推荐系统初学者跟随项目代码理解完整实现流程。压缩包共50个文件包含10个Python源码文件、10个HTML页面模板、4个CSV评分数据集以及图片、配置和说明文档整体仅6.59MB结构清晰便于本地部署与学习。目前已有66人学习下载。资源中提供了可运行的Flask服务、Spark处理脚本、ALS训练代码和项目说明能够帮助读者快速搭建一个完整的电影推荐示例并在现有代码基础上替换数据集或调整算法参数进一步扩展为个性化推荐应用具有较高的实践参考价值。1. 为什么是FlaskSparkALS毕业设计的选型逻辑大部分推荐系统毕业设计翻车不是算法太难而是数据量稍微上来一点就崩。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、timestampmovies.csv 是 movieId、title、genres。直接用 Spark 读取时inferSchema 推断出的 timestamp 是 longrating 是 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, formatcsv, headerTrue, inferSchemaTrue ) 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, truncateFalse)这段代码把 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, formatcsv, headerTrue, inferSchemaTrue ).select( col(movieId).cast(int).alias(movie_id), col(title), col(genres) ) full_df ratings.join(movies, onmovie_id, howleft) 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(ftrain 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 中。字段原始类型清洗动作训练中的作用userIdint重命名 user_id过滤空ALS user 因子输入movieIdint重命名 movie_id过滤空ALS item 因子输入ratingdoublecast 为 float矩阵分解的监督值timestamplongfrom_unixtime 转为时间时间序列切分genresstringsplit 为 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仅 implicitPrefstrue显式评分项目里直接忽略implicitPrefs是否使用隐式反馈false对于 MovieLens 评分数据必须是 falsecoldStartStrategy冷启动处理drop不设置会导致 NaN 预测写入结果特别说明 alpha 只在 implicitPrefsTrue 时生效。原始 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( userColuser_id, itemColmovie_id, ratingColrating, rank15, maxIter10, regParam0.05, coldStartStrategydrop, implicitPrefsFalse, seed42 ) evaluator RegressionEvaluator( metricNamermse, labelColrating, predictionColprediction ) 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( estimatorals, estimatorParamMapsparam_grid, evaluatorevaluator, numFolds3 ) 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测试集里新用户的预测会变成 nullRMSE 计算直接产生 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(host0.0.0.0, port5000, debugFalse)。生产环境不要开 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.enabledfalse可以省去 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, moviesresult)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 与 itemFactorsSpark 输出结果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, ascendingFalse) \ .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.partitions200 \ 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 先 castuser_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 是否为空。本文还有配套的精品资源点击获取