简介:基于协同过滤算法与Hadoop的商品推荐系统完整项目资料,面向计算机相关专业在校生、教师及企业开发人员,适用于毕业设计、课程设计、项目立项演示或推荐系统入门实践。内含可运行源码、项目文档、依赖配置等全套内容,已通过运行测试并获导师认可,可直接用于学习或二次开发。资源共91个文件,以36个Java源码、39个编译后class文件为主,辅以6个XML配置、2个JAR依赖包及说明文档;压缩包约39.66MB,目录层级清晰,便于定位推荐算法实现与Hadoop部署配置。已有62人学习下载。压缩包内包含Maven工程结构、可执行JAR包、授权码及项目相关信息,可帮助理解协同过滤推荐流程和Hadoop环境下商品推荐的实现思路,适合后续按需修改扩展功能。
1. 基于协同过滤的Hadoop商品推荐系统,难的不是算法而是数据流转
基于协同过滤的Hadoop商品推荐系统,在课程设计和真实业务里都处在同一个位置:代码看起来不难,跑起来全是细节。很多人第一次拿到“源码+文档+全部资料(高分项目).zip”之类的压缩包,最想做的事是直接部署上线或照着文档复现,结果卡在评分数据怎么造、相似度怎么算、多个MapReduce阶段怎么串起来。商品推荐本质是在“用户-商品”二部图上做排序,协同过滤负责找相似关系,Hadoop负责把这层相似关系摊到多台机器上并行计算。下面按我自己实现这类项目时的顺序,把ItemCF拆成可以执行的MapReduce阶段,再补上数据倾斜、冷启动和参数调节这些作业文档里通常不会写的内容。适合正在做课程设计、实训项目,或者想离线搭一套简单推荐系统的开发者参考。
2. 协同过滤算法核心与Hadoop上的选型:把“谁和谁像”拆成可并行计算
2.1 先定UserCF还是ItemCF:商品推荐系统默认选后者
协同过滤里最常见的两个分支是UserCF和ItemCF。UserCF先算用户之间的相似度,再把相似用户买过的商品推荐过来;ItemCF先算商品之间的相似度,再按用户买过的商品推荐相似商品。商品推荐系统里我一般默认用ItemCF,原因有两条。第一,用户数量通常比商品数量高一个量级,用户相似矩阵会随用户量二次膨胀,而商品相似矩阵在十万商品这个尺度下还可以交给离线任务维护;第二,商品之间的共现关系在短周期内比较稳定,每天算一次或两天算一次都可以接受,用户兴趣却一直在变,UserCF需要更频繁地更新用户相似度。
这不代表UserCF没有用。新用户只有一两次行为时,ItemCF很难推断偏好,可以用UserCF找到行为相似的老用户,把老用户买过的东西补进候选集。我见过的高分项目里,最稳妥的方案是两条链路各自出结果,再按比例融合,而不是只选一个算法写进文档。
2.2 行为评分表与相似度公式:输入质量决定推荐上限
协同过滤不关心商品名称,只关心“用户-商品”之间的行为值。原始点击日志、订单表、搜索浏览记录都要先统一成一张评分表,常见映射如下:
| 行为 | 评分 | 使用说明 |
|---|---|---|
| 浏览/点击 | 1 | 量最大,噪声也最大 |
| 收藏/关注 | 3 | 明确兴趣,可信度中等 |
| 加购 | 4 | 转化前兆,权重高于收藏 |
| 下单/支付 | 5 | 最强正反馈,几乎无噪声 |
同一用户对同一商品出现多条行为时,不能简单取最大值。我一般保留时间最近的一条,时间相同取评分最高的一条。这样能避免用户先浏览、后下单,结果被浏览行为把评分冲淡。
评分表确定后,商品i和商品j的余弦相似度为:
sim(i,j) = 共同用户评分乘积和 / (商品i向量模长 * 商品j向量模长)
对隐式评分,余弦相似度不用做均值中心化,处理起来简单,也比较贴合“同时高评分购买”的场景。写成MapReduce伪代码时,可以拆成三步:先按用户分组,组内商品两两相乘,最后累加并除以模长。
# 演示“按用户分组 -> 商品两两组合”的相似度计算思路 def user_group_reducer(user_id, item_score_pairs): items = list(item_score_pairs) for i in range(len(items)): item_a, score_a = items[i] for j in range(i + 1, len(items)): item_b, score_b = items[j] # 统一成有序对,避免统计时出现 A#B 与 B#A 两行 pair = (item_a, item_b) if item_a < item_b else (item_b, item_a) emit(pair, score_a * score_b)这段伪代码里,同一用户的商品列表两两枚举,输出“商品对和评分乘积”。先排序再输出的原因很简单:A#B和B#A在reduce端会被当成两个key,最终相似度矩阵翻倍且不稳定。参数方面,乘积非正的对可以直接丢掉,负反馈不应该进入这个阶段。
2.3 用Hadoop而不是Spark:离线批处理作业交代起来更直观
很多做过推荐系统的人会问:Spark算协同过滤更快,为什么还要用Hadoop?关键在项目定位。课程设计和实训项目的评分数据一般在几十万到几百万行,Hadoop Streaming虽然每轮Shuffle都会落盘,但ItemCF拆成三到四个MapReduce就能跑完,不是迭代几十轮的算法。MapReduce的计算过程很容易画成流程图,答辩时能从map、shuffle、reduce各阶段把数据流转讲清楚,这是Spark的DAG没法替代的表达优势。另外,如果你之前已经把hadoop伪分布式搭建跑通,那就只需要在同一个环境里增加HDFS路径,不需要再维护一套Spark依赖。
Hadoop Streaming还可以继续用Python写mapper和reducer。项目源码如果是Java原生MR,可以留着看逻辑;自己动手改参数时,Streaming比Java重新编译快得多。下一章给一套能直接改用的Python Streaming方案。
3. 基于物品协同过滤的MapReduce落地:从评分矩阵到推荐列表
3.1 推荐主链路拆成几个MapReduce阶段
把ItemCF落到Hadoop上,不追求一个复杂Job,而是把计算拆成清晰的小阶段。我一般会拆成下面这张表的结构:
| 阶段 | 输入 | 输出 | 作用 |
|---|---|---|---|
| MR1 | ratings.tsv | user_id -> item:score,item:score | 整理用户向量 |
| MR2a | 用户向量 | itemA#itemB -> 乘积累加 | 计算共现乘积和 |
| MR2b | 用户向量 | item -> 评分平方和 | 计算余弦模长 |
| MR3 | MR2a与MR2b输出 | itemA -> itemB:sim | 得到相似度表 |
| MR4 | 用户历史 + 相似度表 | user_id -> item:score | 生成推荐列表 |
输入表每行固定为三个字段:user_id、item_id、score,分隔符统一用TAB。HDFS路径建议按天或版本管理,比如/rec/input/20250901,这样重算某天数据不会污染其他任务。
3.2 MR1与MR2a的Streaming脚本:先出共现乘积和
MR1只要把原始日志变成“user -> item:score”列表,mapper负责改格式,reducer负责按用户聚合。
#!/usr/bin/env python3 # mr1_map.py import sys for line in sys.stdin: line = line.strip() if not line: continue user_id, item_id, score = line.split("\t", 2) try: score = float(score) except ValueError: continue if score <= 0: continue print(f"{user_id}\t{item_id}:{score:.4f}")这是MR1的mapper,只做三件事:去空白、按TAB切字段、过滤掉非正评分。score保留四位小数是为了控制输出体积,Hadoop Streaming默认按行传输,行太长会拖慢Shuffle。
MR2a的核心是把用户向量内部两两组合成“商品对”。mapper输出时统一商品对的顺序,reduce端再按key累加:
#!/usr/bin/env python3 # mr2a_map.py import sys for line in sys.stdin: user = line.rstrip("\n").split("\t", 1)[0] items_part = line.rstrip("\n").split("\t", 1)[1] items = items_part.split(",") pairs = [] for i in range(len(items)): item_a, score_a = items[i].split(":") score_a = float(score_a) for j in range(i + 1, len(items)): item_b, score_b = items[j].split(":") product = score_a * float(score_b) if product <= 0: continue if item_a < item_b: pairs.append(f"{item_a}#{item_b}\t{product:.4f}") else: pairs.append(f"{item_b}#{item_a}\t{product:.4f}") for p in pairs: print(p)reducer与常见wordcount累加完全一样,只在key变化时输出一次当前商品的共现累积值。运行时用hadoop-streaming提交,命令如下:
hadoop jar /opt/hadoop/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -D mapreduce.job.reduces=8 \ -files mr1_map.py,mr1_reduce.py,mr2a_map.py,mr2a_reduce.py \ -mapper "python3 mr1_map.py" \ -reducer "python3 mr1_reduce.py" \ -input /rec/input/20250901 \ -output /rec/user_vector/20250901参数说明:mapreduce.job.reduces控制reduce并行度,数据量不大时8到16足够;-files会把本地脚本上传到各个节点的工作目录;-output目录在任务启动前不能存在,否则Hadoop会直接报错。
3.3 MR2b和MR3合并:相似度表只保留过阈值的结果
MR2b的任务是统计商品评分平方和。mapper把“user -> item:score列表”转成“item -> score*score”,reducer对同一个item累加。脚本和MR2a_reduce的结构差不多,只是聚合key从itemA#itemB变成item。
MR3负责把共现乘积和与商品模长合并成相似度。常见做法是把商品模长文件放在DistributedCache里,mapper每读一行商品对,就从缓存表里取出两个模长做除法。
#!/usr/bin/env python3 # mr3_map.py import sys SIM_THRESHOLD = 0.1 norm = {} for line in open("item_norm.txt", encoding="utf-8"): item, square = line.rstrip("\n").split("\t", 1) norm[item] = float(square) for line in sys.stdin: pair, total = line.rstrip("\n").split("\t", 1) item_a, item_b = pair.split("#", 1) denom = (norm.get(item_a, 0) ** 0.5) * (norm.get(item_b, 0) ** 0.5) if denom <= 0: continue sim = float(total) / denom if sim >= SIM_THRESHOLD: print(f"{item_a}\t{item_b}:{sim:.4f}")运行MR3前,先把模长结果拉到本地工作目录:
hadoop fs -cat /rec/item_norm/20250901/* > item_norm.txt hadoop jar /opt/hadoop/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files item_norm.txt,mr3_map.py \ -mapper "python3 mr3_map.py" \ -numReduceTasks 0 \ -input /rec/pair_sum/20250901 \ -output /rec/similarity/20250901这里SIM_THRESHOLD=0.1是初始阈值。低于0.1的商品对可以认为噪声大于信号,直接丢弃能减少下一阶段的内存开销。如果推荐结果明显不够用,优先把threshold降到0.05,而不是先加复杂算法。
3.4 MR4:给每个用户生成推荐列表并排序
MR4的mapper同时读相似度表和用户向量。相似度表放在DistributedCache里,用户向量从标准输入读。逻辑并不复杂:用户历史上每个商品都去相似度表里找TopK相似商品,把“相似度×用户对该商品的评分”作为候选得分;如果候选商品已经在用户历史里,直接跳过。
#!/usr/bin/env python3 # mr4_map.py import sys TOP_N = 20 sim_cache = {} for line in open("similarity.txt", encoding="utf-8"): item_a, rest = line.rstrip("\n").split("\t", 1) item_b, sim = rest.split(":", 1) sim_cache.setdefault(item_a, []).append((item_b, float(sim))) for line in sys.stdin: user, items_part = line.rstrip("\n").split("\t", 1) scores = {} for token in items_part.split(","): item, score = token.split(":", 1) scores[item] = float(score) user_score = {} for item, score in scores.items(): candidates = sim_cache.get(item, [])[:TOP_N] for other_item, sim in candidates: if other_item in scores: continue user_score[other_item] = user_score.get(other_item, 0) + sim * score for item, score in user_score.items(): print(f"{user}\t{item}:{score:.4f}")这段mapper没有排序,因为同一用户会出现在多个mapper里,最终文件顺序不可控。实践上我会在reducer里按用户聚合,再按分数降序取前10个输出,代码结构与MR1_reduce几乎一致,只是排序比较的字段要从冒号后截取。
运行参数上,下面几个值值得先写进脚本:
| 参数 | 推荐值 | 说明 |
|---|---|---|
| mapreduce.job.reduces | min(数据块数, 32) | reduce并行度,不是越多越好 |
| mapreduce.reduce.memory.mb | 2048 | reduce容器内存 |
| mapreduce.map.memory.mb | 1024 | map容器内存 |
| mapreduce.task.io.sort.mb | 256 | shuffle排序缓冲,大Shuffle时调大 |
这几个参数在伪分布式实验里不用改太大;等到hadoop集群搭建完成后、数据量变大时,优先调reduce内存。
4. 协同过滤在Hadoop集群上的坑:数据倾斜、冷启动与相似度更新
4.1 用户向量过长时先截断,避免组合爆炸
MR2a最怕的不是商品多,而是某一个用户行为特别多。用户向量长度是L时,内部两两组合是O(L²)。一个买了1万件商品的账号会产生约5000万对商品,直接把mapper输出撑爆。更常见的情况是刷单用户或爬虫账号,行为量比普通用户高两三个量级。
我一般会在MR2a的mapper开头做一次保护性截断:
MAX_VECTOR_LEN = 300 items = items_part.split(",") if len(items) > MAX_VECTOR_LEN: # 按评分从高到低排序,优先保留强反馈商品 items = sorted(items, key=lambda x: -float(x.split(":")[1]))[:MAX_VECTOR_LEN]这样做的代价是丢掉了部分长尾行为,但保住的是评分最高的强反馈商品,对推荐质量的影响可控。如果课程设计的评分数据本身很稀疏,MAX_VECTOR_LEN可以放宽到1000。
4.2 数据倾斜:用户热点会让单个reduce变成瓶颈
数据倾斜是Hadoop面试题里最常被问到的场景,也是协同过滤落地时最容易翻车的地方。热门商品的共现key,比如“手机#充电器”,会被海量用户同时命中,某个reduce处理的数据量可能是其他reduce的几十倍。
常用解法是两阶段聚合加随机盐。在MR2a输出的商品对后面拼一个随机后缀,让同一对商品先分散到不同reduce做局部累加,然后再用一个MapReduce去掉盐,做全局累加。
# 在原来的pair key后面加随机盐 import random salt = random.randint(0, 7) print(f"{item_a}#{item_b}#{salt}\t{product:.4f}")第一阶段的reduce输出itemA#itemB#salt -> 局部和,第二阶段mapper读入时把#salt剥掉,只输出itemA#itemB -> 局部和,第二阶段reduce再按原始key累加,得到完整乘积和。盐的范围取5到8即可,太大会让Shuffle数据量成倍上涨。如果加了盐还是慢,就把热门商品和非热门商品拆成两个任务分别计算,最后合并。
4.3 冷启动兜底:离线Hadoop表加实时缓存
Hadoop按天产出,天然解决不了用户刚刚发生的点击。真正的工程方案是“离线算相似度,在线做召回”。ItemCF的相似度表每天更新,用户最近行为放到Redis,在线服务用用户实时向量去乘当天的相似度表,秒级返回。Hadoop任务只负责定期生成相似度表和全量用户的预计算结果。
| 场景 | 处理方式 |
|---|---|
| 新用户无行为 | 推全局热门TopN |
| 新商品无评分 | 用类目下的热门商品替换,或等积累5个以上行为后再进候选 |
| 行为稀疏用户 | 降低相似度阈值,把候选集扩大2到3倍 |
全局热门表可以单独写一个hot_item任务,输入评分表,reduce端按商品得分累加后取TopN。这个表要输出到独立HDFS路径,方便在线服务每天拉取。
4.4 参数怎么调:阈值、TopN、时间衰减
协同过滤参数不多,但每个参数都直接影响推荐列表的“密度”。常见初始值如下:
| 参数 | 初始值 | 调整方向 |
|---|---|---|
| 相似度阈值 | 0.1 | 调高提升精度,调低提升召回 |
| 每商品TopN | 20 | 调大覆盖更多长尾,但增加内存压力 |
| 共同用户数下限 | 2 | 过滤“恰好同一个人买过”的偶然对 |
| 行为评分权重 | 浏览1、加购4、下单5 | 可按业务调整加购和下单差距 |
如果日志里有行为发生时间,我还会给评分加一个时间衰减,让一个月前的行为权重下降:
score = round(base_score * 0.9 ** days_since, 4)这里的0.9是按天衰减的系数,想衰减快一点就改0.8,慢一点改0.95。注意衰减后score变成小数,MR2a比较score <= 0要保留,防止衰减后仍为负的异常行混进去。
提示:调参时不要同时改多个参数。先固定相似度阈值为0.1,把TopN从20调到50,观察推荐结果变化;再单独调阈值。否则最后说不清是哪个参数贡献了效果。
5. 从zip源码包到集群:伪分布式验收与推荐结果验证
5.1 先用玩具数据集跑通完整链路
拿到一个项目压缩包,不建议立刻上全量数据。先做一个没有任何歧义的小文件ratings.tsv,放进HDFS,确认五个MapReduce阶段都能跑出预期结果。
hdfs dfs -mkdir -p /rec/input hdfs dfs -put ratings.tsv /rec/input/玩具数据建议只写三个用户、六个商品,覆盖“两个用户共同买过A和B”“另一个用户只买过C”这两类情况。只要最终推荐结果里,买过A的用户被推荐了B,说明共现计算和相似度拼接方向是对的。
5.2 用awk和sort检查输出结果
MR4输出完成后,先看整体行数,再抽查指定用户。
hadoop fs -cat /rec/recommend/20250901/* | wc -l # 查看u3用户推荐结果,按得分降序取前5 hadoop fs -cat /rec/recommend/20250901/* | awk -F '\t' '$1=="u3" {print $2}' | sort -t: -k2 -rn | head -5输出格式是user_id \t item_id:score,所以第一条命令统计总行数能判断是否有用户完全没拿到结果。第二条是日常排查最常用的命令,awk先筛用户,sort按冒号后面的得分降序排。
推荐系统还要检查“是不是重复推荐了同一个商品”。MapReduce的多路输入、相似度表去重不彻底都会造成重复,可以用下面这段快速扫一遍:
hadoop fs -cat /rec/recommend/20250901/* | awk -F '\t' '{split($2,a,":"); if (n[$1"#"a[1]]++) print "duplicate:"$0; if (a[2]+0>5.0) print "overflow:"$0}' | head -10这里有三个检查点:一是n[$1"#"a[1]]统计同一个用户是否出现两次同一商品;二是a[2]+0>5.0检查得分是否超过评分上限乘以相似度上限。如果分数冲到5以上,多半是MR4里没有排除用户历史商品,或相似度表里有自环。
5.3 导出到MySQL前的最后一步
在线推荐不会直接读HDFS,通常要把结果导到MySQL或Redis。HDFS输出的是多个part文件,直接逐个下载很麻烦,hadoop自带合并命令:
hadoop fs -getmerge /rec/recommend/20250901 recommend.tsv mysql -u rec_user -p rec_db -e "LOAD DATA LOCAL INFILE '/path/recommend.tsv' INTO TABLE t_recommend(user_id,item_id,score)"getmerge会把几十个part文件按名称顺序拼成一个本地文件,再用LOAD DATA LOCAL INFILE导入。导入前记得清空当天分区表,避免任务重跑后出现重复数据。把getmerge写进run.sh脚本里,后面排查数据质量问题会顺很多。
本文还有配套的精品资源,点击获取