拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Python+Hadoop伪分布式协同过滤电影推荐系统实操

Python+Hadoop伪分布式协同过滤电影推荐系统实操

简介:本资源是一套面向计算机专业本科生的毕业设计实战项目,聚焦大数据环境下的个性化电影推荐系统实现,适用于需完成毕设、夯实Python与Hadoop协同开发能力的学习者。项目以Python为核心语言实现推荐算法(如协同过滤),依托Hadoop分布式框架处理海量用户评分与电影元数据,解决传统单机环境难以应对的大规模稀疏矩阵计算问题。压缩包共10个文件,含4个核心Python脚本(mr1.py/mr2.py/run.py/mrjobTemp.py)、2个CSV数据集(ratings.csv/result.csv)、以及u.data/u.item/u.user三类标准MovieLens结构化数据文件,另含README.md说明文档;整体2.49MB,轻量易部署,目录简洁,突出MapReduce任务划分与结果验证逻辑。目前已有356人学习下载,读者可直接复现从数据预处理、Hadoop作业提交、模型训练到推荐结果生成的完整链路,并参考代码组织方式与HDFS数据交互实践,快速掌握大数据推荐系统的工程落地要点。

1. 这不是个“跑通就完事”的毕设:Python+Hadoop电影推荐系统,真正在伪分布式环境里跑出协同过滤结果的实操闭环

你下载这个毕业设计:基于python+Hadoop的电影推荐系统.zip,不是为了在 PyCharm 里双击 run.py 看个控制台输出“recommendation done”——那叫 Python 小练习。它真正价值在于:用真实 MovieLens 数据(u.data、u.item、u.user),在本地 Hadoop 伪分布式集群上,通过 mrjob 封装 MapReduce,完整走通“用户-物品评分矩阵 → 基于用户的协同过滤 → Top-N 推荐生成 → 结果落盘为 result.csv”的端到端数据流水线。整个流程绕不开 HDFS 文件路径配置、mrjob 作业提交方式、评分稀疏性导致的 KeyError、以及 Python 端对 Hadoop Streaming 协议的隐式适配。适合正在赶毕设 deadline、但又不想交一个“本地单机版推荐算法”的计算机/软件工程本科生;也适合想快速验证 Hadoop 生态下推荐逻辑落地可行性的转行者——它不教你 ZooKeeper 怎么装,但会告诉你hdfs dfs -put u.data /input/后,mrjobTemp.py里哪一行路径写错会导致 job 直接卡在 ACCEPTED 状态。

这不是一个“理论正确但跑不起来”的教学 Demo。它包含可执行的mr1.py(计算用户相似度)、mr2.py(生成邻居预测评分)、run.py(串联调度),还有ratings.csv和result.csv这两个关键中间与终态文件——前者是清洗后的结构化评分表(user_id,item_id,rating,timestamp),后者是你最终能拿去答辩展示的“用户 196 最可能喜欢的 5 部电影 ID 及预测分”。所有代码都基于 Python 3.7–3.9 兼容写法,没用 SparkSQL 或 Dask 这类高阶封装,直面 MapReduce 的 shuffle 语义和序列化约束。如果你的毕设开题写了“采用 Hadoop 分布式框架处理百万级评分数据”,而实际只用了 Pandas 读 CSV 做 SVD,那这份资源就是你补上技术栈缺口的最后一块拼图。


2. 从 MovieLens 到 HDFS:数据准备与 Hadoop 伪分布式环境校验

2.1 MovieLens 数据结构解析与本地清洗必要性

项目自带的u.data、u.item、u.user是经典的 MovieLens 100K 数据集(1998 年采集,10 万条评分记录)。但直接扔进 Hadoop 会翻车——u.data是 tab 分隔的四元组(user_id, item_id, rating, timestamp),但u.item的字段数不固定(含电影标题、年份、类型等),且含大量括号与斜杠,HDFS 默认 TextInputFormat 会因换行符或特殊字符截断记录。必须先做轻量清洗:

# 在项目根目录执行(Linux/macOS)或 Git Bash(Windows) sed 's/[^[:print:]]//g' u.data | sed 's/[[:space:]]\+/\t/g' | awk -F'\t' '$1!="" && $2!="" && $3>=1 && $3<=5 {print $1,$2,$3,$4}' OFS='\t' > ratings.csv

