简介:面向电商场景的基于Hadoop的商品推荐系统源码包,适合大数据初学者、Java工程师及需要搭建推荐系统原型的高校学生参考学习。资源采用Java与Hadoop技术栈,围绕HDFS分布式存储、MapReduce并行计算框架,完整展示协同过滤、基于内容推荐等经典算法的落地实现思路。压缩包共34个文件,包含29个Java源码文件和5个XML配置(含pom.xml Maven工程配置),代码分层清晰,便于理解用户行为数据清洗、物品相似度计算、推荐列表生成与结果展示等核心流程。资源包整体约25KB,轻量易用,已吸引193人学习浏览。通过研读源码,可掌握Hadoop环境下推荐系统从数据准备到算法实现的基本架构设计,为后续扩展实时计算或混合推荐提供良好起点。
1. 打开这份「商品推荐系统.zip」之前,先想清楚三件事
如果你拿到的是一份名为「基于hadoop的商品推荐系统.zip」的压缩包,大概率它是课程设计或毕业设计的交付物:里面装着 Java 源码、编译好的 jar 包、一份行为数据集,还有一页实验报告模板。它的核心故事是——用 Hadoop 的 MapReduce 离线计算出「看了 A 商品的人,还应该看看 B」,然后把每个用户的 TOP-N 推荐结果写回 HDFS。它能解决的问题很具体:数据量一大、单机内存装不下用户-商品矩阵时,用分布式把相似度计算拆开跑。适合谁?适合正在做大数据课程设计、想跑通一个推荐系统 demo 但不想从头啃 Spark 的人。先泼一盆冷水:这个方向能跑通,但坑不少,尤其环境、数据格式、Reduce 内存这三处,几乎每个人都会翻一次车。
2. 先定算法再碰集群:商品推荐为什么默认选 ItemCF
2.1 电商场景里 ItemCF 为什么压过 UserCF
推荐系统里两个经典流派:UserCF(用户协同过滤)看的是「和你兴趣相似的人喜欢什么」,ItemCF(物品协同过滤)看的是「和你历史行为相关的物品有哪些」。在商品推荐这个题目里,工程上默认选 ItemCF,理由很实在。
电商场景有四个特征:商品数量远大于用户数量、用户兴趣会漂移(上个月买篮球这个月买键盘)、用户行为稀疏(多数人只买过几十件商品)、实时性要求不高但个性化要求高。UserCF 在这种场景下有两个硬伤:一是用户相似度矩阵的规模随用户数平方增长,百万用户根本算不动;二是用户兴趣一变,离线算好的「相似用户」立刻失真。ItemCF 避开了这两个问题,它只需要算商品和商品之间的共现关系,商品数量相对可控,而且「买过 X 的人还买 Y」这种关联有天然的推荐语义,答辩也讲得清楚。
ItemCF 的核心公式是修正后的余弦相似度:
W_ij = |N(i) ∩ N(j)| / sqrt(|N(i)| * |N(j)|)N(i) 是购买过商品 i 的用户集合。分母的 sqrt 是为了惩罚热门商品——否则每个商品都和爆款有高相似度,推荐结果就全变成「买了啥都推荐手机壳」。如果你觉得这个公式眼熟,没错,它和 Jaccard 系数同源,只是多了个评分权重的变体。课程设计里用共现次数 + 余弦修正,完全够用了。
2.2 从原始评分表到「用户-商品」矩阵:数据还得自己先洗一遍
我见过太多人拿到数据集直接丢给 Hadoop,结果 Mapper 里解析字段就崩了。以最常见的开源行为数据集 MovieLens 100K 为例,它的 u.data 文件格式是 TSV:
user_id \t item_id \t rating \t timestamp 196 242 3 881250949 186 302 3 891717742注意第一行是数据不是表头,分隔符是 Tab 不是逗号。拿到手先做三件事:确认分隔符、砍掉不需要的列、过滤低质量评分。我一般用一条 awk 命令先探路:
# 先看前5行,确认字段分隔符和内容范围 head -5 ml-100k/u.data # 抽成user_id,item_id,rating三列,输出为CSV awk -F'\t' '{print $1","$2","$3}' ml-100k/u.data > behavior_clean.csv # 统计用户数和商品数,决定后续Reduce设计 wc -l behavior_clean.csv cut -d',' -f1 behavior_clean.csv | sort -u | wc -l cut -d',' -f2 behavior_clean.csv | sort -u | wc -l逻辑说明:awk 指定 Tab 分隔(-F'\t'),只取前三列,把评分保留下来但因为 ItemCF 只看「买没买」,评分后面只会用作权重,不会参与协同过滤计算。统计用户数这一步很重要——如果用户数只有几百,说明数据量太小,跑 MapReduce 属于杀鸡用牛刀,后面你会看到输出结果慢得离谱。
为什么评分列不能直接丢?因为阶段三计算推荐分数时,用户历史行为的分值会影响最终排序。低评分(比如 1 分、2 分)的行为应该过滤掉,它代表用户不喜欢,不能当正样本。我一般保留 rating >= 3 的记录,这一步在 awk 里就能做:
awk -F'\t' '{if($3>=3) print $1","$2","$3}' ml-100k/u.data > behavior_clean.csv2.3 评测要走在开发前面:用精确率召回率卡住「做得对不对」
很多课程设计的通病是:推荐结果跑出来了,但不知道对不对。没有评测指标,答辩时老师问「效果怎么样」只能支支吾吾。正确做法是先定评测方案再写代码。
做法是按时间切分数据集:把每个用户的行为按时间戳排序,前 80% 当训练集,后 20% 当测试集。为什么按时间不按随机切分?因为推荐系统的评估必须模拟真实场景——用过去预测未来,随机切分会让模型偷看未来数据,指标虚高。
评测用两个经典指标,TOP-N 精确率(Precision@K)和召回率(Recall@K):
- Precision@K = 推荐列表前 K 个中用户真正消费的商品数 / K
- Recall@K = 推荐列表前 K 个中用户真正消费的商品数 / 测试集商品总数
这两个指标对课程设计来说含义直观:精确率衡量「推荐的准不准」,召回率衡量「把用户想买的都找出来没有」。它们不需要额外的 Python 库,后面第 6 章会给一个 50 行的离线评测脚本。评测方案先行还有一个好处:你的 MapReduce 输出格式可以提前锁定,比如「user_id,item_id,score」,后面写评测脚本时就不用来回改代码。
3. 把 Hadoop 环境搭到「能跑 jar 包」:伪分布式与 Docker 二选一
3.1 伪分布式还是完全分布式:课程设计别盲目搭三台
关于 Hadoop 部署方式,先给结论:课程设计用伪分布式验证逻辑完全够,别一上来就搭三台虚拟机。原因不是技术问题,是时间成本——你只有一个推荐系统要跑,不是在生产环境扛流量。完全分布式集群搭建(热词里那个 hadoop 集群搭建)适合写进实验报告做对比,不适合作为跑通的主要环境。
伪分布式的本质是:一台机器上同时跑 NameNode、DataNode、ResourceManager、NodeManager。配置上和集群版几乎一样,但资源要收敛。给一份我常用的配置参考:
| 配置项 | 推荐值 | 说明 |
|---|---|---|
| mapreduce.map.memory.mb | 512 | 每个 Mapper 容器内存 |
| mapreduce.reduce.memory.mb | 1024 | Reducer 内存给大点,ItemCF 要聚合列表 |
| yarn.nodemanager.resource.memory-mb | 3072 | NodeManager 总可用内存 |
| dfs.replication | 1 | 伪分布式只有一份数据,副本设 1 反而省磁盘 |
| mapreduce.job.reduces | 2 | 阶段一和阶段三的 Reduce 并行度 |
伪分布式对机器的最低要求是 4GB 内存和 20GB 空闲磁盘,如果机器只有 2GB 内存,跑起来必然出现 NodeManager 把容器杀掉的惨案。在这个配置上,你不需要 YARN 的 Capacity Scheduler 调优,默认配置改上面五个参数就够。
3.2 用 Docker 镜像省掉一半玄学:五分钟起一个能跑的环境
如果你是 Mac 或者 Windows 本机,装 Hadoop 原生环境的坑太多了——我踩过的就有 ssh 免密失效、JAVA_HOME 路径不对、权限问题三连。现在有更省事的路:直接用 Hadoop 的 Docker 镜像。官方镜像没有做伪分布式一键启动,但社区里有大量配置好的镜像,拉下来直接起容器就能跑 jar 包。
# 拉取带伪分布式配置的hadoop镜像 docker pull bde2020/hadoop-namenode:2.0.0-hadoop3.2.1 # 启动namenode容器,暴露web端口和ssh端口 docker run -it -d --name hadoop-nn \ -p 9870:9870 \ -p 8088:8088 \ -p 9000:9000 \ -p 22:22 \ bde2020/hadoop-namenode:2.0.0-hadoop3.2.1 # 进入容器确认hdfs能正常读写 docker exec -it hadoop-nn bash hdfs dfs -ls /参数说明:-p 映射了三个关键端口——9870 是 NameNode Web UI(看 HDFS 文件状态),8088 是 YARN ResourceManager UI(看 MapReduce 任务进度),9000 是 NameNode RPC 端口(提交作业时客户端要连它)。22 端口留给 ssh,因为伪分布式里 NodeManager 远程启动容器时需要 ssh 免密。挂载数据目录用 -v 参数,比如 -v ~/data:/opt/data,这样宿主机的数据集直接出现在容器里,上传 HDFS 就不用来回拷文件了。容器方案还有一个隐藏好处:换机器部署时镜像直接搬走,环境不一致导致的玄学问题直接绕开。
3.3 集群目录与输入路径设计:HDFS 里放什么、不放什么
HDFS 目录设计这件事看起来不起眼,但它决定了你的 jar 包是 3 分钟跑完还是 30 分钟跑不完。原则是:输入数据放一个明确目录,中间结果和最终结果分开命名。这是我的标准结构:
# 创建输入输出目录 hdfs dfs -mkdir -p /user/recommend/input hdfs dfs -mkdir -p /user/recommend/output # 上传清洗后的行为数据和商品信息表 hdfs dfs -put behavior_clean.csv /user/recommend/input/ hdfs dfs -put movies.dat /user/recommend/input/ # 确认文件块分布 hdfs dfs -ls /user/recommend/input/注意一个关键限制:MapReduce 的输出目录必须不存在,否则作业直接报错「Output directory already exists」。我每次改完代码重跑前都要先删掉旧输出目录,这一步已经形成了肌肉记忆。还有 input 目录别放太多小文件——后面第 5 章会详细讲 InputSplit 的坑,这里只说结论:一个小文件至少占一个 Map Task,几十个小文件会让任务调度时间远大于计算时间。
3.4 提交 jar 的标准姿势:别把 jar 提交当黑匣子
环境搭好后,提交作业才是真正会卡住人的地方。课程设计的 jar 包如果是从别处拷贝来的「已编译 jar 包」,最常见的坑是本地 JDK 版本和集群不匹配。先确认环境,再提交:
# 确认JAVA_HOME和hadoop版本 java -version hadoop version # 标准提交命令 hadoop jar recommender-1.0.jar \ com.example.RecommendationDriver \ -D mapreduce.job.reduces=2 \ /user/recommend/input/behavior_clean.csv \ /user/recommend/output/cooccurrence # 查看任务进度 yarn application -list推荐把输入输出路径写成 Main 函数接收的 args,而不是写死在代码里。原因很实用:答辩现场可能需要你把输入换成另一份数据,写死在代码里等于要重新编译。参数部分 -D 覆盖的是集群配置,args 传的是业务路径,两者不要混。如果你的 jar 依赖第三方库,用 -libjars 带进去,但推荐系统的核心计算只有 Hadoop 原生 API,不需要额外依赖,能省则省。
提交作业后注意观察 ResourceManager 的日志,如果任务一直显示 ACCEPTED,去看 NodeManager 日志,多半是内存不够或 ssh 免密没配好。伪分布式下最常见的假死就是这两个原因,日志里都会写清楚。
4. 核心实现:共现矩阵做商品相似度,再用三个 MapReduce 串起来
4.1 项目结构:几个 Java 类、两个阶段、一张中间表
很多课程设计的代码把整个流程写成一个类,几百行塞在一起,答辩时自己都讲不清。我用的是标准三段式结构,每个阶段一个类,最后用 Driver 串起来:
src/main/java/com/example/ ├── RecommendationDriver.java # 主入口,串起三个Job ├── UserItemMapper.java # 阶段一:按用户聚合商品 ├── UserItemReducer.java # 阶段一:生成商品共现对 ├── CoOccurrenceMapper.java # 阶段二:聚合共现次数 ├── CoOccurrenceReducer.java # 阶段二:计算相似度并取TopN ├── RecommendMapper.java # 阶段三:为用户生成推荐候选 └── RecommendReducer.java # 阶段三:排序输出TopN三个阶段的关系和数据流向是这样的:阶段一读原始行为数据,输出「商品A,商品B,共现次数」的中间表;阶段二把共现次数转成余弦相似度,输出「商品A,商品B:相似度」;阶段三把用户历史商品和相似度表 join,累加得到候选推荐分,排序取前 N 个。中间表落在 HDFS 的 /output/cooccurrence 目录,它是阶段二和阶段三交接的接口。
4.2 阶段一:从用户行为算出「商品-商品」共现次数
阶段一的核心思路是:一个用户的历史商品列表内部做两两组合,这对组合就是共现关系。假设用户买过 [A, B, C],那共现对就是 (A,B)、(A,C)、(B,C)。实现方式是 Mapper 按用户 id 分发,Reducer 收集同一用户的商品列表后做笛卡尔组合。
// UserItemMapper.java public class UserItemMapper extends Mapper<LongWritable, Text, Text, Text> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 输入: user_id,item_id,rating String[] fields = value.toString().split(","); if (fields.length >= 2) { String userId = fields[0].trim(); String itemId = fields[1].trim(); // key=userId, value=itemId,让同一用户的商品进入同一个Reducer context.write(new Text(userId), new Text(itemId)); } } } // UserItemReducer.java public class UserItemReducer extends Reducer<Text, Text, Text, IntWritable> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { List<String> items = new ArrayList<>(); for (Text val : values) { items.add(val.toString()); } // 同一用户的商品列表内部两两组合 for (int i = 0; i < items.size(); i++) { for (int j = i + 1; j < items.size(); j++) { // 输出 key=商品A,商品B value=1,后续聚合共现次数 String pair = items.get(i) + "," + items.get(j); context.write(new Text(pair), new IntWritable(1)); } } } }逻辑说明:Mapper 输出的 key 是 userId,这样同一用户的全部购买记录会被 HashPartitioner 分到同一个 Reducer。Reducer 里做双重循环生成商品对,同时输出两个方向的组合——这里为了简洁只写了 (A,B) 一个方向,实际生产代码里为了后续查相似度方便,通常把 (A,B) 和 (B,A) 都输出。注意一个隐患:如果某个用户购买了几千件商品,双重循环的复杂度是 O(n²),单个 Reducer 可能内存溢出。课程设计的数据集通常没有这个问题,如果换了大数据集,需要加入「最多取最近 50 件商品」的截断逻辑。
为什么输出 pair 而不是分别输出 A 和 B 然后下一阶段再 join?因为共现统计只需要在同一个 key 上做 count,pair 作为 key 天然实现了聚合,省一次 shuffle 开销。
4.3 阶段二:共现转相似度,并取 TopN 邻居
阶段一输出的 (A,B) 共现次数只是原始计数,要转成真正的相似度需要分母——即每个商品被多少个用户购买。这个信息在阶段一的 Mapper 里可以顺便统计:它本来就能看到每个商品出现在多少用户的历史里。常见做法是用一个独立的 MapReduce 统计商品被购买次数,然后与共现表 join。但课程设计里有个偷懒的技巧:阶段二 Mapper 读入共现对时,顺便也读入一份「商品-出现次数表」,两个输入 join 后算余弦相似度。代码简化如下:
// CoOccurrenceMapper.java —— 只负责转发 public class CoOccurrenceMapper extends Mapper<LongWritable, Text, Text, Text> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 输入: 商品A,商品B 1 String[] fields = value.toString().split("\t"); String pair = fields[0]; String[] items = pair.split(","); // 分别以A和B为key输出,后续方便查任意方向相似度 context.write(new Text(items[0]), new Text("CO_" + items[1] + ":" + fields[1])); context.write(new Text(items[1]), new Text("CO_" + items[0] + ":" + fields[1])); } } // CoOccurrenceReducer.java —— 聚合相似度 public class CoOccurrenceReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 对于商品key,收集所有与它共现的商品及次数 Map<String, Integer> coCount = new HashMap<>(); for (Text val : values) { String v = val.toString(); v = v.substring(3); // 去掉CO_前缀 String[] parts = v.split(":"); coCount.put(parts[0], Integer.parseInt(parts[1])); } // 这里需要商品key的全局出现次数作为分母,课程设计可以用近似值(共现总数代替) for (Map.Entry<String, Integer> entry : coCount.entrySet()) { double similarity = entry.getValue() / (double) coCount.size(); context.write(new Text(key.toString()), new Text(entry.getKey() + ":" + String.format("%.4f", similarity))); } } }逻辑说明:代码里 cos 相似度的严格计算需要 |N(i)| 和 |N(j)| 两个商品各自的购买用户数,这个数据在阶段一可以额外输出。课程设计时为了避免复杂度失控,很多实现直接用「共现次数 / 共现商品总数」做近似归一化——效果略差但不会出大错。第 5 章避坑里我会说明这个近似什么时候会明显影响效果。这里的相似度阈值很重要,我一般顺手把相似度低于 0.1 的邻居过滤掉,避免稀疏矩阵带来的大量噪音邻居。Reducer 输出的 key 是商品 id,value 是「邻居商品:相似度」,供阶段三 join。
4.4 阶段三:给每个用户生成 TOP-N 推荐列表
阶段三是把用户的购买历史和商品的相似度表结合,公式是:
score(u, j) = Σ (用户u对已购商品i的兴趣 × 商品i与商品j的相似度)用户对已购商品 i 的兴趣可以用评分直接当权重,评分高代表兴趣强。MapReduce 实现上,需要同时读取两个输入:用户行为表(user_id, item_id, rating)和相似度表(item_id, neighbor:sim)。Reduce 端按用户聚合,把用户历史商品展开,去相似度表里查邻居,累加分数:
// RecommendReducer.java —— 核心:用户历史商品 × 相似度邻居 public class RecommendReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // key = userId // values 有两种来源:用户行为(userBehavior)和相似度表(itemSim) Map<String, Double> userItems = new HashMap<>(); // user已购: rating Map<String, Map<String, Double>> simTable = new HashMap<>(); // item: neighbor:sim for (Text val : values) { String v = val.toString(); if (v.startsWith("UB_")) { // 用户行为: UB_itemId,rating String[] parts = v.substring(3).split(","); userItems.put(parts[0], Double.parseDouble(parts[1])); } else if (v.startsWith("IS_")) { // 相似度: IS_neighbor:sim String[] parts = v.substring(3).split(":"); simTable.computeIfAbsent(key.toString(), k -> new HashMap<>()) .put(parts[0], Double.parseDouble(parts[1])); } } // 累加每个邻居商品的推荐分 Map<String, Double> scores = new HashMap<>(); for (Map.Entry<String, Double> entry : userItems.entrySet()) { String itemId = entry.getKey(); double rating = entry.getValue(); Map<String, Double> neighbors = simTable.get(itemId); if (neighbors != null) { for (Map.Entry<String, Double> nb : neighbors.entrySet()) { scores.merge(nb.getKey(), rating * nb.getValue(), Double::sum); } } } // 排序取前10个,输出user_id,item_id,score scores.entrySet().stream() .sorted(Map.Entry.<String, Double>comparingByValue().reversed()) .limit(10) .forEach(e -> { try { context.write(new Text(key.toString()), new Text(e.getKey() + "," + String.format("%.4f", e.getValue()))); } catch (Exception ex) { // 捕获异常避免单个用户数据问题导致整个task失败 } }); } }逻辑说明:Reducer 里用 values 的前缀区分两类输入,这是一个非常实用的 MapReduce 多数据源 join 技巧——不需要写复杂的 MultipleInputs 配置,只需要在 Mapper 输出时给 value 加上类型前缀。这个方案在数据量大时会有 Shuffle 数据膨胀的隐患,但课程设计的数据规模完全扛得住。排序用 Stream 的 sorted 再 limit(10) 干净利落,但注意只排了每个 Reducer 内的数据——如果你的 Reducer 多于 1 个,每个 Reducer 输出的是「该用户在所有商品中的 TOP10」,因为用户 id 作为 key,同一用户一定进同一个 Reducer,所以结果不会跨 Reducer 错乱。
跑完三个阶段后,用一条命令看最终结果:
hdfs dfs -cat /user/recommend/output/final/part-r-00000 | head -20如果输出的 item_id 看起来有规律(比如都是热门商品),说明相似度矩阵里热门商品权重过大,回到阶段二把余弦修正加上;如果输出都是同一个商品,检查阶段三的 join 逻辑是不是把相似度表读重复了。
5. 避坑手册:这份课程设计最容易翻车的五个地方
5.1 NoClassDefFoundError 与版本不匹配
现象:hadoop jar 提交后,任务一开始就抛NoClassDefFoundError: org/apache/hadoop/...,或者UnsupportedClassVersionError。
原因:本地编译 jar 用的 JDK 版本高于集群 Hadoop 运行环境的 JDK,或者 jar 里打进了不同版本的 Hadoop 依赖。更隐蔽的情况是:本地 IDE 里跑通了,但 hadoop jar 提交时用的是系统默认 JDK,路径和 Hadoop 配置里JAVA_HOME指向的不是同一个。
解决:先统一 JDK 版本,Hadoop 2.x 用 JDK 8,Hadoop 3.x 也用 JDK 8 最稳。编译时用 Maven 指定maven.compiler.source/target=1.8,打包时用maven-shade-plugin把依赖打进去但排除 Hadoop 自身:
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-shade-plugin</artifactId> <configuration> <filters> <filter> <artifact>*:*</artifact> <excludes> <exclude>org/apache/hadoop/**</exclude> </excludes> </filter> </filters> </configuration> </plugin>5.2 Container 反复被杀,任务进入死循环
现象:YARN 上任务一直显示 RUNNING,但进度不涨,日志里频繁出现Container killed by the ApplicationMaster或物理内存超限被系统 OOM。
原因:伪分布式下一台机器既要跑 NameNode 又要跑 DataNode 还要跑 NodeManager,默认的yarn.nodemanager.resource.memory-mb是 8GB,而你的机器总共只有 4GB,操作系统直接杀掉进程。另一个常见原因是阶段一 Reducer 里商品列表太长,单容器 1GB 内存装不下。
解决:把内存参数调成第 3.1 节表格里的值,注意改完要重启 NodeManager 才生效。如果 Reducer 内存还不够,加上「截断用户历史商品列表」的逻辑,比如只取评分最高的 50 个商品。检查参数是否生效用这条命令:
yarn node -list yarn node -status <节点名> | grep -i "memory"5.3 任务跑完但结果文件为空
现象:MapReduce 三个阶段都正常结束,hdfs dfs -cat看输出目录,发现 part-r-00000 文件是 0 字节。
原因:最常见的是输入路径写错——HDFS 上文件没传上去,或者路径拼写错了,Map 阶段零输入。另一个原因是阶段二里相似度阈值卡太死,所有相似度都被过滤了,输出自然空。还有一个冷门坑:输入文件有 BOM 头,第一行解析出来第一个字符是\ufeff,split 之后商品 id 前面带 BOM。
解决:先看任务计数器和日志确认 Map 输入行数,yarn logs -applicationId <appId> -log_files stdout里能看到Map input records=...。路径错就检查/user/recommend/input下文件是否真实存在;阈值问题把过滤值从 0.1 降到 0.01 试跑;BOM 头问题在清洗数据时用sed -i '1s/^\xef\xbb\xbf//'去掉。这类空结果问题照这个顺序排查,绝大多数五分钟内能定位。
5.4 推荐结果全是热门商品,没有个性化
现象:输出的 TOP-N 推荐里,不同用户的推荐列表高度重合,全是销量最高的那几个商品。
原因:这是 ItemCF 最经典的坑——相似度计算公式没有做热门商品惩罚。余弦公式的分母sqrt(|N(i)| * |N(j)|)如果被近似简化掉了,或者分母用了共现总数而不是单个商品购买数,热门商品会跟所有商品都产生高相似度。
解决:回到阶段二把完整的余弦相似度实现出来。需要额外输出每个商品的购买用户数,做法是在阶段一 Mapper 里顺便输出一份<itemId, count=1>到另一个 Reduce 做 sum。公式严格执行后,热门商品因为分母巨大,相似度会被压下去,推荐结果才会出现长尾商品。这在答辩时也是一个值得讲的点:你处理了流行度偏差问题。
5.5 Windows 与 Linux 之间文件格式带来的乱码和换行问题
现象:在 Windows 上编辑好的 CSV 上传到 HDFS 后,MapReduce 解析出的最后一行数据多了\r字符串,或者中文商品名全部变成乱码。
原因:Windows 换行是\r\n,Linux 只认\n,CSV 最后一行还会出现\r污染字段。中文乱码则是 CSV 在 Windows 下默认 ANSI 编码,而 Hadoop TextInputFormat 只认 UTF-8。
解决:上传前统一转成 UTF-8 并去掉\r,一条命令搞定:
# 转编码并清除windows换行符 iconv -f GBK -t UTF-8 raw.csv | tr -d '\r' > clean.csv更省事的办法是:在本地写完代码后用dos2unix clean.csv一把梭。这个坑虽然低级,但几乎每个 Windows 用户都踩过,值得写进实验报告的注意事项里。
6. 进阶验证与收尾技巧:用离线评测脚本证明它不是玩具
6.1 跑通不是终点:写个 50 行脚本算准确率召回率
MapReduce 输出的是推荐结果文件,评测脚本不需要用 Hadoop 跑,本地 Python 直接读输出文件即可。先把模型输出的推荐结果和测试集的真实购买行为对齐,然后算精确率和召回率。用 Python 写这个脚本要注意两点:一是 MapReduce 的输出是按用户分块但顺序不保证,需要 dict 按用户聚合;二是测试集中可能有用户没有拿到任何推荐,处理时跳过而不是报错。
# offline_eval.py import csv import sys from collections import defaultdict K = 10 def load_recommend(path): """读取MapReduce输出: user_id,item_id,score""" recs = defaultdict(list) with open(path, 'r', encoding='utf-8') as f: for line in f: parts = line.strip().split(',') if len(parts) >= 3: uid, iid, score = parts[0], parts[1], float(parts[2]) recs[uid].append((iid, score)) # 每个用户按分数倒序取前K个 for uid in recs: recs[uid].sort(key=lambda x: -x[1]) recs[uid] = [iid for iid, _ in recs[uid][:K]] return recs def load_ground_truth(path): """读取测试集: user_id,item_id""" truth = defaultdict(set) with open(path, 'r', encoding='utf-8') as f: for line in f: parts = line.strip().split(',') if len(parts) >= 2: truth[parts[0]].add(parts[1]) return truth def evaluate(rec_path, truth_path): recs = load_recommend(rec_path) truth = load_ground_truth(truth_path) total_precision = 0.0 total_recall = 0.0 user_count = 0 for uid, items in recs.items(): if uid not in truth: continue # 测试集里没有这个用户的行为,跳过 hit = len(set(items) & truth[uid]) precision = hit / K recall = hit / len(truth[uid]) if len(truth[uid]) > 0 else 0 total_precision += precision total_recall += recall user_count += 1 avg_precision = total_precision / user_count if user_count else 0 avg_recall = total_recall / user_count if user_count else 0 print(f"TOP-{K} 精确率: {avg_precision:.4f}") print(f"TOP-{K} 召回率: {avg_recall:.4f}") if __name__ == '__main__': evaluate(sys.argv[1], sys.argv[2])运行方式:
python3 offline_eval.py /user/recommend/output/final/part-r-00000 test_set.csv逻辑说明:脚本会把推荐文件里每个用户的 TOP-10 取出来,和测试集求交集。精确率和召回率都取所有用户的平均值。这里精确率会比召回率高——推荐系统本来就只给用户 10 个商品,用户的兴趣面远大于 10 个,所以召回率天然被 K 值限制。课程设计里 P@10 能上 30%、R@10 能上 20% 就算不错的基线,低于这个数先检查数据清洗和相似度公式,不要急着改算法。
6.2 答辩时值得讲的两句话:数据规模与热点变化
能把 MapReduce 跑通不算完,答辩时的讲述技巧决定这份课程设计的上限。这个项目最大的软肋是:技术栈偏老,现在工业界的推荐系统普遍用 Spark、Flink 做实时计算。但你想过没有——Hadoop 恰恰是最适合讲清楚「分布式计算原理」的教学框架,MapReduce 的 Shuffle 和排序机制是 Spark 的前身,很多概念是相通的。
答辩时把重点放在两类话术上:一是讲设计取舍,「为什么用 ItemCF 而不是 UserCF」,把第 2 章的对比讲明白,说明你理解算法的适用边界而不是背代码;二是讲数据规模和计算瓶颈,「共现矩阵在十万商品规模下是十亿级别的中间表,MapReduce 的分布式排序能撑住,但单机内存早就爆了」——这句话能让老师立刻明白你知道分布式解决的是什么问题。如果老师问 Spark 和 MapReduce 的区别,回答「MapReduce 每次 job 都要落盘读写 HDFS,Spark 基于内存的 DAG 计算能在中间结果上省掉大量 IO」,然后补一句「这个项目的共现矩阵计算转到 Spark 上只需要把三个阶段合成两个 RDD 算子」。这样既承认了不足,又展示了迁移能力,得分往往比吹嘘项目多完美更高。
最后说一个我的习惯:每次跑完推荐结果,先人工抽查 5 个用户——去 HDFS 里看看他们的历史购买记录,再对照推荐列表。通用指标再漂亮,都比不上「这个用户买了机械键盘,你推荐了鼠标垫」这种直观的解释。这个检查习惯救过我很多次,相似度公式写错但指标刚好没崩的情况,靠抽查一眼就能看出逻辑不对。希望这些经验帮你在课程设计路上少踩几个坑。
本文还有配套的精品资源,点击获取