资讯动态

基于Hadoop的协同过滤视频推荐系统:手工实现MapReduce推荐链路

发布时间:2026/10/9 14:33:56 来源:尧图企业网站定制
简介一份面向大数据与推荐系统学习者的Hadoop协同过滤视频推荐系统完整项目包针对视频数量激增带来的推荐性能瓶颈与扩展性问题演示如何借助HDFS与MapReduce处理海量用户行为数据并实现基于用户/项目的协同过滤推荐适用于毕业设计、课程设计或Hadoop技术入门实践。压缩包共383个文件大小12.1MB包含58个Java源码、88个JavaScript、54个CSS与13个HTML前端页面、78个PNG/66个JPG等图片资源另有SQL初始化脚本、XML及YAML配置文件等整体结构覆盖后端逻辑、前端展示与数据存储。已有24人学习下载适合希望查看完整代码思路、数据库建表语句及Hadoop项目实际组织方式的读者借助该包可快速搭建可演示的推荐系统原型并加深对Hadoop生态与协同过滤算法结合的理解。1. 基于Hadoop的协同过滤视频推荐系统与其调库不如亲手拆一遍推荐链路做视频推荐的同学多半被问过这样一个问题“你用Spark或者PyTorch跑协同过滤那如果把Spark去掉只给你一套Hadoop你还能不能把推荐跑出来”这不是面试官故意刁难而是很多课程设计和毕业设计里题目写的就是“基于Hadoop的协同过滤视频推荐系统”。这种项目拿到的往往是一个zip包里面有代码、数据集、文档但真正的价值不在压缩包里的那一坨文件而在于你重启一遍这套离线推荐流程时踩过的每一个MapReduce的坑。这篇文章就是要带你从头搭一遍用伪分布式Hadoop配合Python Streaming实现协同过滤不依赖Spark、不调现成推荐库把相似度计算和TopN推荐老老实实跑出来。适合准备大数据课程设计、想把MapReduce彻底搞懂、或者简历上需要一个“亲手实现的推荐系统”的从业者。2. 系统架构与协同过滤选型先把推荐逻辑拆清楚2.1 协同过滤的两种路线基于用户 vs 基于物品视频场景选谁协同过滤的核心是“物以类聚、人以群分”。基于用户的协同过滤UserCF先找与当前用户口味相似的其他用户再用这些用户看过的视频来推荐基于物品的协同过滤ItemCF则先计算视频和视频之间的相似度然后推荐和用户历史中高分视频相似的视频。在视频场景里用户数通常是百万级甚至亿级而视频数可能是几十万级且视频的更新速度远低于新用户的涌入速度。如果做UserCF需要实时维护一个用户相似度矩阵这个矩阵的规模是用户数的平方在离线批处理里非常容易“爆掉”。而ItemCF的相似度矩阵规模是物品数的平方相对可控并且视频的相似关系变化没那么快可以提前算好每天或每周更新一次。所以这个项目我最终选择基于物品的协同过滤。这也是YouTube早期论文中提到的离线推荐思路先根据用户行为找出视频间的关联再为用户生成TopN候选。相似度计算方式有很多种常见的有余弦相似度、皮尔逊相关系数、Jaccard相似度。对于隐式反馈如用户是否观看过我们通常用“共现次数”作为基础同时被同一个用户看过的两个视频就在它们之间产生一次共现。为了抵消热门视频的干扰简单做法是拿共现次数除以两个视频各自被观看次数的几何平均得到类似Jaccard归一化后的分值。这个项目里我们第一版就用最基础的共现次数后面在进阶章再讲怎么加归一化和时间衰减。2.2 数据模型与HDFS目录设计评分数据怎么组织不要一上来就想复杂。既然是基于Hadoop的离线推荐输入数据就是一张很简单的评分表字段包括用户ID、视频ID、用户行为分播放完成率、点赞、划走等都折算成0到1的浮点数。如果没有真实数据集可以用MovieLens的评分数据但注意评分是1-5的整数我们需要把它归一化到0-1。或者自己造一份CSV几万条记录足够跑通全流程。在HDFS上我习惯这样规划目录/recsys/input/ratings.csv原始评分数据。/recsys/recjob1/output第一个MapReduce任务输出物品共现矩阵。/recsys/recjob2/output第二个MapReduce任务输出每个用户的推荐候选及分数。/recsys/cache/similarity.txt第一个任务的结果存入分布式缓存供第二个任务读取。这里要说明为什么一定要把中间结果放在HDFS而不是本地路径因为后续Job需要从分布式缓存中加载相似度矩阵虽然单机伪分布式下本地文件也能用-files上传但正式环境里HDFS路径是最稳妥的。如果你把中间结果写到本地换到集群环境就会立刻翻车。2.3 离线推荐的完整Pipeline从原始日志到TopN列表整个推荐流程可以拆成四个阶段数据清洗把原始行为日志解析成用户ID、视频ID、行为分并过滤掉播放时长过短的记录。计算物品相似度利用MapReduce统计物品共现矩阵再归一化得到相似度。生成推荐候选对每个用户的历史视频从相似度矩阵中找出最相似的K个视频按分数加权累加。结果入库把TopN列表写回MySQL或者导出CSV供线上接口查询。在这个zip项目里最核心的是第2和第3阶段。很多人以为用Hadoop做推荐需要写很复杂的Java类其实用Hadoop Streaming加Python脚本几十行就能完成。而且Python在处理文本和字典时有天然优势方便在Mapper和Reducer里做内存缓存。下一章先把Hadoop环境跑起来这步过不了后面全白搭。3. 在虚拟机里搭好Hadoop伪分布式一步步跑通WordCount3.1 为什么用伪分布式而非单机或完全分布式伪分布式是指所有Hadoop进程NameNode、DataNode、JobTracker或ResourceManager、NodeManager都在同一台机器上但它使用HDFS作为底层文件系统走完整的MapReduce流程。相比单机本地模式Local Mode它能暴露更多分布式环境的问题相比完全分布式它只需一台机器和一个配置文件特别适合在自己笔记本的虚拟机里折腾。而且伪分布式是面试最爱问的“hadoop伪分布式搭建”场景很多热词也都指向这一步。我曾经遇到过一个同学直接在公司三台生产服务器上搭完全分布式结果用了两天还没搞定其实他连伪分布式的配置都没跑通过。我的建议是先在自己电脑上用VMware或VirtualBox开一台Ubuntu虚拟机配2核4G内存把伪分布式跑熟再去扩机器。3.2 配置文件的修改与核心参数说明Hadoop 2.x和3.x的配置基本一致。我这里以Hadoop 2.10.1为例你也可以用3.x但注意端口和部分参数名不一样。用hadoop用户登录解压后进入etc/hadoop目录。先改core-site.xmlconfiguration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/home/hadoop/hadoop_tmp/value /property /configurationfs.defaultFS指定了NameNode的地址9000是RPC端口。hadoop.tmp.dir是HDFS元数据存储的根目录默认在/tmp下系统重启会丢必须改到用户目录。这一步不改后续格式化NameNode时会莫名其妙找不到数据目录。再改hdfs-site.xmlconfiguration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/home/hadoop/hadoop_tmp/dfs/name/value /property property namedfs.datanode.data.dir/name value/home/hadoop/hadoop_tmp/dfs/data/value /property /configuration伪分布式只有一台机器副本数不设成1的话DataNode会一直报“块副本不足”的警告。dfs.replication1是必须的。接着是mapred-site.xml默认没有这个文件要把mapred-site.xml.template复制一份cp mapred-site.xml.template mapred-site.xml然后编辑configuration property namemapreduce.framework.name/name valueyarn/value /property /configuration最后是yarn-site.xml。需要配置ResourceManager的地址和NodeManager的辅助服务configuration property nameyarn.resourcemanager.hostname/name valuelocalhost/value /property property nameyarn.nodemanager.aux-services/name valuemapreduce_shuffle/value /property /configurationaux-services不配置的话MapReduce任务会卡在“Default Chained Executor”那里。3.3 启动集群并跑通第一个MapReduce任务配置写好后先格式化NameNodehdfs namenode -format格式化只是在第一次启动前需要做以后不要乱执行否则会把元数据清空。然后启动HDFS和YARNstart-dfs.sh start-yarn.sh启动后用jps检查进程如果是以下5个进程说明启动成功NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager。接下来在HDFS上创建输入目录把本地一个测试文件传上去跑自带的WordCounthdfs dfs -mkdir -p /input echo hello world | hdfs dfs -put - /input/word.txt hadoop jar $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-examples-*.jar wordcount /input /output跑完用hdfs dfs -cat /output/part-r-00000查看结果。如果这里出问题大部分是环境变量没配好检查~/.bashrc里有没有export HADOOP_HOME/home/hadoop/hadoop export PATH$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin另外虚拟机的内存至少要给1.5G不然启动时多个Java进程会互相挤死。这一套跑通你的伪分布式环境就算稳了。4. 用MapReduce实现协同过滤两个Job搞定相似度与推荐4.1 Mapper与Reducer设计把评分矩阵转成同现矩阵现在开始写推荐逻辑。第一个Job的目标是统计视频之间的共现次数。输入是ratings.csv每一行是userId,videoId,score。Map阶段按用户分组把同一个用户看过或评分高于阈值的视频两两配对输出。这里有个关键优化如果某个用户看了上千个视频两两配对会产生O(n^2)的中间结果容易造成网络风暴。一般做法是只取该用户评分最高的N个视频再配对比如80个把中间数据量压下来。我们的Mapper里加一个过滤操作。下面是用Hadoop Streaming写Python脚本文件是job1_mapper.py#!/usr/bin/env python # job1_mapper.py import sys MAX_ITEMS_PER_USER 50 def parse_line(line): # 期望格式: userId, videoId, score parts line.strip().split(,) if len(parts) 3: return None try: uid parts[0] vid parts[1] score float(parts[2]) return uid, vid, score except ValueError: return None current_uid None videos [] for line in sys.stdin: parsed parse_line(line) if parsed is None: continue uid, vid, score parsed if uid ! current_uid: # 切换用户时处理上一个用户的视频列表 if current_uid is not None: # 只取分数最高的MAX_ITEMS_PER_USER个视频 videos.sort(keylambda x: x[1], reverseTrue) top videos[:MAX_ITEMS_PER_USER] n len(top) for i in range(n): for j in range(i1, n): a top[i][0] b top[j][0] # 输出物品对小的在前大的在后避免重复 if a b: print(f{a}\t{b}\t1) else: print(f{b}\t{a}\t1) current_uid uid videos [] videos.append((vid, score)) # 处理最后一个用户 if current_uid is not None: videos.sort(keylambda x: x[1], reverseTrue) top videos[:MAX_ITEMS_PER_USER] n len(top) for i in range(n): for j in range(i1, n): a top[i][0] b top[j][0] if a b: print(f{a}\t{b}\t1) else: print(f{b}\t{a}\t1)这段代码的核心是用current_uid攒一个用户的所有观看记录等用户切换时再统一配对。为什么要攒着而不是每行都输出因为每个用户需要两两配对必须在一个Mapper实例内把同key的数据收集齐。Hadoop Streaming的Mapper是按行读取的所以我们需要手动维护用户缓冲。MAX_ITEMS_PER_USER是关键参数我一般设为50。视频领域中一个用户真正认真看过的视频很少超过50个再多的都对相似度计算没帮助反而会造成数据膨胀。如果这个值设太大比如500Mapper输出的中间数据会猛增Reducer压力也变大。Reducer端就是简单的累加job1_reducer.py#!/usr/bin/env python # job1_reducer.py import sys current_pair None count 0 for line in sys.stdin: parts line.strip().split(\t) if len(parts) ! 3: continue a, b, c parts key (a, b) if key current_pair: count int(c) else: if current_pair is not None: print(f{current_pair[0]}\t{current_pair[1]}\t{count}) current_pair key count int(c) if current_pair is not None: print(f{current_pair[0]}\t{current_pair[1]}\t{count})这里没有做归一化输出的是纯共现次数。你可以在Reducer里直接除以两个视频的总观看次数但那样需要额外Job统计视频流行度课程设计阶段先不搞复杂。运行第一个Job的命令hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input /recsys/input/ratings.csv \ -output /recsys/recjob1/output \ -mapper python job1_mapper.py \ -reducer python job1_reducer.py \ -file job1_mapper.py \ -file job1_reducer.py \ -numReduceTasks 2-file参数会把本地脚本分发到集群的每个节点。伪分布式下只有一个节点所以也可以不加但为了以后迁移到多机最好加上。-numReduceTasks设2可以直观看到有多个part文件输出。4.2 Reducer中的相似度计算归一化与阈值第二个Job是真正的推荐计算。输入是用户的评分记录我们需要为每个用户计算他可能喜欢的视频。思路对每个用户遍历他高质量观看过的视频比如评分大于0.6的从第一个Job生成的相似度矩阵中查出与这些视频相似的视频把相似度乘以用户评分作为加权分最后按总分排序取TopN。这里需要一点技巧Hadoop Streaming的Mapper是逐行读取用户评分的但我们还需要在每个Mapper里加载完整的相似度矩阵。如果把相似度矩阵直接读进内存Mapper每处理一行都要全矩阵扫描太慢。所以我们要使用Hadoop的分布式缓存。先把第一个Job的输出上传到HDFS缓存目录然后把它加到DistributedCache中。第二个Job的Mapper逻辑#!/usr/bin/env python # job2_mapper.py import sys # 加载相似度矩阵到字典 # 分布式缓存会把这个文件放到本地工作目录 SIMILARITY_FILE similarity.txt similarity {} # 用全局变量缓存加载结果避免每个Mapper实例重复加载 _cache_loaded False def load_similarity(): global _cache_loaded if _cache_loaded: return with open(SIMILARITY_FILE, r) as f: for line in f: parts line.strip().split(\t) if len(parts) ! 3: continue a, b, score parts[0], parts[1], float(parts[2]) # 存两遍这样查a的邻居和b的邻居都方便 if a not in similarity: similarity[a] [] similarity[a].append((b, score)) if b not in similarity: similarity[b] [] similarity[b].append((a, score)) _cache_loaded True # 用户历史观看记录 user_history {} current_uid None load_similarity() for line in sys.stdin: uid, vid, score line.strip().split(,) if float(score) 0.6: # 只考虑有效观看 continue if uid not in user_history: user_history[uid] [] user_history[uid].append((vid, float(score))) # 这里不能立刻输出因为我们要等所有输入行读取完才能为每个用户完整计算 # 但Streaming Mapper是每行输出一次所以需要把用户历史攒到最后统一处理 # 实际上我们需要存buffer # 由于Streaming的限制这里简化先缓存所有用户数据最后一次性输出问题来了Hadoop Streaming的Mapper是每读一行输出一行没法像前面Job那样在current_uid切换时输出。因为推荐计算要等一个用户的所有历史记录都读完才能和相似度矩阵做乘法。一个Map任务可能处理多个用户需要把所有用户的数据攒在内存里。对于几十万条数据放在内存中还是可行的但要小心内存溢出。一种更稳妥的做法是把第二个Job的输入变成“按用户分组后的一行”的形式这需要第三个Job来做预处理。但既然叫“两个Job搞定”我们可以这样折中在Mapper里仍按用户攒数据等所有输入读完后在脚本末尾统一输出。Hadoop Streaming允许Mapper在进程结束时输出这虽然破坏了“流式”的精神但在数据处理量不大时完全可行。改一下job2_mapper.py#!/usr/bin/env python # job2_mapper.py import sys SIMILARITY_FILE similarity.txt similarity {} _cache_loaded False def load_similarity(): global _cache_loaded if _cache_loaded: return with open(SIMILARITY_FILE, r) as f: for line in f: parts line.strip().split(\t) if len(parts) ! 3: continue a, b, score parts[0], parts[1], float(parts[2]) if a not in similarity: similarity[a] [] similarity[a].append((b, score)) if b not in similarity: similarity[b] [] similarity[b].append((a, score)) _cache_loaded True users {} load_similarity() for line in sys.stdin: uid, vid, score line.strip().split(,) s float(score) if s 0.6: continue if uid not in users: users[uid] [] users[uid].append((vid, s)) # 所有输入读完后统一为每个用户做推荐 for uid, history in users.items(): # 候选视频字典video - 累计分数 rec_score {} for vid, s in history: # 找出相似视频 neighbors similarity.get(vid, []) for nid, sim in neighbors: if nid vid: # 去掉自己 continue # 如果用户已经看过可以跳过或者保留作为“再次推荐” if any(nid hv for hv, _ in history): continue rec_score[nid] rec_score.get(nid, 0) s * sim # 取TopN这里N10 topn sorted(rec_score.items(), keylambda x: x[1], reverseTrue)[:10] for rid, rscore in topn: print(f{uid}\t{rid}\t{rscore})注意这个Mapper把整个用户历史和候选矩阵都放在内存里所以第二个Job的Map数不能太大否则容易内存不足。伪分布式下一个内存4G的虚拟机跑几万条数据没问题。如果你的真实数据量大需要把“同现矩阵转置”或用多次Join但那是生产级设计这里先不展开。Reducer就简单了把同一个用户的推荐结果合并起来按分数排序输出#!/usr/bin/env python # job2_reducer.py import sys current_uid None recs [] for line in sys.stdin: parts line.strip().split(\t) if len(parts) ! 3: continue uid, vid, score parts if uid ! current_uid: if current_uid is not None: recs.sort(keylambda x: x[1], reverseTrue) for v, s in recs[:10]: print(f{current_uid}\t{v}\t{s}) current_uid uid recs [] recs.append((vid, float(score))) if current_uid is not None: recs.sort(keylambda x: x[1], reverseTrue) for v, s in recs[:10]: print(f{current_uid}\t{v}\t{s})运行第二个Job之前要把第一个Job的结果放到HDFS缓存目录并设置缓存hdfs dfs -mkdir -p /recsys/cache hdfs dfs -cp -f /recsys/recjob1/output/part-* /recsys/cache/similarity.txt hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input /recsys/input/ratings.csv \ -output /recsys/recjob2/output \ -mapper python job2_mapper.py \ -reducer python job2_reducer.py \ -file job2_mapper.py \ -file job2_reducer.py \ -files /recsys/cache/similarity.txt#similarity.txt \ -numReduceTasks 1-files参数后面用#给远程文件起了个别名否则Mapper里需要用完整的HDFS路径来打开文件。Streaming会把similarity.txt下载到每个Map任务的工作目录。4.3 包装成可执行的Driver脚本输入输出路径与运行配置每次敲这么长的hadoop jar命令容易出错我一般把它们写成一个Shell脚本放在项目根目录下叫run_recsys.sh。里面提前定义好路径变量并加上了Job失败时清理输出目录的逻辑#!/bin/bash # run_recsys.sh INPUT/recsys/input/ratings.csv SIM_OUT/recsys/recjob1/output REC_OUT/recsys/recjob2/output CACHE_DIR/recsys/cache # 清理输出目录避免重复运行时报错 hdfs dfs -rm -r -f $SIM_OUT $REC_OUT $CACHE_DIR echo Step 1: 计算物品共现矩阵 hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input $INPUT \ -output $SIM_OUT \ -mapper python job1_mapper.py \ -reducer python job1_reducer.py \ -file job1_mapper.py \ -file job1_reducer.py \ -numReduceTasks 2 echo Step 2: 复制相似度矩阵到缓存 hdfs dfs -mkdir -p $CACHE_DIR hdfs dfs -cp -f $SIM_OUT/part-* $CACHE_DIR/similarity.txt echo Step 3: 生成用户推荐 hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input $INPUT \ -output $REC_OUT \ -mapper python job2_mapper.py \ -reducer python job2_reducer.py \ -file job2_mapper.py \ -file job2_reducer.py \ -files $CACHE_DIR/similarity.txt#similarity.txt \ -numReduceTasks 1 echo Done! 推荐结果在 $REC_OUT执行bash run_recsys.sh。结束后用hdfs dfs -cat /recsys/recjob2/output/part-00000 | head -20查看结果。每一行是“用户ID\t视频ID\t推荐分”。如果你能看到输出说明你自己的基于Hadoop的协同过滤推荐系统已经跑通了。5. 避坑指南Hadoop排序、数据倾斜与路径配置的三个翻车现场5.1 相似度结果全是0整数整除的玄学现象第一个Job跑完后查看part文件发现所有分数都是0但共现次数明明不是0。原因在Reducer里如果写了类似score count / (count1 * count2)当所有变量都是整数时Python 2会直接截断为0。如果你用的是Python 2运行Streaming这就是最典型的错误。或者你在Java里用了int类型除法。解决在计算相似度前把分母或分子强制转成float。在job1_reducer.py里如果做归一化不要写count / pop_a而是写count * 1.0 / pop_a。更稳妥的做法是直接用float(count)。另外一个隐蔽点CSV里的score如果写成0代表用户看过但没有打分需要提前过滤否则会被当成负样本。5.2 第二个Job启动后一直卡住YARN任务超时现象运行run_recsys.sh时第二个Job停在任务进度98%日志里看不到Mapper输出。原因最常见的是-files路径写错或者HDFS缓存目录里没有similarity.txt。我一开始就是直接在-files后写/recsys/cache/similarity.txt但没加#similarity.txt别名导致Mapper打开文件时找不到。解决检查hdfs dfs -ls /recsys/cache/确保有similarity.txt。执行Hadoop命令时把-files参数改成/recsys/cache/similarity.txt#similarity.txt这样Streaming会把它放到工作目录名字就是similarity.txt。另一点是个好习惯在所有Job启动前先清理输出目录否则Hadoop会报“目录已存在”而直接退出。5.3 数据倾斜某个Reducer处理了90%的数据其他Reducer闲等现象第一个Job设置了2个Reducer但part-00000有50MBpart-00001只有5KB。推荐结果集中在几个热门视频上。原因热门视频和其他视频的共现次数特别高导致 hash在计算时把大量key分到了同一个Reducer。这是协同过滤在MapReduce实现上的经典坑也叫“推荐系统的冷门与热门分布不均问题”。解决在Mapper输出配对时对热门视频做降权或裁剪。我在job1_mapper.py里加了一个MAX_ITEMS_PER_USER这已经能抑制一部分如果还不够可以把热门视频出现次数前1%直接排除出相似度计算或者给它们的共现次数加一个衰减因子1/log(1count)。还有一种工程更重的做法是加随机前缀打散key但那样Reducer需要二次聚合复杂度高课程设计阶段我不建议做。5.4 伪分布式下DataNode没启动Web界面打不开现象start-dfs.sh执行完毕但jps看不到DataNode进程。访问http://localhost:500702.x是500703.x是9870页面显示“NameNode active but DataNode not running”。原因多半是hdfs-site.xml里的dfs.datanode.data.dir没有创建或者格式化NameNode后DataNode的data目录和NameNode的namespaceID不一致。常见的翻车点是你重新格式化NameNode但DataNode的data目录里已经有旧的VERSION文件导致集群ID不匹配。解决第一次格式化前确定dfs.namenode.name.dir和dfs.datanode.data.dir的目录不存在或者手动删除/home/hadoop/hadoop_tmp/dfs下的所有内容然后再直接执行hdfs namenode -format。格式化完成后再start-dfs.sh。如果还报错可以登录到DataNode日志中查看hadoop-datanode.log里面会明确指出是namenode地址不通还是clusterID冲突。5.5 中文用户名或路径导致本地文件无法上传现象准备用hdfs dfs -put ratings.csv /recsys/input/上传中文命名的CSV时抛异常java.io.FileNotFoundException但文件明明存在。原因Hadoop原生对中文文件名支持不佳尤其是Python脚本里open()函数默认编码和系统编码不一致。我的虚拟机是Ubuntu默认Locale是英文而CSV用记事本存成了UTF-8带BOMBOM会被当作字符串的一部分。解决所有文件名和字段改成英文小写用iconv -f GBK -t UTF-8转换编码Python脚本开头加# -*- coding: utf-8 -*-并且用codecs模块打开文件时可以utf-8-sig来清除BOM。还有一个小血泪如果是在Windows下编辑的脚本传到Linux后要执行dos2unix否则脚本首行会出现#!/usr/bin/env python\r导致“No such file or directory”。6. 验证与进阶用测试集评估推荐效果并给系统加一个参数开关前面跑通了推荐流程但推荐质量好不好不能拍脑袋。最常见的做法是留一法评估把每个用户最高评分的那条记录藏起来用剩下的数据训练看推荐结果里能不能包含被藏起来的那条。如果命中计数加一。最后用“命中数 / 总用户数”作为召回率这是离线推荐最朴素的指标。给我自己的项目写一个小的评估脚本把第二个Job的输出和原始测试集对比。假设你已经把数据划分好为train.csv和test.csv先在HDFS上跑完训练流程再用hdfs dfs -cat把结果拉下来hdfs dfs -cat /recsys/recjob2/output/part-* predictions.txt python evaluate.py predictions.txt test.csvevaluate.py的核心是import sys # 加载预测结果: userId, videoId, score pred {} with open(predictions.txt, r) as f: for line in f: uid, vid, score line.strip().split(\t) pred.setdefault(uid, []).append(vid) # 加载测试集: userId, videoId hit 0 total 0 with open(sys.argv[1], r) as f: for line in f: uid, vid, _ line.strip().split(,) total 1 if uid in pred and vid in pred[uid]: hit 1 print(fhit{10} {hit/total:.4f})我的经验是如果命中率能到1%以上对于一个纯离线MapReduce实现已经很不错了。不要期望Spark能到10%以上的场景因为这里没有加用户个性化特征、时间窗口等。接着我给系统加一个参数开关。在run_recsys.sh里设置两个环境变量MAX_ITEMS_PER_USER50 # Mapper里每个用户最多取多少个视频 TOPN10 # 每个用户推荐多少个视频 export MAX_ITEMS_PER_USER TOPN然后在job1_mapper.py和job2_mapper.py里通过os.environ读取它们。这样调优时不用改代码只改脚本顶部的数值就行。我做过的很多次实验发现MAX_ITEMS_PER_USER设太小比如10会导致相似度矩阵太稀疏设太大比如200会产生大量低质量候选。而TOPN在5到20之间召回率会有明显差异。每个数据集都有自己的最优值这个开关就是我留下的后悔药。最后说一个我自己的习惯每次调完参数都把当时的HDFS输出目录完整复制一份命名成experiment_vXX。这样跑出来的结果可回溯出了问题能直接对比旧数据而不用重新算一遍。这个习惯救过我很多次因为MapReduce任务一旦跑错重新洗数据要花几十分钟。如果你也用这套基于Hadoop的协同过滤方案建议跑之前先想好怎么记录实验版本。希望这篇能帮到你。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