提示:这行命令干三件事:① 删除不可见控制字符(MovieLens 原始文件含 DOS 换行符 \r);② 统一空格/制表符为单个\t;③ 过滤掉 user_id 或 item_id 为空、评分不在 1–5 区间的脏数据。最终ratings.csv是严格四列纯数字文本,无 header,这是 mrjob 要求的输入格式。

u.user和u.item不参与 MapReduce 计算,仅作后续结果解释用。u.user中的 age、occupation 字段在协同过滤中未使用,但答辩时可说明“预留用户画像扩展接口”。

2.2 Hadoop 伪分布式环境最低可行验证清单

别急着改core-site.xml——先确认你的 Hadoop 已处于可提交作业状态。本项目依赖hadoop-client和hadoop-common,而非全集群。验证步骤必须全部通过:

步骤命令预期输出关键检查点
1. Java 与 Hadoop 版本兼容java -version
hadoop version
Java ≥ 1.8
Hadoop ≥ 3.2.0
mrjob 6.x 与 Hadoop 3.x 的 RPC 协议不兼容旧版
2. HDFS 是否可读写hdfs dfs -mkdir -p /input
hdfs dfs -put ratings.csv /input/
hdfs dfs -ls /input/
显示ratings.csv文件大小若报Connection refused,说明 NameNode 未启动
3. YARN ResourceManager 是否就绪yarn node -list至少显示Total Nodes:1且状态为RUNNINGACCEPTED状态卡住的根源常在此

注意:Windows 用户若用 WSL2,务必关闭 Windows 防火墙;Mac 用户若用 Homebrew 安装 Hadoop,需手动设置HADOOP_HOME并将$HADOOP_HOME/bin加入 PATH。hadoop-env.sh中JAVA_HOME必须指向 JDK 路径(非 JRE),否则yarn进程启动失败。

2.3 mrjob 配置文件.mrjob的核心参数绑定

项目未提供.mrjob配置文件,但mrjobTemp.py依赖它指定 Hadoop 运行模式。必须在项目根目录创建该文件:

# .mrjob runners: hadoop: hadoop_streaming_jar: /opt/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.6.jar hadoop_bin: /opt/hadoop/bin/hadoop yarn_bin: /opt/hadoop/bin/yarn hdfs_home: hdfs://localhost:9000 python_archives: [] setup_cmds: - pip install numpy==1.21.6

参数说明:

  • hadoop_streaming_jar:Hadoop 3.x 的 streaming jar 路径(版本号需与你安装的 Hadoop 一致,常见位置/share/hadoop/tools/lib/);
  • hdfs_home:必须是hdfs://协议,不能是file://,否则 mrjob 会尝试本地模式而非 YARN 提交;
  • setup_cmds:指定 worker 节点需预装的 Python 包(本项目仅需 numpy,避免在 mapper/reducer 中 import 失败);
  • python_archives:留空即可,本项目无自定义模块打包需求。

3. 协同过滤的 MapReduce 实现:mr1.py 与 mr2.py 的数据流拆解

3.1 mr1.py:基于用户的相似度计算(User-Based CF 第一阶段)

mr1.py实现的是“用户两两共评电影数 + 余弦相似度分子”计算。其 Map 阶段输出格式决定了 Reduce 阶段能否聚合:

# mr1.py 关键 map 方法(简化版) def mapper(self, _, line): user_id, item_id, rating, _ = line.strip().split('\t') # 输出:(user_id, item_id) -> rating,用于后续 join yield (user_id, item_id), float(rating) # 同时输出:(item_id, user_id) -> rating,构建物品-用户倒排索引 yield (item_id, user_id), float(rating)

逻辑说明:此设计是协同过滤 MapReduce 的经典 trick——同一行输入产生两条 KV 对,第一条(user_id, item_id)用于后续按用户分组;第二条(item_id, user_id)用于找出“哪些用户共同评价过同一部电影”。Reduce 阶段收到(item_id, user_id)的所有 rating 后,就能统计共评用户对(如 user1 和 user2 都评了 item123,则计数 +1)。

mr1.py的 Reduce 阶段不直接算相似度,只输出(user1,user2)→共评电影数。因为余弦相似度分母(用户各自评分向量模长)需全局统计,放在mr2.py中统一计算更高效。

