简介这份资源是一套基于Hadoop的电影推荐系统完整实现面向计算机、电子信息工程、数学等专业的大学生可用于课程设计、期末大作业或毕业设计。项目在Windows 10环境下搭建采用Hadoop 2.8.3、Python 3.x、VSCode与MySQL 8.0通过编写代码完成HDFS文件操作与数据处理的实战训练帮助读者掌握分布式平台下的推荐算法落地流程。压缩包共10个文件包含4个Python脚本、2个CSV数据文件以及u.user、u.data、u.item等经典MovieLens数据集文件和一份README说明整体约2.49MB结构紧凑、便于快速上手。代码采用参数化编程参数修改方便思路清晰且注释详细并附有运行结果均经过测试运行成功后才上传。目前已有231人学习适合希望理解Hadoop与推荐系统结合、需要可运行参考实现的学习者。1. 电影推荐系统遇上 Hadoop一套能扛住百万级评分的 Python 工程骨架单机跑推荐算法几万条评分数据用 pandas 还能撑住一旦评分表涨到千万行内存直接爆掉ALS 训练动辄几小时。这个标题讲的不是玩具 demo而是一套用 Python 做算法层、Hadoop 做存储与分布式计算层的电影推荐系统配套源代码和文档说明。它解决的核心问题是把 MovieLens 这类评分数据放进 HDFS用 MapReduce 或 Spark 做离线特征统计与相似度计算再由 Python 侧完成推荐排序和服务输出。适合谁正在做 Hadoop 课程设计的学生、需要给中小型视频站搭离线推荐管线的后端工程师以及想搞懂「推荐系统怎么和 Hadoop 生态拼起来」的开发者。下面按选型、搭建、算法、避坑、调优的顺序拆开讲。2. 为什么是 Python Hadoop 这套组合选型逻辑与数据流拆解2.1 推荐系统为什么绕不开 Hadoop推荐系统的计算量集中在两处一是用户-物品评分矩阵的构建与统计二是相似度矩阵或隐向量的迭代求解。以 MovieLens 25M 数据集为例2500 万条评分、16 万用户、6 万部电影单机用 pandas 读入后光评分表就占约 600MB 内存做物品相似度时中间矩阵会膨胀到几十 GB。Hadoop 的价值在于把这两步拆到多台机器上HDFS 负责存原始评分和中间结果MapReduce 或 Spark 负责分布式聚合。常见做法是用 MapReduce 做「按电影分组统计评分人数和均分」这类可并行任务用 Spark MLlib 的 ALS 做矩阵分解。Python 在这里不是替代 Hadoop而是作为调度层和算法胶水层——用 PySpark 提交作业用 pandas 做小规模结果的后处理用 Flask 或 FastAPI 暴露推荐接口。提示如果数据量在百万条以下单机 scikit-learn 的 NMF 或 implicit 库足够不必上 Hadoop。上 Hadoop 的门槛是数据量持续增长且单机训练时间超过可接受范围。2.2 整体数据流从 HDFS 到推荐结果一条完整的离线推荐链路分四段原始评分数据ratings.csv、movies.csv通过hdfs dfs -put上传到 HDFS 的/movie/raw/目录。用 MapReduce 或 Spark SQL 做清洗与统计输出到/movie/clean/包括每部电影的评分人数、均分、评分分布。用 Spark MLlib 的 ALS 在清洗后的数据上训练隐向量模型模型文件存到/movie/model/。Python 侧读取模型和统计结果对每个用户生成 Top-N 推荐列表写入 HDFS 或 MySQL供接口层查询。这个分层的意义在于原始数据、中间统计、模型、服务结果各自独立存储任何一段出问题都能单独重跑不用全链路重来。2.3 最小可跑环境Hadoop 伪分布式 PySpark 安装先确认 Java 和 SSH 可用这是 Hadoop 的硬依赖。下面是在 Ubuntu 22.04 上的最小步骤# 安装 Java 8Hadoop 3.x 对 Java 8 兼容最稳 sudo apt update sudo apt install openjdk-8-jdk -y java -version # 配置 SSH 免密伪分布式必须 ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys ssh localhost逻辑说明Hadoop 的 NameNode 和 DataNode 之间通过 SSH 启动免密没配好会出现Permission denied导致 DataNode 起不来。java -version要确认输出是 1.8.x如果是 11 或 17部分 Hadoop 3.2 以下版本会报UnsupportedClassVersionError。接着下载 Hadoop 并配置核心文件# 解压到 /usr/local sudo tar -xzf hadoop-3.3.6.tar.gz -C /usr/local/ sudo mv /usr/local/hadoop-3.3.6 /usr/local/hadoop # 配置环境变量写入 ~/.bashrc export HADOOP_HOME/usr/local/hadoop export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64参数说明HADOOP_HOME是后续所有脚本的基准路径JAVA_HOME必须指向 JDK 而非 JRE否则hadoop namenode -format会失败。改完source ~/.bashrc生效。然后编辑core-site.xml和hdfs-site.xml!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration !-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property /configuration伪分布式下副本数必须设为 1设成 3 会因为只有一个 DataNode 而一直处于副本不足状态。格式化并启动hdfs namenode -format start-dfs.sh jps # 应看到 NameNode、DataNode、SecondaryNameNodejps是排查 Hadoop 启动问题的第一命令缺哪个进程就去看对应日志日志在$HADOOP_HOME/logs/。2.4 PySpark 接入与数据上传pip install pyspark3.5.0 pandas numpy hdfs dfs -mkdir -p /movie/raw hdfs dfs -put ratings.csv /movie/raw/ hdfs dfs -ls /movie/raw/pyspark版本要和 Hadoop 版本匹配3.5.x 对应 Hadoop 3.3。上传后用hdfs dfs -ls确认文件存在文件大小和本地一致再继续。3. 用 PySpark 实现 ALS 推荐从评分表到 Top-N 列表3.1 数据清洗与评分统计原始 ratings.csv 通常有 userId、movieId、rating、timestamp 四列。第一步是去掉评分次数过少的用户和电影否则 ALS 会为这些冷门对象生成噪声向量。from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg spark SparkSession.builder \ .appName(MovieRecommend) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 8) \ .getOrCreate() # 读取 HDFS 上的评分数据 ratings spark.read.csv(hdfs://localhost:9000/movie/raw/ratings.csv, headerTrue, inferSchemaTrue) # 统计每个用户和每部电影的评分次数 user_counts ratings.groupBy(userId).agg(count(rating).alias(cnt)) movie_counts ratings.groupBy(movieId).agg(count(rating).alias(cnt)) # 过滤用户至少评 20 部电影至少被评 50 次 valid_users user_counts.filter(col(cnt) 20).select(userId) valid_movies movie_counts.filter(col(cnt) 50).select(movieId) clean ratings.join(valid_users, userId).join(valid_movies, movieId) clean.cache() print(清洗后评分条数, clean.count())逻辑说明spark.sql.shuffle.partitions默认 200本地伪分布式下会启动 200 个任务拖慢速度设成 CPU 核数的 2 到 4 倍即可。cache()把清洗结果留在内存后续 ALS 训练会多次读取。阈值 20 和 50 是经验值数据稀疏时可降到 10 和 30但太低会让模型学到不可靠的向量。3.2 ALS 模型训练与参数设置ALS交替最小二乘把用户-物品评分矩阵分解成两个低秩矩阵通过交替固定一方求解另一方来逼近原始评分。from pyspark.ml.recommendation import ALS from pyspark.ml.evaluation import RegressionEvaluator # 划分训练集和测试集 train, test clean.randomSplit([0.8, 0.2], seed42) als ALS( userColuserId, itemColmovieId, ratingColrating, rank50, # 隐向量维度 maxIter10, # 迭代次数 regParam0.1, # 正则化系数 coldStartStrategydrop, # 丢弃测试集中冷启动样本 nonnegativeTrue # 评分非负约束向量非负 ) model als.fit(train) # 在测试集上评估 predictions model.transform(test) evaluator RegressionEvaluator(metricNamermse, labelColrating, predictionColprediction) rmse evaluator.evaluate(predictions) print(RMSE , rmse)参数说明rank控制隐向量维度50 是 MovieLens 上的常用起点调到 100 可能提升精度但训练时间翻倍regParam防过拟合0.1 偏保守RMSE 偏高时可试 0.05maxIter10 次通常够收敛观察 RMSE 不再下降即可停coldStartStrategydrop必须加否则测试集中新用户会产生 NaN 预测值导致 RMSE 为 NaN。RMSE 在 0.85 到 0.95 之间属于正常范围低于 0.8 要警惕数据泄漏。3.3 生成 Top-N 推荐并落库# 为每个用户生成 10 部推荐电影 user_recs model.recommendForAllUsers(10) # 展开成 (userId, movieId, rating) 三列 from pyspark.sql.functions import explode flat_recs user_recs.select( userId, explode(recommendations).alias(rec) ).select( userId, col(rec.movieId).alias(movieId), col(rec.rating).alias(score) ) # 关联电影名输出到 HDFS movies spark.read.csv(hdfs://localhost:9000/movie/raw/movies.csv, headerTrue, inferSchemaTrue) result flat_recs.join(movies, movieId).select( userId, movieId, title, score ) result.write.mode(overwrite).parquet(hdfs://localhost:9000/movie/output/recs)逻辑说明recommendForAllUsers(10)返回的是数组列必须用explode展开才能和电影表 join。输出用 parquet 而非 csv因为 parquet 带 schema 且压缩率高后续 Python 侧用 pandas 读取时不用再推断类型。mode(overwrite)保证重跑不报目录已存在。Python 侧读取结果做接口import pandas as pd recs pd.read_parquet(hdfs://localhost:9000/movie/output/recs) def get_user_recs(user_id, top_n10): sub recs[recs[userId] user_id].nlargest(top_n, score) return sub[[title, score]].to_dict(records)这段代码把 HDFS 上的 parquet 读进 pandas按用户过滤后返回 Top-N。生产环境应把结果同步到 MySQL 或 Redis避免每次请求都读 HDFS。4. 避坑指南Hadoop 伪分布式 PySpark 推荐系统的 5 个血泪翻车点4.1 DataNode 启动后立刻消失现象start-dfs.sh后jps只看到 NameNode没有 DataNode。原因多次执行hdfs namenode -format导致 NameNode 的 clusterID 和 DataNode 的 clusterID 不一致。解决停掉所有进程删除dfs.namenode.name.dir和dfs.datanode.data.dir指向的目录重新格式化一次之后不要再重复格式化。4.2 PySpark 报 Python worker exited unexpectedly现象提交 ALS 作业时抛Python worker exited unexpectedly或Connection refused。原因PySpark 的 Python 版本和集群节点不一致或PYSPARK_PYTHON未设置。解决在spark-env.sh里加export PYSPARK_PYTHON/usr/bin/python3并确认所有节点的 Python 路径一致。伪分布式下只有一个节点但仍需显式指定。4.3 ALS 预测结果全是 NaN现象evaluator.evaluate(predictions)返回 NaN。原因测试集中存在训练集里没出现过的用户或电影ALS 无法为其生成向量。解决设置coldStartStrategydrop或在划分数据集前先做一次全局过滤保证测试集的用户和电影都在训练集中出现过。4.4 内存不足导致 shuffle 失败现象作业跑到某个 stage 报Container killed by YARN for exceeding memory limits或java.lang.OutOfMemoryError。原因spark.sql.shuffle.partitions太小导致单分区数据量过大或spark.driver.memory不够。解决把 shuffle partitions 调到 16 到 32driver 内存设 2g 以上executor 内存按机器实际内存的 70% 设置。4.5 中文电影名乱码现象输出结果里电影标题显示为?????。原因原始 csv 是 UTF-8但 Spark 读取时默认编码或 HDFS 输出编码不一致。解决读取时显式指定encodingUTF-8输出 parquet 不受影响但如果导出 csv 给前端要在 Python 侧用encodingutf-8-sig写文件。5. 进阶调优与验证让推荐结果从「能跑」到「可信」5.1 用交叉验证替代单次划分单次 8:2 划分的 RMSE 波动可能达到 0.05不足以判断参数好坏。用 Spark 的CrossValidator做 3 折交叉验证同时搜索 rank 和 regParamfrom pyspark.ml.tuning import ParamGridBuilder, CrossValidator param_grid ParamGridBuilder() \ .addGrid(als.rank, [20, 50, 100]) \ .addGrid(als.regParam, [0.05, 0.1, 0.2]) \ .build() cv CrossValidator(estimatorals, estimatorParamMapsparam_grid, evaluatorevaluator, numFolds3) cv_model cv.fit(clean) best_rank cv_model.bestModel.rank print(最佳 rank, best_rank)逻辑说明9 组参数乘 3 折等于 27 次训练本地伪分布式下可能跑 1 到 2 小时。建议先用小数据集比如 10% 采样筛出候选参数再用全量数据验证。bestModel.rank能直接读出最优维度省去手动比对。5.2 离线指标之外加一层业务校验RMSE 低不代表推荐好看。加两个业务指标一是覆盖率即推荐列表里出现的不同电影数占总电影数的比例太低说明推荐集中在头部二是新颖度用推荐电影的平均流行度倒数衡量越高说明越能推冷门好片。# 覆盖率 rec_movies flat_recs.select(movieId).distinct().count() total_movies movies.count() coverage rec_movies / total_movies print(覆盖率, coverage) # 新颖度推荐电影的平均评分人数倒数 movie_pop clean.groupBy(movieId).agg(count(rating).alias(pop)) rec_with_pop flat_recs.join(movie_pop, movieId) novelty rec_with_pop.select(avg(1.0 / col(pop))).collect()[0][0] print(新颖度, novelty)覆盖率低于 0.3 说明推荐同质化严重可以调大 rank 或对热门电影做降权。新颖度没有绝对标准但同一批参数下对比才有意义。5.3 一个我常犯的错误早期我总盯着 RMSE 调参把 rank 从 50 加到 200RMSE 降了 0.02但训练时间从 8 分钟涨到 40 分钟线上接口的推荐结果反而更集中在几部热门片上。后来我养成习惯每次调参同时记录 RMSE、覆盖率、训练时长三个数只有三者都在可接受范围才采纳。推荐系统不是精度竞赛是精度、多样性、成本三者的平衡。希望帮到你。本文还有配套的精品资源点击获取