
簡介這是一套面向計算機專業本科生的畢業設計與課程大作業實戰資源基于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, [])中添加默認空列表防御。本文還有配套的精品資源點擊獲取