3.2 mr2.py:邻居预测与 Top-N 生成(User-Based CF 第二阶段)

mr2.py接收mr1.py的输出,并关联原始评分数据(ratings.csv):

# mr2.py 中的关键 reduce 方法(片段) def reducer(self, user_pair, values): # values 是 [common_count, user1_rating, user2_rating, ...] common_items = [] for v in values: if isinstance(v, tuple): # 来自 mr1.py 的共评数 common_count = v[0] else: # 来自 ratings.csv 的原始评分 user_id, item_id, rating = v.split(',') if user_id == user_pair[0]: common_items.append((item_id, float(rating))) # 对 user_pair[0] 的每个未评电影,用 user_pair[1] 的评分加权预测 for item_id, pred_rating in self._predict(user_pair, common_items): yield user_pair[0], (item_id, pred_rating)

参数说明:

  • user_pair是元组(user1, user2),表示候选邻居;
  • self._predict()内部实现标准协同过滤公式:
    $$\hat{r}{ui} = \bar{r}u + \frac{\sum{v \in N(u)} sim(u,v) \cdot (r{vi} - \bar{r}v)}{\sum{v \in N(u)} |sim(u,v)|}$$
    其中N(u)是 user u 的 Top-K 相似用户,sim(u,v)来自mr1.py输出;

  • 最终yield的(user_id, (item_id, pred_rating))会被 mrjob 自动按 user_id 分组,供run.py汇总。

3.3 run.py:作业调度与结果合并控制流

run.py不是简单顺序执行,而是构建 DAG 依赖:

# run.py 核心逻辑 if __name__ == '__main__': # Step 1: 执行 mr1.py,输出到 /tmp/mr1_output mr1_job = MRUserSimilarity(args=['-r', 'hadoop', '--output-dir', '/tmp/mr1_output']) with mr1_job.make_runner() as runner: runner.run() # Step 2: 将 mr1 输出与 ratings.csv 合并,作为 mr2 输入 hdfs_cmd = f"hdfs dfs -cat /tmp/mr1_output/part-* > /tmp/mr1_merged" subprocess.run(hdfs_cmd, shell=True) # Step 3: 执行 mr2.py,指定输入为合并后数据 mr2_job = MRRecommendation(args=['-r', 'hadoop', '--input', '/tmp/mr1_merged', '--output-dir', '/output/result']) with mr2_job.make_runner() as runner: runner.run()

关键点:mr2.py的输入不是 HDFS 路径,而是本地临时文件/tmp/mr1_merged。这是因为 mrjob 的--input参数不支持 HDFS 路径直接读取(需用hdfs dfs -cat导出)。/output/result是 HDFS 路径,result.csv将在此目录下生成。


4. 避坑指南:Hadoop 伪分布式下协同过滤的五个血泪经验

4.1 现象:mrjob 提交后 YARN Web UI 显示ACCEPTED状态长期不变成RUNNING

原因:ResourceManager 未分配 Container,常见于yarn.scheduler.capacity.root.default.maximum-capacity设置过低(默认 100,但若集群内存不足仍会拒绝)。
解决:编辑$HADOOP_HOME/etc/hadoop/capacity-scheduler.xml,将该值设为100,并确保yarn.nodemanager.resource.memory-mb≥ 4096(伪分布式至少需 4GB 内存)。

4.2 现象:mr1.pyReduce 阶段报KeyError: 'user1'

原因:ratings.csv中 user_id 为字符串(如"196"),但mr1.py的 mapper 解析时未 strip 引号,导致(user_id, item_id)键含多余空格或引号。
解决:在mapper中强制清理:user_id = line.split('\t')[0].strip('"\' ')。

4.3 现象:result.csv为空,HDFS 中/output/result/part-00000文件大小为 0

原因:mr2.py的 reducer 未触发——因为mr1.py输出的 key 格式与mr2.py期望的user_pair不匹配。mr1.py输出(user1,user2)是字符串拼接(如"196,234"),而mr2.py试图用tuple(key.split(','))解析,但若mr1.py输出含空格则失败。
解决:统一mr1.py的 key 输出为f"{u1},{u2}",mr2.py中用key.split(',')后strip()每个元素。

4.4 现象:run.py执行时报subprocess.CalledProcessError,提示hdfs dfs -cat: No such file or directory

