简介:基于Hadoop的商品推荐系统课程设计项目以zip压缩包形式发布,适合大数据、分布式计算方向的高校学生以及推荐系统入门开发者学习使用,主要解决海量用户行为数据场景下商品推荐系统的工程搭建问题。压缩包共含19个文件,包括10个Java源码、5个XML配置文件,并附有pom.xml、README.md、properties配置等工程说明,整体大小仅26KB,目录结构简洁清晰,可以直接导入主流开发工具运行和调试。项目从Hadoop基础入手,涵盖用户行为数据处理、协同过滤与基于内容推荐、混合推荐策略、数据输入层、处理层、训练层与推荐服务层等多个模块,借助完整的源码与说明文档,可系统梳理分布式推荐系统的整体实现流程,理解MapReduce在用户画像与商品关联分析中的实际用途,也可在课程设计、毕业设计或大数据岗位面试准备中作为可直接改造的工程参考。目前已有3435人学习浏览,具有较高的参考热度。
1. 为什么Hadoop商品推荐系统成了课程设计的标配
很多人的Hadoop课程设计死在环境搭建,而不是推荐算法。商品推荐系统看着简单,给一堆商品算相似度、输出推荐列表;但挂上Hadoop后,要处理HDFS存储、MapReduce调度和YARN资源问题。算法本身反而只需二三十行核心逻辑。
这个标题对应大数据方向最常见的综合训练:HDFS负责存储,MapReduce负责计算,推荐系统负责业务价值。做完它,等于同时演示存储、计算和应用三层能力,比WordCount有说服力得多。适合正在选题的学生,也适合准备大数据面试、想补真实场景的工程师。
接下来的推进顺序是:先拆推荐模型,再搭Hadoop环境,跑通MapReduce协同过滤,最后给答辩前的三个调优参数。不依赖Spark和Flink,一个伪分布式Hadoop就能复现。按顺序走完,你会看到一条从CSV到推荐列表的完整链路。
2. 把商品推荐系统拆成Hadoop能算的模型
2.1 基于物品的协同过滤为什么比基于用户更合适
商品推荐系统里,协同过滤是最早被大规模落地的算法。基于用户的协同过滤先计算用户之间的相似度,把相似用户喜欢的商品推荐给目标用户;基于物品的协同过滤则反过来,先计算商品之间的相似度,再根据用户历史行为推荐相似商品。课程设计里我默认选基于物品的协同过滤(ItemCF),原因很实际:物品间相似度相对稳定,可以离线算好存起来,用户来的时候直接查表,符合MapReduce批处理的特性。
| 对比维度 | 基于用户(UserCF) | 基于物品(ItemCF) |
|---|---|---|
| 适用场景 | 新闻、社交内容,兴趣变化快 | 电商、视频,物品关系稳定 |
| 实时性 | 用户相似度算完要更新 | 物品相似度可定时离线计算 |
| 计算开销 | 用户量小时可用 | 物品量小于用户量时更划算 |
| 课程设计友好度 | 需要二次排序 | Map端拆商品对,Reduce累加 |
表格里这四点可以拿去做答辩开场。电商场景下用户数远大于商品数,物品相似度矩阵比用户相似度矩阵小得多,伪分布式下跑起来也更快。更重要的是,ItemCF的整个计算过程可以完全用MapReduce表达:输入是用户行为,输出是商品对共现次数,中间不需要一台能装下全部数据的机器。
2.2 用MapReduce算商品相似度的核心思路
假设用户行为数据有三列:用户ID、商品ID、行为值。行为值可以是1表示点击,也可以是自己定的评分。要计算商品A和B的相似度,第一步是找到所有同时行为过A和B的用户,第二步统计这样的用户有多少个。放到MapReduce里,Map阶段按用户聚合商品列表,再输出所有商品对;Reduce阶段对相同商品对求和,得到的就是共现次数。
u001,p001,1 u001,p003,1 u002,p001,1 u002,p002,1这是最原始的三行数据。第一条和第二条说明u001同时行为过p001和p003,所以p001和p003的共现次数加1;第三条和第四条说明p001和p002共现次数加1。如果同一用户对同一商品有多条行为,先按用户和商品去重,否则会把这个商品的权重人为放大。去重可以在Map阶段做,也可以在导入HDFS前用一条sort | uniq完成。
这里有个容易被忽略的点:Hadoop默认在shuffle阶段按key排序和分组。如果直接把商品对拆成两个字段输出,Reduce端拿到的是同一个"商品A"下的所有商品B,不是同一个(A,B)组合。因此课程设计里更稳妥的做法是把"A,B"拼成一个字符串当作key,Reduce端再按逗号拆分。这个细节会在第4章的代码里体现。
2.3 数据模型设计:三张表搞定需求
课程设计不需要设计复杂数仓,三张表足够。第一张是用户行为表,记录user_id、item_id、score、ts;第二张是商品信息表,记录item_id、item_name、category;第三张是推荐结果表,记录user_id、item_id、rk、score。行为表是原始输入,商品信息表用于最后把ID翻译成可读名称,结果表是最终输出。
CREATE EXTERNAL TABLE user_behavior ( user_id STRING, item_id STRING, score INT, ts BIGINT ) ROW FORMAT DELIMITED FIELDS TERMINATED BY ',' STORED AS TEXTFILE LOCATION '/user/hadoop/recommend/input';这段Hive建表语句把HDFS路径直接挂到表上,之后可以用SQL随时查行为量、商品分布,对排查数据质量非常有用。参数方面,FIELDS TERMINATED BY ','要和上传的CSV分隔符一致,LOCATION指向HDFS目录,如果路径写错,查出来的数据全是NULL。Hive只是辅助工具,没有它也能跑推荐,但有了它,答辩时能顺手展示Hadoop生态里的SQL层。
从Hadoop生态图的角度看,这套方案用到的是HDFS、YARN、MapReduce、Hive四块,它们组合起来就是一个最小的离线推荐管道。下一章先把环境搭起来,让这些组件能真正跑起来。
3. Hadoop开发环境搭建与最小数据管道
3.1 Hadoop伪分布式集群的安装与配置要点
课程设计阶段强烈建议先搭伪分布式Hadoop。伪分布式的意思是NameNode、DataNode、ResourceManager、NodeManager这些进程都跑在同一台机器上,配置文件却和真实集群完全一样,只是没有多节点。它最大的价值是让你用单机资源走通整个Hadoop集群搭建流程,答辩时可以直接说"我把集群从单机扩展到了N台,核心配置不变"。
安装顺序是JDK、SSH免密、Hadoop解压、环境变量、四个配置文件。SSH免密很多人会跳过,但start-dfs.sh会通过SSH登录本机启动DataNode,没有配好会一直卡着让输入密码。配置文件的改动集中在下面这个表里,不同版本参数名基本稳定。
| 配置文件 | 参数 | 建议值 | 说明 |
|---|---|---|---|
| core-site.xml | fs.defaultFS | hdfs://localhost:9000 | 默认文件系统地址 |
| hdfs-site.xml | dfs.replication | 1 | 伪分布式只有1个DataNode |
| yarn-site.xml | yarn.nodemanager.aux-services | mapreduce_shuffle | 指定MapReduce的Shuffle服务 |
| mapred-site.xml | mapreduce.framework.name | yarn | 让MapReduce跑在YARN上 |
如果是真实集群搭建,dfs.replication改成3,fs.defaultFS换成NameNode所在机器的IP,另外还要在slaves或workers文件里列出DataNode主机名。伪分布式最容易犯的错是把dfs.replication配成3,因为单节点放不下3份副本,会导致启动后DataNode上报异常,HDFS进入安全模式。
# 第一次启动前格式化文件系统,之后不要重复执行 hdfs namenode -format # 启动HDFS和YARN start-dfs.sh start-yarn.sh # 验证进程 jpshdfs namenode -format只执行一次,重复格式化会把NameNode的namespaceID改掉,DataNode还是旧的,最后报version mismatch。jps能看到NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager五个进程里至少四个,就说明Hadoop安装与配置成功。Web界面方面,NameNode管理页在新版本是9870端口,ResourceManager是8088端口,打开能看到节点和任务状态,这些信息在答辩现场很有用。
3.2 把商品行为数据导入HDFS
环境起来以后,第一步是把本地CSV上传到HDFS。推荐在HDFS上建一个专属于本次课程设计的目录结构,后续作业都在这个目录下输入输出,避免和系统目录混在一起。这里先做一次排序,是为了让下一章的流式Mapper能正确按用户分组,后面会看到好处。
# 按用户ID排序,保证流式Mapper能顺序处理用户 sort -t, -k1,1 behavior.csv -o behavior_sorted.csv # 创建数据目录 hdfs dfs -mkdir -p /user/hadoop/recommend/input # 上传排序后的数据 hdfs dfs -put behavior_sorted.csv /user/hadoop/recommend/input/behavior.csv # 查看文件块分布 hdfs fsck /user/hadoop/recommend/input/behavior.csv -files -blockssort -t, -k1,1指定逗号作为分隔符,按第一列用户ID排序;上传时把排序后的文件改名为behavior.csv,后续作业不用关心源文件名。-mkdir -p会递归创建目录,第一次用的时候不要省略-p。-put如果目标是目录,会把文件放到目录下;如果目标路径不存在,HDFS会尝试把它当文件名处理,所以目录必须提前建好。fsck输出里能看到文件被分成几个Block、每个Block在哪个DataNode上,这正好可以回答"HDFS底层怎么存数据"这种高频面试题。
数据导入HDFS后,用hdfs dfs -ls -R /user/hadoop/recommend/检查目录树。如果看到文件大小和本地一致,说明上传完整。注意HDFS上文件默认有用户权限,上传时用的用户名会变成文件owner,后续MapReduce作业要以同用户名执行,否则会报Permission denied。所以在集群上提交作业时,最好用统一的运行账号。
3.3 用自带示例作业验证MapReduce链路
Hadoop安装目录里自带了一组mapreduce示例jar,其中pi作业能快速验证从客户端提交到YARN到Map任务执行整条链路是否正常。不需要写一行代码就能检查环境。
hadoop jar $HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-examples-*.jar pi 2 10这个命令会启动2个Map任务,每个任务跑10个采样点,最终输出一个估算的圆周率值。关键不是数值多准,而是能看到任务进度从0%走到100%,说明ResourceManager成功分配了容器,MapReduce框架正常工作。如果这步失败,去检查YARN日志和/tmp/hadoop-yarn目录,伪分布式环境下绝大多数失败都出在权限、端口、内存这三类问题上。
配套的验证命令是yarn node -list和hdfs dfsadmin -report。前者能看到NodeManager是否注册,后者能看到NameNode管理的存储空间和DataNode状态。把这四条命令串起来,就是课程设计开题阶段的环境验收清单。
提示:
output目录不能提前创建。MapReduce作业在输出前会自动创建结果目录,如果目录已存在会直接报FileAlreadyExistsException。
4. 用MapReduce实现商品推荐系统的协同过滤核心
4.1 用Hadoop Streaming写共现矩阵作业
环境就绪后,最直接的办法是用Hadoop Streaming提交Python脚本。相比Java MapReduce,Python版本不需要编译,算法逻辑看得见摸得着,适合课程设计快速迭代;如果你需要交Java代码,思想完全一样,只需要把map和reduce逻辑搬到Java类里。我一般优先用Streaming把流程跑通,再决定要不要翻译成Java。
共现矩阵的核心逻辑在Mapper。因为第3章已经按用户ID排序,这个脚本可以流式处理:一旦发现用户切换,就说明上一个用户的全部行为已经读完,立刻输出该用户的商品两两组合,键是p001,p002而不是两个单独字段,这样才能保证shuffle时同一个商品对进同一个Reducer。
#!/usr/bin/env python import sys current_user = None items = [] for line in sys.stdin: parts = line.strip().split(",") if len(parts) < 3: continue user, item = parts[0], parts[1] if user != current_user: if current_user is not None: # 去重后两两组合商品对 unique_items = sorted(set(items)) for i in range(len(unique_items)): for j in range(i + 1, len(unique_items)): print(f"{unique_items[i]},{unique_items[j]}\t1") current_user = user items = [item] else: items.append(item) # 处理最后一个用户的分组 if current_user is not None: unique_items = sorted(set(items)) for i in range(len(unique_items)): for j in range(i + 1, len(unique_items)): print(f"{unique_items[i]},{unique_items[j]}\t1")代码里current_user用来判断用户是否切换,因为行为数据按用户排序后可以流式处理。sorted(set(items))先排序再去重,保证p001,p002和p002,p001不会变成两个键。内层两个循环里,外层从0到倒数第二个,内层从外层后一项到最后一个,这样每个商品对只输出一次。如果用户买了100个商品,这里会产生4950个组合,所以在数据预处理阶段过滤掉行为数过少的用户是一个好习惯。
Reducer端只需要把相同商品对的计数累加:
#!/usr/bin/env python import sys current_pair = None total = 0 for line in sys.stdin: parts = line.strip().split("\t") if len(parts) < 2: continue pair = parts[0] count = int(parts[1]) if pair != current_pair: if current_pair is not None: print(f"{current_pair}\t{total}") current_pair = pair total = count else: total += count if current_pair is not None: print(f"{current_pair}\t{total}")因为shuffle排序保证了相同key连续到达,所以Reducer可以用"当前键和上一个键是否相等"的方式流式累加。current_pair用于保存正在统计的商品对,一旦发现新key,就把旧的输出并把计数重新置为新key的第一次值。这里不能用内存把所有键存下来,否则数据量一大就会OutOfMemory。
提交到Hadoop的命令:
hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -input /user/hadoop/recommend/input \ -output /user/hadoop/recommend/cooccur \ -mapper "python mapper.py" \ -reducer "python reducer.py" \ -file mapper.py \ -file reducer.py-mapper和-reducer支持命令串,-file必须写,否则脚本不会被打包到分布式缓存里,节点上找不到文件。-numReduceTasks没有写,默认只启动1个Reducer;如果你的数据量到达几十万行,可以加-numReduceTasks 4让多个Reducer并行,输出结果是多个part文件,读取时用-getmerge合并。
4.2 从共现次数到相似度分数
共现次数只能说明两个物品经常一起出现,还不能直接当推荐分数。一个用户买了100次商品A,又买商品A的同时买了B,那么A-B共现次数会虚高。课程设计里常用的修正方式是余弦相似度:
similarity(A,B) = cooccur(A,B) / sqrt(count(A) * count(B))cooccur(A,B)是上一步算出的共现次数,count(A)和count(B)是商品A和B各自被行为的总次数。两种方法对比看下面这个表。
| 相似度方法 | 计算方式 | 特点 |
|---|---|---|
| 共现次数 | cooccur(A,B) | 简单但热门商品权重过高 |
| 余弦相似度 | cooccur(A,B) / sqrt(count(A) * count(B)) | 对热门商品有惩罚,课程设计推荐 |
算出相似度后,对每个用户的每个历史商品,找出与之相似度最高的K个未购买商品,按相似度排序取TopN。下面的命令把共现矩阵拉回本地,再生成推荐结果,最后放回HDFS。
# 把共现矩阵从HDFS合并拉回本地 hdfs dfs -getmerge /user/hadoop/recommend/cooccur cooccurrence.txt # 用Python脚本计算相似度并生成推荐结果 python recommend.py cooccurrence.txt behavior_sorted.csv > topn.txt # 把推荐结果放回HDFS hdfs dfs -put topn.txt /user/hadoop/recommend/result/getmerge会把part-*文件合并成一个本地文件,适合后续调试。recommend.py里读入共现矩阵,统计每个商品的总出现次数,再套用余弦公式。这个脚本不依赖集群,本地单机跑就能出结果;如果你希望整个链路都跑在Hadoop上,也可以把这一段写成第二个MapReduce作业,Map端读共现矩阵,Reduce端按商品聚合后计算相似度。
4.3 推荐结果落HDFS并用Hive校验
推荐结果以文本形式放回HDFS,每一行是用户ID、商品ID、排名、相似度,用tab分隔。为了让答辩更丰满,可以再建一张外部表指向结果目录,然后用SQL验证推荐逻辑。
CREATE EXTERNAL TABLE recommend_result ( user_id STRING, item_id STRING, rk INT, score DOUBLE ) ROW FORMAT DELIMITED FIELDS TERMINATED BY '\t' STORED AS TEXTFILE LOCATION '/user/hadoop/recommend/result'; SELECT * FROM recommend_result WHERE user_id = 'u001' ORDER BY score DESC LIMIT 10;FIELDS TERMINATED BY '\t'要严格对上topn.txt的输出分隔符。如果发现字段错位,优先检查生成文件中是否出现了空格或多余逗号。这道SQL跑通后,你就有了一个从原始CSV到推荐列表的完整数据流:HDFS存储原始数据,MapReduce产出共现矩阵,本地脚本生成推荐,Hive提供SQL查询层。答辩时把这个数据流画出来,比贴几百行代码更高效。
5. 答辩前必调的3个Hadoop参数与常见坑
5.1 数据倾斜:热门商品把Reducer压垮
推荐系统的输入数据天然倾斜:少数热门商品占据了大量行为记录。倾斜的直接表现是MapReduce作业所有任务都结束,剩下一个Reducer迟迟不完成。定位方法是用yarn logs -applicationId查看任务日志和Counter,如果某个Reducer的输入字节数远大于其他Reducer,就基本能确认是倾斜。
解决数据倾斜的第一招是加Combiner。当前Reducer做的是求和,满足交换律和结合律,完全可以用同一个脚本作为Combiner在Map端提前做部分累加,减少Shuffle数据量。在提交命令里加-combiner "python reducer.py"即可。第二招是提高Reducer数量,-numReduceTasks从1调到4或8,能让热门商品对分散到多个Reducer,但只是缓解,不能根治。答辩时能说出这两招,已经超过多数课程设计。
5.2 小文件太多:Map数量虚高
课程设计的数据集往往只有几十KB到几MB,上传时如果切分成几百个小文件,Hadoop会把每个文件当成一个Split,启动几百个Map任务,每个任务几秒钟就结束,大部分时间浪费在进程启动和容器分配上。调整这个问题的常用做法是上传前先合并本地文件,或者用hdfs dfs -put一次上传一个大文件。
如果数据源已经有很多小文件,可以先用Hadoop自带的hdfs balancer或者写一个MapReduce作业把小文件合并成SequenceFile。Streaming作业默认用TextInputFormat,无法自动合并小文件,最简单的办法是在上传前用cat *.csv > behavior.csv合成一个文件。这个技巧在真实集群中同样有效,面试官问"小文件有什么危害"时,你可以直接拿这个例子回答。
5.3 内存参数:Container kill是伪分布式最高频错误
伪分布式机器内存通常不大,默认的MapReduce内存限制却按多节点集群来设计。常见的错误是作业运行到一半,所有Map任务同时失败,日志里写着Container killed by the ApplicationMaster。解决方法是在mapred-site.xml里调整下面几个参数。
| 参数 | 默认值 | 建议值 | 作用 |
|---|---|---|---|
| mapreduce.map.memory.mb | 1024 | 2048 | 单个Map容器内存 |
| mapreduce.reduce.memory.mb | 1024 | 2048 | 单个Reduce容器内存 |
| yarn.nodemanager.resource.memory-mb | 8192 | 按机器实际内存调 | NodeManager可用内存 |
改完后重启YARN才生效。注意yarn.nodemanager.resource.memory-mb不能小于所有容器内存之和,否则容器排不上队。如果机器内存只有4GB,建议把Map和Reduce的memory.mb都保持在1024,优先减少同时运行的任务数,而不是强行调大。答辩现场调试时,用yarn logs -applicationId看具体错误,先确认是数据倾斜还是内存不足,再决定动哪个参数。
本文还有配套的精品资源,点击获取