Hadoop商品推荐系统实战:离线协同过滤与MapReduce实现

发布时间:2026/10/9 21:29:53
Hadoop商品推荐系统实战:离线协同过滤与MapReduce实现 简介基于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 工程核心代码集中在四个包下。拿到项目包以后第一步不是急着跑命令而是把类名和依赖关系理清楚。包路径类名职责modelUserBehaviorWritable用户行为记录的结构化对象实现 Writable 接口similarityItemCooccurrenceMapper把同一用户购买的商品两两组合输出商品对similarityItemCooccurrenceReducer累加商品对同现次数输出同现矩阵scoreRatingMapper读取同现矩阵和用户行为计算预测评分scoreTopNReducer对每个用户的候选商品按评分排序取前 N 个jobDriver组装两个作业设置输入输出路径和参数代码里依赖只有一个 hadoop-client版本跟着集群走。拿到包以后先改 pom 里的 Hadoop 版本号改成和你集群一致的版本不然提交作业时容易报依赖冲突。3.2 第一步生成模拟数据并上传 HDFS项目里没有提供现成的大规模行为数据时先用脚本生成一份结构正确的模拟数据。下面这个 Python 脚本生成 5 万条行为记录覆盖 500 个用户和 200 个商品评分分布偏向高分段模拟真实用户更愿意给好评的行为。import random import csv users [fu{str(i).zfill(4)} for i in range(1, 501)] items [fp{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 MapperLongWritable, Text, Text, IntWritable { private MapString, ListString 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 (ListString 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.reduces16 \ -D mapreduce.map.memory.mb2048 \ com.example.recommend.job.Driver \ /recommend/input \ /recommend/output/similarity \ /recommend/output/score这里有个新手最容易看漏的地方-D 参数必须写在 jar 包和主类之间放在命令末尾会被当作 main 方法收到的普通参数传给 Driver导致参数不生效。遇到过不止一个同事在这个地方栽过作业跑起来后 Reduce 数量永远是默认的 1数据一倾斜直接卡死。作业参数按下面的建议值调能避开大部分性能问题参数默认值建议值说明mapreduce.job.reduces18 到 16超过集群可用核数反而浪费调度时间mapreduce.map.memory.mb10242048行为数据量大时 map 容易 OOMmapreduce.reduce.memory.mb10244096同现矩阵聚合开销大默认值不够mapreduce.map.speculativetruefalse数据倾斜场景下推测执行会导致重复计算提示修改 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_k10, total_item_count200): 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.lambda0.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% 说明参数合适。从那以后我每次交付推荐项目都会先拿历史一周的日志做回放观察推荐结果的空窗率和指标波动确认稳定后再交给定时调度不然上线第一天就会被新的数据波动打懵。这套流程虽然多花半小时但能省掉后面几天的排查时间。希望帮到你。本文还有配套的精品资源点击获取

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询