原因:/tmp/mr1_output目录在 HDFS 中不存在,或part-*文件名不匹配(Hadoop 3.x 默认输出为part-r-00000,非part-00000)。
解决:将hdfs_cmd改为hdfs dfs -cat /tmp/mr1_output/part-r-* > /tmp/mr1_merged。

4.5 现象:result.csv中出现重复 user_id 行,且预测分异常高(如 12.5)

原因:协同过滤公式中未对sim(u,v)归一化,当某用户与多个邻居相似度极高时,分子爆炸。mr2.py的_predict方法缺少sim截断(如sim = max(-1.0, min(1.0, sim)))。
解决:在_predict中添加相似度钳位:sim = max(-0.99, min(0.99, sim)),避免除零和数值溢出。


5. 结果验证与答辩级可视化:从 result.csv 到可演示的推荐看板

5.1 result.csv 结构解析与可信度校验

result.csv是 HDFS 输出的文本文件,需先导出本地:

hdfs dfs -get /output/result/part-r-00000 result_local.csv

其格式为:user_id<TAB>item_id<TAB>predicted_rating。验证三要素:

  1. 覆盖率:统计user_id去重数 ÷ 总用户数(u.user行数),应 ≥ 85%(MovieLens 100K 有 943 用户);
  2. 合理性:predicted_rating应在 0.5–5.0 区间,超界值占比 < 0.1%;
  3. 多样性:对任一 user_id,其 top-5item_id应覆盖不同电影类型(查u.item第 5–24 列的类型 bit 位)。
# quick_validate.py import pandas as pd df = pd.read_csv('result_local.csv', sep='\t', names=['user','item','pred']) print(f"覆盖用户数: {df['user'].nunique()}/943") print(f"预测分范围: [{df['pred'].min():.2f}, {df['pred'].max():.2f}]") print(f"Top-10 用户平均推荐数: {df.groupby('user').size().nlargest(10).mean():.1f}")

5.2 构建答辩演示页:用 Flask 快速搭推荐看板

无需 React/Vue,50 行 Flask 足够:

# demo_app.py from flask import Flask, render_template, request import pandas as pd app = Flask(__name__) df = pd.read_csv('result_local.csv', sep='\t', names=['user','item','pred']) items = pd.read_csv('u.item', sep='|', encoding='ISO-8859-1', usecols=[0,1], names=['item_id','title']) @app.route('/') def index(): return render_template('index.html', users=df['user'].unique()[:20]) @app.route('/recommend') def recommend(): uid = int(request.args.get('user')) recs = df[df['user']==uid].sort_values('pred', ascending=False).head(5) recs = recs.merge(items, left_on='item', right_on='item_id') return render_template('result.html', user=uid, recommendations=recs.to_dict('records')) if __name__ == '__main__': app.run(debug=True)

配套templates/index.html用<select>下拉选用户,result.html用<ul>展示电影标题+预测分。启动后访问http://localhost:5000,输入196(MovieLens 标准测试用户),即可看到“Star Wars (1977)”、“Contact (1997)”等高分推荐——这才是答辩时能点击演示的“活系统”。

5.3 毕设文档关键页:如何把技术细节转化成论文图表

答辩 PPT 中避免贴代码,用三张图讲清技术价值:

图表类型内容要点制作工具
数据流图u.data→ratings.csv→mr1.py→mr1_output→mr2.py→result.csv,标注各环节耗时(用time命令实测)draw.io 或 PowerPoint SmartArt
相似度热力图取前 50 用户,计算两两相似度矩阵,用 seaborn heatmap 展示(代码见plot_similarity.py)Python + seaborn
推荐质量对比表本系统 vs Pandas 单机版 vs 随机推荐,在 RMSE、Coverage、Novelty 三指标对比Excel 或 Markdown 表格

我的习惯:每次答辩前,我会用hadoop fs -du -s /tmp/*清理所有临时目录,并重新跑一遍run.py,确保result.csv时间戳最新。从那以后我每次提交毕设代码,都强制走一遍hdfs dfs -ls /output/+hdfs dfs -cat /output/result/part-r-* | head -n 5,确认输出真实存在——这比任何文档描述都有说服力。希望帮到你。

本文还有配套的精品资源,点击获取

返回列表