简介这是一套面向计算机专业本科生的毕业设计与课程大作业实战资源基于Hadoop生态实现电影推荐系统采用Python语言开发兼顾算法逻辑与工程部署适合零基础入门者快速上手分布式推荐实践。压缩包共10个文件4个核心Python脚本含详细注释、2个CSV格式数据集、1个README.md文档、以及u.user/u.item/u.data等标准MovieLens结构化数据文件总大小仅2.49MB轻量易部署涵盖数据预处理、MapReduce任务编写、协同过滤算法实现及结果分析全流程。已有356人学习下载资源结构清晰、模块职责明确——mr1.py与mr2.py实现分步MapReduce计算run.py封装主流程result.csv直观呈现推荐结果配套文档说明部署步骤与运行验证方法。读者可直接复现完整推荐链路掌握HadoopPython协同开发范式积累分布式系统调试与推荐算法落地经验。1. 这不是纯 Python 推荐系统而是用 Python 写 MapReduce 逻辑、跑在 Hadoop 上的真实批处理链路很多同学拿到“Python 实现的电影推荐系统”压缩包第一反应是装个scikit-learn或surprise库读 CSV、调模型、出结果——但这个项目完全不是这条路。它本质是一套面向 Hadoop 生态的批处理推荐流水线用户行为日志u.data和电影元数据u.item作为输入通过mrjob框架将 Python 编写的 Mapper/Reducer 逻辑提交到本地伪分布式或远程 Hadoop 集群执行最终生成基于物品协同过滤Item-Based CF的相似度矩阵与 Top-N 推荐结果result.csv。它不依赖 Spark MLlib也不用 PySpark DataFrame API而是回归 MapReduce 原语——这对理解推荐系统底层数据流、调试分布式计算瓶颈、应对课程设计答辩中“你为什么不用 Spark”的追问有不可替代的价值。适合需要展示“大数据平台集成能力”而非仅“算法调包能力”的本科毕设、大数据课程设计或期末大作业场景尤其当老师明确要求“必须运行在 Hadoop 环境下”时这套代码比纯 Python 版本更具说服力。2. 为什么选 mrjob 而非原生 Java 或 PySpark轻量级 Python-Hadoop 桥接的工程权衡2.1 mrjob 的定位与不可替代性在 Hadoop 生态中Python 开发者面临三类主流选择原生 Java MapReduce性能高、控制粒度细但开发成本陡增对 Python 背景学生极不友好PySparkAPI 抽象层高DataFrame 操作简洁但需完整 Spark 环境且默认不直接兼容 Hadoop YARN 的资源调度细节mrjob核心价值在于零配置桥接 Python 逻辑与 Hadoop Streaming。它将 Python 脚本自动打包为符合 Hadoop Streaming 协议的可执行文件即标准输入/输出流式处理无需编译.jar不依赖 Spark 集群甚至能在单机伪分布式 Hadoop 上直接验证逻辑正确性。本项目中mr1.py和mr2.py就是典型的两阶段 mrjob 作业前者计算用户-电影评分共现频次后者基于共现矩阵计算物品相似度。这种分阶段设计正是传统协同过滤中“先统计、再计算”的经典范式落地。提示mrjob 并非生产级大数据框架但它是教学场景下的黄金折中——既避开 Java 的语法门槛又绕过 Spark 的环境复杂度让学生把精力聚焦在“数据如何被切分、聚合、传递”这一本质问题上。2.2 项目核心文件职责解耦与执行顺序从main文件夹结构可还原完整数据流文件名类型核心职责关键依赖u.data原始数据用户 ID、电影 ID、评分、时间戳Tab 分隔必须存在Hadoop 输入源u.item元数据电影 ID、标题、类型竖线 分隔run.py控制脚本封装mr1.py→mr2.py串行执行逻辑含 Hadoop 参数注入mrjob,subprocessmr1.pyMapReduce 作业1Mapper 输出(movie_id, user_id)Reducer 统计每对电影被同一用户评分的次数共现矩阵mrjob.job.MRJobmr2.pyMapReduce 作业2Mapper 解析mr1输出的共现对Reducer 计算余弦相似度并归一化同上需读取mr1输出目录result.csv最终输出格式为movie_id1,movie_id2,similarity_score按相似度降序排列由mr2.py的--output-dir指定该流程严格遵循 Hadoop Streaming 的“输入→Map→Shuffle→Reduce→输出”五阶段模型。run.py中的关键命令如下# 执行第一阶段生成共现对 python mr1.py -r hadoop \ --hadoop-bin /usr/local/hadoop/bin/hadoop \ --hadoop-streaming-jar /usr/local/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.6.jar \ --output-dir hdfs://localhost:9000/user/output/mr1 \ hdfs://localhost:9000/user/input/u.data # 执行第二阶段基于共现对计算相似度 python mr2.py -r hadoop \ --hadoop-bin /usr/local/hadoop/bin/hadoop \ --hadoop-streaming-jar /usr/local/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.6.jar \ --output-dir hdfs://localhost:9000/user/output/mr2 \ hdfs://localhost:9000/user/output/mr1/part-00000注意--hadoop-bin必须指向你本地 Hadoop 安装路径下的hadoop可执行文件--hadoop-streaming-jar的版本号如3.3.6需与你的 Hadoop 版本严格一致否则会报ClassNotFoundException。若使用 Hadoop 3.xJAR 包路径通常在share/hadoop/tools/lib/下Hadoop 2.x 则多在contrib/streaming/目录。2.3 mrjob 作业的核心编码模式与参数解析以mr1.py的关键片段为例说明 Python 如何映射 MapReduce 语义from mrjob.job import MRJob from mrjob.step import MRStep class MRMovieCooccurrence(MRJob): def mapper(self, _, line): # 解析 u.data格式为 user_id\tmovie_id\trating\ttimestamp fields line.strip().split(\t) if len(fields) 2: user_id, movie_id fields[0], fields[1] # 输出(movie_id, user_id)为后续按 movie_id 分组做准备 yield movie_id, user_id def reducer(self, movie_id, user_ids): # 收集所有给该电影打分的用户列表 users list(user_ids) # 两两组合生成共现对(movie_id1, movie_id2) 表示被同一用户评分 for i in range(len(users)): for j in range(i 1, len(users)): # 注意此处实际应关联用户评分的所有电影但本项目简化为同用户ID的电影对 # 真实实现需先构建 user-movies 映射此处为教学精简 yield (users[i], users[j]), 1 def steps(self): return [ MRStep(mapperself.mapper, reducerself.reducer) ]mapper方法接收原始行数据按\t切分后提取user_id和movie_id并以movie_id为 key、user_id为 value 输出。这步看似简单实则决定了 Shuffle 阶段的数据分区逻辑——Hadoop 会将相同movie_id的所有(movie_id, user_id)发送给同一个 Reducer。reducer方法接收movie_id和其对应的所有user_id列表然后对用户列表做两两组合生成(user_id_i, user_id_j)作为新 key并赋予计数1。这是共现统计的起点后续mr2.py会基于此 key 进一步聚合。steps()方法定义作业执行流程支持多阶段串联如mr1→mr2是 mrjob 区别于裸 Hadoop Streaming 的关键抽象。注意本项目mr1.py的 reducer 实现存在教学简化——真实物品协同过滤需先构建user → [movie1, movie2, ...]映射再对每个用户评分的电影两两组合。当前代码若直接运行会产生逻辑错误。修正方案见第 4 章排错部分。3. 从零部署Hadoop 伪分布式环境搭建与项目运行全流程3.1 Hadoop 单机伪分布式环境最小化配置本项目不强制要求全集群Hadoop 伪分布式Pseudo-Distributed Mode即可满足所有功能验证。以下是 Ubuntu 22.04 下的精简配置步骤以 Hadoop 3.3.6 为例步骤 1安装 Java 11 并配置环境变量sudo apt update sudo apt install openjdk-11-jdk -y echo export JAVA_HOME/usr/lib/jvm/java-11-openjdk-amd64 ~/.bashrc echo export PATH$JAVA_HOME/bin:$PATH ~/.bashrc source ~/.bashrc java -version # 验证输出包含 openjdk version 11.步骤 2下载并解压 Hadoopcd /opt sudo wget https://downloads.apache.org/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz sudo tar -xzf hadoop-3.3.6.tar.gz sudo chown -R $USER:$USER hadoop-3.3.6步骤 3配置核心 XML 文件仅修改关键项编辑/opt/hadoop-3.3.6/etc/hadoop/core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration编辑/opt/hadoop-3.3.6/etc/hadoop/hdfs-site.xmlconfiguration property namedfs.replication/name value1/value !-- 单节点设为1 -- /property property namedfs.namenode.name.dir/name valuefile:/opt/hadoop-3.3.6/data/namenode/value /property property namedfs.datanode.data.dir/name valuefile:/opt/hadoop-3.3.6/data/datanode/value /property /configuration注意namenode.name.dir和datanode.data.dir对应的目录需手动创建mkdir -p /opt/hadoop-3.3.6/data/{namenode,datanode}步骤 4格式化 NameNode 并启动服务# 格式化文件系统仅首次运行 /opt/hadoop-3.3.6/bin/hdfs namenode -format # 启动 HDFS /opt/hadoop-3.3.6/sbin/start-dfs.sh # 验证进程应看到 NameNode 和 DataNode jps # 创建 HDFS 输入目录并上传数据 /opt/hadoop-3.3.6/bin/hdfs dfs -mkdir -p /user/input /opt/hadoop-3.3.6/bin/hdfs dfs -put /path/to/your/u.data /user/input/3.2 项目依赖安装与数据预处理在 Python 环境中安装 mrjob 及其依赖pip install mrjob0.7.4 # 本项目适配 0.7.x 版本避免与新版 mrjob 不兼容 pip install pandas numpy # 用于 run.py 中的结果解析检查数据文件格式是否合规u.data必须为 Unix 换行符LF且无 BOM 头。可用file u.data验证输出应含CRLF或LF若u.data来自 Windows用dos2unix u.data转换u.item中电影类型字段以|分隔需确保无多余空格或转义字符。3.3 执行 run.py 并监控任务状态进入项目根目录执行主控脚本python run.pyrun.py内部会依次调用mr1.py和mr2.py并自动处理 HDFS 路径。关键监控点终端输出观察Running step 1 of 1...及Streaming final output from ...日志确认无IOException或ClassNotFoundExceptionHadoop Web UI浏览器访问http://localhost:9870NameNode UI点击 “Utilities” → “Browse the file system”导航至/user/output/mr2/确认part-00000文件存在且非空结果验证下载part-00000并检查前几行1,2,0.8571428571428571 1,3,0.7071067811865475 2,4,0.9128709291752769格式为movie_id1,movie_id2,similarity_score符合预期。提示若run.py报错No module named mrjob请确认当前 Python 环境which python与pip install mrjob的环境一致若报错Connection refused检查 HDFS 是否已启动jps是否显示 DataNode。4. 关键逻辑修正与性能调优解决共现统计偏差与小数据集冷启动问题4.1 修复 mr1.py 中的共现逻辑缺陷原始mr1.py的reducer方法存在根本性错误它将movie_id作为 mapper 的 key却在 reducer 中对user_id列表做两两组合这导致同一用户对不同电影的评分无法被关联。正确做法是mapper 应输出(user_id, movie_id)reducer 收集每个用户的全部电影列表再进行两两组合。修正后的mr1.py核心代码如下def mapper(self, _, line): fields line.strip().split(\t) if len(fields) 2: user_id, movie_id fields[0], fields[1] # 关键修正以 user_id 为 keymovie_id 为 value yield user_id, movie_id def reducer(self, user_id, movie_ids): movies list(movie_ids) # 对该用户评分的所有电影两两组合 for i in range(len(movies)): for j in range(i 1, len(movies)): # 输出共现对(movie_i, movie_j) 作为 key计数为 1 yield (movies[i], movies[j]), 1此修正确保了“同一用户评分的任意两部电影”均被计入共现是物品协同过滤的数学基础。若跳过此步mr2.py计算的相似度将完全失真。4.2 mr2.py 中的余弦相似度实现与参数调优mr2.py的 reducer 负责将mr1.py输出的共现对(m1,m2)聚合计算余弦相似度def reducer(self, movie_pair, counts): count_list list(counts) co_occurrence sum(count_list) # 共现次数 # 获取 m1 和 m2 各自的总评分次数需提前统计本项目通过辅助脚本 precompute_counts.py 完成 # 假设已存入全局字典 movie_count {1: 120, 2: 95, ...} m1, m2 movie_pair count_m1 self.movie_count.get(m1, 1) # 防止除零 count_m2 self.movie_count.get(m2, 1) # 余弦相似度 co_occurrence / sqrt(count_m1 * count_m2) similarity co_occurrence / (count_m1 * count_m2) ** 0.5 yield movie_pair, similarity为提升计算稳定性建议在mr2.py初始化时加载预计算的电影评分频次def __init__(self, *args, **kwargs): super(MRMovieSimilarity, self).__init__(*args, **kwargs) # 从 HDFS 读取预计算的 movie_count.json self.movie_count self.load_movie_counts()预计算脚本precompute_counts.py可用以下命令生成# 统计每部电影被评分的总次数 /opt/hadoop-3.3.6/bin/hdfs dfs -cat /user/input/u.data | \ awk -F\t {print $2} | sort | uniq -c | \ awk {print \ $2 \: $1 ,} movie_count.json4.3 小数据集下的冷启动优化技巧MovieLens 的u.data仅含 10 万条记录直接计算全量物品相似度效率低且稀疏。实用优化策略阈值过滤在mr2.py的 reducer 中增加最小共现阈值丢弃co_occurrence 5的电影对Top-K 截断每个电影只保留相似度最高的 20 个邻居减少result.csv体积缓存元数据将u.item加载为内存字典在run.py中将result.csv的数字 ID 替换为电影标题提升可读性import pandas as pd items pd.read_csv(u.item, sep|, encodingISO-8859-1, headerNone) item_dict dict(zip(items[0].astype(str), items[1])) result_df[movie1_title] result_df[movie_id1].map(item_dict) result_df[movie2_title] result_df[movie_id2].map(item_dict)注意u.item的编码为ISO-8859-1非 UTF-8直接用pd.read_csv会报错必须显式指定encoding参数。5. 结果验证与推荐效果评估用 Python 脚本快速生成用户推荐列表5.1 从 result.csv 构建物品相似度索引result.csv是扁平化的相似度对需转换为可查询的字典结构。以下脚本build_similarity_index.py将生成similarity_index.pklimport pandas as pd import pickle # 读取 result.csv注意处理引号 df pd.read_csv(result.csv, names[movie_id1, movie_id2, similarity], skiprows1, # 跳过可能的 header quotechar, enginepython) # 构建双向索引{movie_id: [(similar_movie_id, score), ...]} sim_index {} for _, row in df.iterrows(): m1, m2, score str(row[movie_id1]), str(row[movie_id2]), row[similarity] if m1 not in sim_index: sim_index[m1] [] if m2 not in sim_index: sim_index[m2] [] sim_index[m1].append((m2, score)) sim_index[m2].append((m1, score)) # 按相似度降序排列每个电影的邻居 for movie in sim_index: sim_index[movie].sort(keylambda x: x[1], reverseTrue) # 保存为 pickle供推荐脚本调用 with open(similarity_index.pkl, wb) as f: pickle.dump(sim_index, f)运行后生成similarity_index.pkl体积小、加载快是后续推荐的基石。5.2 为指定用户生成 Top-10 推荐列表generate_recommendations.py脚本演示如何结合用户历史行为与相似度索引生成推荐import pickle import pandas as pd # 加载相似度索引和用户历史 with open(similarity_index.pkl, rb) as f: sim_index pickle.load(f) # 读取用户历史假设用户 196 的历史评分为 {movie_id: rating} user_history {} u_data pd.read_csv(u.data, sep\t, headerNone, names[user_id,movie_id,rating,timestamp]) user_196 u_data[u_data[user_id] 196] for _, row in user_196.iterrows(): user_history[str(row[movie_id])] row[rating] # 生成推荐对用户看过的每部电影取其 top-5 相似电影加权累加相似度 recommendations {} for watched_movie, rating in user_history.items(): if watched_movie in sim_index: for similar_movie, similarity in sim_index[watched_movie][:5]: if similar_movie not in user_history: # 过滤已评分电影 score rating * similarity recommendations[similar_movie] recommendations.get(similar_movie, 0) score # 按综合得分排序取 top-10 top10 sorted(recommendations.items(), keylambda x: x[1], reverseTrue)[:10] # 加载电影标题并打印 items pd.read_csv(u.item, sep|, encodingISO-8859-1, headerNone) item_dict dict(zip(items[0].astype(str), items[1])) print(User 196 Top-10 Recommendations:) for movie_id, score in top10: title item_dict.get(movie_id, Unknown Movie) print(f{title} (ID:{movie_id}) - Score: {score:.4f})运行此脚本你将看到类似输出User 196 Top-10 Recommendations: Star Wars (1977) (ID:2) - Score: 4.2187 Contact (1997) (ID:286) - Score: 3.9521 ...这证明整个 Hadoop 批处理链路产出的相似度数据能被下游 Python 应用直接消费完成端到端推荐闭环。提示若generate_recommendations.py报错KeyError说明result.csv中未覆盖某些电影 ID可在sim_index.get(movie_id, [])中添加默认空列表防御。本文还有配套的精品资源点击获取