简介:基于Hadoop的商品推荐系统是一份面向大数据开发者与算法初学者的完整项目资源,重点解决用户行为数据采集、清洗与个性化推荐落地问题。资源依托HDFS集群与MapReduce计算模型,通过Step1至Step6共6个MapReduce作业串起完整数据流水线,覆盖数据清洗、用户评分统计、商品相似度计算与推荐结果生成等关键环节,帮助理解分布式环境下推荐系统的工程化流程。资料包共16个文件,包含7个Java源码文件、2个XML配置、2个Eclipse偏好设置、1个project工程文件、1个说明文档及1个CSV样本数据等,整体压缩包仅94KB,小巧但结构清晰,便于快速部署与阅读。docx说明文档与CSV示例数据可直接对照学习,配置文件和Java类覆盖了从环境搭建到结果输出的关键环节。目前已有5307人学习下载,适合用于课程设计、毕业设计或Hadoop入门实战。
1. 基于 Hadoop 的商品推荐系统:离线批量推荐为何仍是中小平台的最优解
在大数据项目里,Hadoop 经常被说成“重”和“慢”,一提到推荐系统,大家下意识想到实时计算。实际上,绝大多数商品推荐场景根本不需要秒级更新:白天攒行为日志,夜间跑批,凌晨把 TopN 列表推给业务库,就够了。这个基于 Hadoop 的商品推荐系统项目,覆盖从用户行为清洗、物品协同过滤、相似度计算到 TopN 结果落库的完整链路,整个工程依赖少、部署简单,适合做大数据的在校学生打通分布式编程模型,也适合小团队的后端工程师在现有 Hadoop 集群上快速接一条离线推荐管线。
2. 选型与数据链路:HDFS 存什么、MapReduce 算什么、Hive 清洗什么
2.1 算法选型:基于物品的协同过滤为什么更适合 Hadoop 批处理
推荐算法有很多种,这个项目选的是基于物品的协同过滤。核心逻辑一句话:如果大量用户同时买过商品 A 和商品 B,那么 A 和 B 是相似的;用户买过 A,就把和 A 相似的 B 推给他。和基于用户的协同过滤相比,ItemCF 在电商场景里有一个很现实的优势:商品数量通常比用户数量低一个数量级,算出来的相似度矩阵是“商品 × 商品”,规模可控。用户量涨到几十万以后,UserCF 的用户相似度矩阵几乎存不下,小时级跑批也扛不住。
ItemCF 的相似度计算可以用余弦公式表示:sim(i, j) 等于同时评价过商品 i 和 j 的用户评分之积,除以各自评分向量模长的乘积。落到 MapReduce 上,整个过程分成两步——第一步统计商品共现次数,两个商品被同一个用户买过,就记一次同现;第二步把共现次数归一化成相似度,存回 HDFS。在线推荐时,只需要读取每个商品的 TopN 相似列表,再按用户历史上买过的商品做一次加权求和,就能得到候选推荐列表。这个链路天然适合离线批量计算,这也是它放在 Hadoop 上最顺的原因。
2.2 数据模型:行为表、评分表和 Hive 分区表怎么设计
数据模型决定了下游算得顺不顺。这个项目里最核心的表是用户行为表,字段只有四个:user_id、item_id、rating、ts。rating 不是简单的 1 或 0,而是 1 到 5 的评分,这样相似度计算能用到分值信息,而不是只统计有没有共现。ts 存 Unix 时间戳,不用字符串日期,因为后续做时间衰减时要直接参与运算,字符串解析一次就是浪费一轮 Map 周期。
CREATE EXTERNAL TABLE IF NOT EXISTS dwd_user_behavior ( user_id STRING COMMENT '用户ID', item_id STRING COMMENT '商品ID', rating TINYINT COMMENT '行为评分 1-5', ts BIGINT COMMENT '行为时间戳' ) PARTITIONED BY (dt STRING COMMENT '日期分区') STORED AS ORC LOCATION '/warehouse/dwd_user_behavior';字段说明:rating 用 TINYINT 而不是 INT,评分区间 1 到 5,一个字节足够,列存底下能省不少空间;ts 用 BIGINT 是为了后续在 Hive SQL 里做时间衰减时可以直接比较和运算。分区字段 dt 按天挂,每天跑批只需要读一个分区,避免每次全表扫描。ORC 是列式存储,压缩率高,读取评分列时只扫需要的列。
这里有一个很多课程项目都不会讲的细节:原始日志和清洗后的表最好分开。原始日志放在一个不落分区的目录里,清洗时先过滤异常数据、去重、做时间衰减,再写入 dwd_user_behavior。这样相似度作业读到的数据永远是干净的,排查问题时也能回到原始目录对账。
2.3 数据流转链路:一天的行为日志如何变成推荐列表
数据在 Hadoop 上转一圈,环节比实时推荐多,但每一步职责都很清楚。
第一步,前端埋点或者业务库导出用户行为日志,落到 HDFS 的原始目录。第二步,Hive 跑定时清洗任务,做去重、过滤异常评分、按时间窗口截取数据。第三步,MapReduce 作业一读清洗后的数据,统计商品共现矩阵。第四步,MapReduce 作业二读共现矩阵,计算相似度,生成每个用户的 TopN 候选列表。第五步,结果写回 HDFS 的结果目录,用定时导出任务同步到业务数据库。第六步,线上推荐服务缓存 TopN 列表,按用户请求组装推荐位数据。
这套链路里 HDFS 既是输入也是输出,承担存储职责;MapReduce 做两次核心计算;Hive 做 ETL 清洗;ZooKeeper 负责集群协调。如果业务要求推荐位每天更新一次,这套链路跑完正好赶上早高峰前上线。至于为什么不用 Spark,如果你的集群本来就只有 Hadoop 发行版,再引一套 Spark 依赖,运维成本直接翻倍,离线场景下小时级延迟完全能接受,MapReduce 稳定、好排查,够用了。
3. 从源码到集群:搭建商品推荐系统的完整落地步骤
3.1 项目结构与模块职责
先看工程结构,整个项目是一个标准 Maven 工程,核心代码集中在四个包下。拿到项目包以后,第一步不是急着跑命令,而是把类名和依赖关系理清楚。
| 包路径 | 类名 | 职责 |
|---|---|---|
| model | UserBehaviorWritable | 用户行为记录的结构化对象,实现 Writable 接口 |
| similarity | ItemCooccurrenceMapper | 把同一用户购买的商品两两组合,输出商品对 |
| similarity | ItemCooccurrenceReducer | 累加商品对同现次数,输出同现矩阵 |
| score | RatingMapper | 读取同现矩阵和用户行为,计算预测评分 |
| score | TopNReducer | 对每个用户的候选商品按评分排序,取前 N 个 |
| job | Driver | 组装两个作业,设置输入输出路径和参数 |
代码里依赖只有一个 hadoop-client,版本跟着集群走。拿到包以后先改 pom 里的 Hadoop 版本号,改成和你集群一致的版本,不然提交作业时容易报依赖冲突。
3.2 第一步:生成模拟数据并上传 HDFS
项目里没有提供现成的大规模行为数据时,先用脚本生成一份结构正确的模拟数据。下面这个 Python 脚本生成 5 万条行为记录,覆盖 500 个用户和 200 个商品,评分分布偏向高分段,模拟真实用户更愿意给好评的行为。
import random import csv users = [f"u{str(i).zfill(4)}" for i in range(1, 501)] items = [f"p{str(i).zfill(4)}" for i in range(1, 201)] with open("user_behavior.csv", "w", newline="") as f: writer = csv.writer(f) writer.writerow(["user_id", "item_id", "rating", "ts"]) random.seed(42) for i in range(50000): user = random.choice(users) item = random.choice(items) rating = random.choice([1, 2, 3, 3, 4, 4, 5, 5, 5]) ts = str(1609430400 + i * 327) writer.writerow([user, item, rating, ts])脚本逻辑:random.seed(42) 保证每次生成的数据完全一致,方便复现;评分用 random.choice 带权重,5 分出现的概率最高,低分是少数;时间戳从 2021 年开始递增。跑完以后用下面命令上传:
hdfs dfs -mkdir -p /recommend/input hdfs dfs -put user_behavior.csv /recommend/input/ hdfs dfs -cat /recommend/input/user_behavior.csv | head上传前先确认 HDFS 目录存在,put 之后用 cat 检查文件前几行,确认没有空行和乱码再继续。这一步虽然基础,但跳过检查直接跑作业,后面出错时你很难判断是数据问题还是代码问题。
3.3 第二步:理解相似度计算的 MapReduce 核心代码
作业一的核心是 ItemCooccurrenceMapper。它的思路是:每一个用户的购买记录里,任意两个商品组成一对,作为 key 输出,value 固定为 1。同一个商品对如果被多个用户共现过,就会在 Reduce 阶段被累加,得到共现次数。
public class ItemCooccurrenceMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private Map<String, List<String>> userItemsCache = new HashMap<>(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields = value.toString().split(","); if (fields.length < 3) return; String userId = fields[0]; String itemId = fields[1]; userItemsCache.computeIfAbsent(userId, k -> new ArrayList<>()).add(itemId); } @Override protected void cleanup(Context context) throws IOException, InterruptedException { for (List<String> items : userItemsCache.values()) { for (int i = 0; i < items.size(); i++) { for (int j = i + 1; j < items.size(); j++) { context.write(new Text(items.get(i) + ":" + items.get(j)), new IntWritable(1)); } } } } }代码逻辑:map 阶段只做数据读取和缓存,不输出中间结果;cleanup 阶段在每个 Mapper 完成任务前,把缓存的用户商品列表做两两组合输出。这种写法适合课程项目和中小数据量,能显著减少 Shuffle 阶段的数据量。生产环境更严谨的做法是先按 user_id 排序分组,保证同一个用户的记录被同一个 Mapper 完整读取,但项目里数据量不大时,缓存方案跑起来更简单。
Reducer 侧反而是最简单的部分,对 key 相同的 value 做累加。它不需要太复杂的逻辑,真正的优化点应该在 Combiner——在 map 端先做一轮本地累加,减少写入磁盘的中间结果。数据量一大,Combiner 能省下将近一半的 Shuffle 开销。
3.4 第三步:编译打包、提交 YARN 与参数调优
代码改好以后,直接打包提交。下面是完整命令:
mvn clean package -DskipTests hadoop jar target/recommend-1.0.jar \ -D mapreduce.job.reduces=16 \ -D mapreduce.map.memory.mb=2048 \ com.example.recommend.job.Driver \ /recommend/input \ /recommend/output/similarity \ /recommend/output/score这里有个新手最容易看漏的地方:-D 参数必须写在 jar 包和主类之间,放在命令末尾会被当作 main 方法收到的普通参数传给 Driver,导致参数不生效。遇到过不止一个同事在这个地方栽过,作业跑起来后 Reduce 数量永远是默认的 1,数据一倾斜直接卡死。
作业参数按下面的建议值调,能避开大部分性能问题:
| 参数 | 默认值 | 建议值 | 说明 |
|---|---|---|---|
| mapreduce.job.reduces | 1 | 8 到 16 | 超过集群可用核数反而浪费调度时间 |
| mapreduce.map.memory.mb | 1024 | 2048 | 行为数据量大时 map 容易 OOM |
| mapreduce.reduce.memory.mb | 1024 | 4096 | 同现矩阵聚合开销大,默认值不够 |
| mapreduce.map.speculative | true | false | 数据倾斜场景下推测执行会导致重复计算 |
提示:修改 reduce.memory.mb 时,同时要保证集群里每个节点可用内存足够,不然容器会一直等待调度,表现为作业提交成功但迟迟跑不起来。
提交完作业以后,用yarn application -status <application_id>查看进度。如果某个 reduce 长时间卡在 99%,不要急着 kill,先去第 4 章的避坑记录里对照排查。
4. 避坑记录:五个影响推荐结果的真实问题
推荐系统容易翻车的点,九成不在算法,在数据。下面五个坑来自不同项目里的真实排查经历,每一条都按现象、原因、解决的方式整理,希望对得上号。
4.1 数据倾斜:一个 Reduce 卡了几个小时
现象:同现矩阵作业里,99% 的 Reduce 几十秒跑完,唯独有一个 Reduce 跑了两个小时还在 99%。打开 YARN 日志,发现某个 key 处理的记录数是其他 key 的几百倍。
原因:热门商品和任何商品都有共现。比如一个爆款手机壳,几乎每个用户都买过,它和其他几百个商品组成商品对,全部被 hash 到同一个 Reduce,单个 key 的记录数直接爆炸。
解决:三个手段配合使用。第一,在 Mapper 的 cleanup 里对同一个用户的商品列表先去重,同一个用户重复购买同一商品不重复计共现;第二,加 Combiner 做本地累加,让 Shuffle 数据量降下来;第三,对热门商品做降权,相似度公式改为score = raw_count / (1 + sqrt(hot_i * hot_j)),热门商品对之间的相似度会被压下去,长尾商品才有机会浮上来。
4.2 评分矩阵稀疏:推荐结果全空
现象:把几千条测试数据喂进去,跑完作业后输出目录里只有几十行结果,大部分用户的推荐列表是空的。第一反应是代码写错了,翻了一天代码,最后发现算法逻辑没问题。
原因:行为数据太少,用户之间几乎没有共同的商品购买记录,同现矩阵本身就很稀疏,算出来的相似度大部分还是 0。冷启动阶段的推荐系统,稀疏矩阵是绕不开的问题。
解决:三个手段并用。第一,在生成候选集时设置一个极小值兜底,相似度为 0 的商品对赋一个 0.01 的平滑值,避免向量全零;第二,清洗时不要把数据源局限在“购买”行为,点击、收藏、加购都纳入进来,只是权重不同;第三,如果 TopN 列表还是为空,直接按商品热度排序填充,宁可用热门商品占位,也不让推荐位空着。
4.3 时间窗口太宽:推荐永远在推“昨天”
现象:用户两周前买了跑步鞋,这周打开首页,推荐位第一位还是跑步鞋。业务方跑来问是不是集群没更新,其实作业每天都在跑,是逻辑层面的问题。
原因:相似度计算把所有历史行为当成同等权重,没有时间衰减。半年前的数据对今天的推荐还在起作用,新商品永远没有机会进入相似列表。
解决:在 Hive 清洗层加时间衰减,超过 90 天的行为直接过滤,窗口内的行为按衰减系数降权:
SELECT user_id, item_id, rating * POW(0.95, DATEDIFF(CURRENT_DATE, FROM_UNIXTIME(ts, 'yyyy-MM-dd'))) AS decayed_rating FROM dwd_behavior_raw WHERE DATEDIFF(CURRENT_DATE, FROM_UNIXTIME(ts, 'yyyy-MM-dd')) <= 90;POW(0.95, n) 让消息按天指数衰减,90 天前的数据权重已经接近 0,可以直接截断。另外在生成最终推荐结果时,把用户已经购买过的商品排除掉,这是很多人会漏掉的一步。
4.4 HDFS 小文件:拖慢的其实是作业初始化
现象:几百万条行为数据,输入切片却有几百个,Map Task 数量上千,作业启动花了几分钟,NameNode 的 GC 时间也明显变长。
原因:上传日志时把一天的数据拆成了几百个小文件,HDFS 里每个文件都要占用元数据内存,每个小文件至少产生一个 Map 切片。小文件问题是 Hadoop 离线任务里最容易被忽视的性能杀手。
解决:上传前先合并小文件。用 getmerge 把多个小文件合并成一个大文件再上传:
hdfs dfs -getmerge /recommend/input_raw /tmp/all_logs.csv hdfs dfs -put /tmp/all_logs.csv /recommend/input/combined.csv如果数据源还在持续产生小文件,可以在 Hadoop 配置里打开合并输入格式,让多个小文件共享一个切片,减少 Map 数量。建议数据文件不超过 128MB 的切割阈值时就先合并。
4.5 本地跑通、集群翻车:环境差异排查
现象:在 IDE 里用本地模式跑数据,结果完全正常;打包提交到集群后,输出文件缺了一截,偶尔还报 Container killed,日志里是内存溢出。
原因:本地模式用的是 LocalJobRunner,不真正走 YARN 容器,也没有严格的内存限制;集群模式下容器内存受限,代码里如果初始化了过大的堆内存,或者中间结果写入磁盘太多,很容易触到容器上限。另一个常见问题是本地模式默认只有一个 Reduce,数据倾斜在本地根本暴露不出来。
解决:提交集群前先跑一个小数据集,确认路径、权限、依赖都没问题。重点检查两处:一是代码里不要写死本地文件路径,统一从 Driver 的参数读取;二是看一下提交命令里的 reduce.memory.mb 和容器内存是否匹配。用yarn logs -applicationId <app_id>查具体报错,不要只看 Container killed 就盲目加内存。
5. 验证与增量更新:让推荐系统从“能跑”到“可信”
5.1 离线评估:精确率、召回率、覆盖率
作业跑完不代表推荐做完了。推荐结果好不好,需要用离线指标量化。常见做法是把行为日志按时间切分,前 80% 作为训练集,后 20% 作为测试集,然后对比推荐列表和用户真实产生的行为,计算三个指标:精确率、召回率、覆盖率。
def evaluate(reco_result, test_data, top_k=10, total_item_count=200): hit = 0 total_precision = 0.0 total_recall = 0.0 all_reco_items = set() for user, reco_list in reco_result.items(): reco_items = reco_list[:top_k] test_items = test_data.get(user, set()) hit_count = len(set(reco_items) & test_items) hit += hit_count total_precision += hit_count / top_k total_recall += hit_count / len(test_items) if test_items else 0 all_reco_items.update(reco_items) precision = total_precision / len(reco_result) recall = total_recall / len(reco_result) coverage = len(all_reco_items) / total_item_count return precision, recall, coverage代码逻辑:精确率算的是推荐列表里有多少是用户真实点击过的;召回率算的是用户真实点击过的商品有多少被推荐出来了;覆盖率反映推荐系统是不是只集中在热门商品上。三个指标合在一起看,才能判断推荐质量,单个指标高没有意义——精确率很高但覆盖率很低,说明系统只推爆款,个性化基本没生效。
5.2 进阶技巧:增量相似度更新,避免天天全量重算
全量重算的逻辑很简单,每天把历史所有行为重新读一遍,缺点是越往后数据量越大,跑批时间越来越长。更常见的做法是增量更新:每天只计算当天新增行为涉及的商品对,和历史相似度矩阵做加权合并。
# 每日只算当天新增行为对应的商品对 hadoop jar target/recommend-1.0.jar \ com.example.recommend.job.DailyIncrement \ /recommend/input/$(date +%Y%m%d) \ /recommend/output/increment_sim # 把增量结果和存量结果合并,历史权重取 0.8 hadoop jar target/recommend-1.0.jar \ com.example.recommend.job.IncrementMerge \ -D merge.lambda=0.8 \ /recommend/output/history_sim \ /recommend/output/increment_sim \ /recommend/output/merged_sim新相似度的计算公式是new_sim = 0.8 * history_sim + 0.2 * increment_sim。lambda 越大,历史信息保留得越多,推荐结果越稳定;lambda 越小,新行为影响权重越高,热点反应越快。一般从 0.8 开始调,观察评估指标的波动幅度,每天波动不超过 5% 说明参数合适。
从那以后,我每次交付推荐项目,都会先拿历史一周的日志做回放,观察推荐结果的空窗率和指标波动,确认稳定后再交给定时调度,不然上线第一天就会被新的数据波动打懵。这套流程虽然多花半小时,但能省掉后面几天的排查时间。希望帮到你。
本文还有配套的精品资源,点击获取