简介:基于Hadoop的商品推荐系统完整Java工程,适合课程设计、毕业设计或自学参考,主要面向大数据初学者与电商推荐系统开发者,用来解决在分布式环境下处理海量用户行为数据并生成个性化推荐的问题。压缩包内含Maven工程结构,共34个文件,包括29个Java源文件与5个XML配置,包体仅25KB,代码紧凑清晰,适合直接导入IDE阅读和运行。目前已有193人学习下载,可结合HDFS存储、MapReduce并行计算、协同过滤等知识点进行源码级对照学习。Java源码覆盖用户行为数据清洗、购买记录统计、用户相似度与物品相似度计算、TopN推荐列表生成等核心环节,并通过Hadoop作业配置XML文件关联输入输出路径及运行参数;项目还体现了基于内容的推荐与混合推荐的实现思路。对想掌握Hadoop生态实战、快速搭建推荐系统原型的读者,这份工程提供了从数据预处理到推荐结果输出的完整参考链路,也能帮助理解MapReduce在真实推荐场景中的落地方式。
1. 基于Hadoop的商品推荐系统:先把规模和边界说清楚
当一条埋点日志一天产出几个GB,单机内存已经装不下协同过滤要用的全量评分矩阵时,你以为的“加内存”其实只是推迟了问题。这套基于Hadoop的商品推荐系统,解决的是规模化之后的事:把用户行为日志落到HDFS,用Java写MapReduce清洗数据、算相似度、生成推荐列表,跑完整条离线推荐链路。它不是一个在线实时推荐引擎,而是一套经过实战拆解的批处理项目,拿到手能改、能跑、能往课程设计或简历里写。
项目工程名GRMS-master,压缩包打开是标准Maven骨架,pom.xml、src/main、src/test都在。核心代码分三块:数据预处理、协同过滤算法、推荐结果生成。适合两类人——课程设计需要“大数据+推荐系统”标签的学生,以及刚接触Hadoop生态、想在一个完整业务场景里把MapReduce流程走一遍的Java开发。下文按Hadoop的存储与计算角色、算法拆解、环境复现、踩坑记录、优化方向展开,全程对着这份代码讲。
2. Hadoop在推荐系统里的角色:从HDFS存储到MapReduce计算边界
2.1 为什么是Hadoop而不是一台大内存机器
推荐系统的数据链路,归根结底是“读日志 → 算相似度 → 出推荐列表”。当用户量在十万级、商品在万级时,相似度矩阵几十亿条记录,单机用HashMap勉强能塞;但一旦输入是原始埋点日志,问题就变了。日志自带脏数据,字段缺失、无效点击、爬虫刷量都混在一起,单机处理时要么把整批数据load进内存然后OOM,要么自己写多线程逐行扫描,代码复杂度一下子失控。
Hadoop在这里的核心价值,不是让推荐算法变得更快,而是把“数据容量上限”这件事变得不焦虑。HDFS把大文件切分成128MB的块,散到多个DataNode上,每个块默认三副本,单节点宕机不丢数据。MapReduce把计算逻辑推送到数据所在的节点执行,“移动计算而不是移动数据”,框架自动调度Mapper读取split,shuffle后交给Reducer聚合,你只需要实现map和reduce两个方法,不需要管线程池、不需要管分布式文件锁。
但边界必须说清楚:MapReduce的批处理延迟是分钟级的,不适合“用户刚下完单,下一秒推荐列表就要刷新”的场景。这种实时需求要交给Spark Streaming或者Flink,而不是Hadoop MapReduce。这套项目里的定位就是每天凌晨跑一轮全量重算,产出当天推荐,这在很多中小电商场景里完全够用。
2.2 GRMS项目结构:pom.xml与src/main下的工程骨架
解压zip后第一眼看到的是GRMS-master目录,标准的Maven结构。根目录下pom.xml声明依赖和打包方式,src/main/java下面是主代码,src/main/resources放日志和配置文件。我拆包时最关注的是pom.xml里那几个坐标,直接决定了作业能不能提交到你本地的Hadoop集群上。
| 文件/目录 | 作用 |
|---|---|
| pom.xml | Maven项目描述,定义hadoop-client、junit等依赖版本 |
| src/main/java | 主代码目录,Mapper/Reducer/Driver都在这里 |
| src/main/resources | 配置文件目录,如log4j.properties |
| src/test/java | 单元测试,针对清洗逻辑和相似度函数做本地验证 |
pom.xml里值得注意的依赖项,整理成一张表:
| 依赖 | 典型版本区间 | 用途 |
|---|---|---|
| hadoop-client | 2.7.x或3.x | 提供HDFS、MapReduce、YARN客户端API |
| junit | 4.x | 本地跑通清洗和相似度逻辑测试 |
| commons-lang3 | 3.x | 字符串处理,处理埋点字段转义和截断 |
依赖版本有个大坑:hadoop-client的版本必须和你集群安装的Hadoop大版本对齐。本地用3.3.6打包,集群还是2.7.5,提交作业会直接报Protocol版本不一致,这个我在第5.3节详细讲。
2.3 数据预处理:从原始埋点日志到评分矩阵
推荐系统的起点是数据,但原始日志不能直接喂算法,里面有无效点击、字段缺失、重复上报。这套项目的第一步MapReduce就是清洗和评分转化。先看典型的清洗Mapper:
public class LogCleanMapper extends Mapper<LongWritable, Text, Text, Text> { private Text userId = new Text(); private Text outValue = new Text(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line = value.toString(); // 埋点日志格式:userId,itemId,action,timestamp String[] fields = line.split(","); if (fields.length < 3) { return; // 脏数据直接跳过 } String user = fields[0].trim(); String item = fields[1].trim(); String action = fields[2].trim(); // 行为映射评分:浏览1分,加购3分,下单5分 int score; switch (action) { case "view": score = 1; break; case "cart": score = 3; break; case "buy": score = 5; break; default: return; } userId.set(user); outValue.set(item + ":" + score); context.write(userId, outValue); } }这段Mapper做三件事:按逗号拆字段,长度不足的丢弃,这是最基础的脏数据过滤;把行为类型映射成数字评分,浏览1分、加购3分、下单5分,权重可以按业务调整;以userId为key输出,Reduce阶段按用户聚合出评分向量。评分映射直接影响后续相似度计算的方向,如果只关心“是否购买”,向量值全变成0/1,余弦相似度就退化成Jaccard系数,损失行为强度信息。
可能有人问,数据清洗直接用Hive SQL不就行了?固定结构的日志确实可以,一句INSERT INTO ... SELECT就能替代大半个Mapper。但埋点日志常出现嵌套JSON、URL特殊字符、时间格式不统一,Hive的正则调起来很痛苦,Java代码里用Gson或手写解析更容易控制逻辑。这个项目选Java做清洗是合理的,可读性也更好。
Reducer侧把同一个用户的评分聚合到一起,输出“userId \t item1:score,item2:score”格式,为协同过滤准备输入:
public class LogCleanReducer extends Reducer<Text, Text, Text, Text> { private Text outValue = new Text(); @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 用LinkedHashMap去重并保留插入顺序 Map<String, Integer> itemScoreMap = new LinkedHashMap<>(); for (Text val : values) { String[] parts = val.toString().split(":"); String item = parts[0]; int score = Integer.parseInt(parts[1]); // 同一商品保留最高分,避免重复行为造成歧义 itemScoreMap.merge(item, score, Math::max); } StringBuilder sb = new StringBuilder(); for (Map.Entry<String, Integer> entry : itemScoreMap.entrySet()) { if (sb.length() > 0) sb.append(","); sb.append(entry.getKey()).append(":").append(entry.getValue()); } outValue.set(sb.toString()); context.write(key, outValue); } }这个Reducer里做的去重很关键。用户可能会在一天内反复浏览同一商品,日志里会出现多条“view”记录,如果不去重,评分向量里同一维度出现多个值,后续相似度计算会被放大。这里的merge策略是保留最高分,比累加更合理——加购3分再浏览1分,最高分3分,符合“用户对该商品兴趣的最高强度”这个语义。输出格式保持简洁,一行一个用户,后面的相似度计算直接按行读取即可。
3. 商品推荐算法在MapReduce上的拆解:协同过滤的两种实现路径
3.1 用户-用户协同过滤:找相似的人再找他们买过的商品
协同过滤的核心假设是:过去行为相似的人,未来偏好也相似。用户-用户协同过滤(User-Based CF)是这句话的直接实现:先找到和目标用户行为最相似的K个用户,把这K个人买过但目标用户没买过的商品,按“相似用户+商品”的加权热度排序,截取前N个作为推荐。
在MapReduce上跑,通常分成两轮作业。第一轮把评分矩阵转成“用户-用户”相似度矩阵,Reducer端对每一对用户计算相似度;第二轮用相似度矩阵和评分矩阵做乘法,生成每个用户的候选商品得分。这轮的关键是:相似度计算不能把所有评分向量都load进内存再算,那样单机内存瓶颈又回来了。正确做法是让MapReduce的shuffle机制把公共评分项相同的用户对聚合到同一个Reducer处理。
实现上要控制输出规模。用户量为U时,两两相似度是U的平方级别,十万用户就是百亿条记录,直接落盘会撑爆HDFS。常见做法是先过滤掉行为数过少的用户——只有一条购买记录的用户,算出来的相似度没有统计意义,直接在Map端丢弃。
3.2 物品-物品协同过滤:从共现频率到关联推荐
物品-物品协同过滤(Item-Based CF)思路反过来,不关心人像不像,关心商品之间有没有共现关系。用户同时买了A和B,A和B之间就产生一条关联边,共现次数越多关联越强。它在电商场景有个天然优势:物品之间的相似关系比用户关系稳定得多,一个商品的上架周期内相似度变化不大,一天一算和三天一算差别很小,全量重算的成本摊销下来是划算的。
实现上,物品-物品比用户-用户更省资源。第一轮Map阶段把每个用户的购买列表输出成“itemA:itemB”共现对,Reduce阶段统计商品对出现次数;第二轮把共现矩阵归一化成相似度,再针对每个用户的历史商品,在相似商品集合里做累加排序。
工程上常见的做法是只保留共现次数超过阈值的商品对,比如最少5次共现才进入候选集。否则一个冷门商品因为两三次偶然共现就被强行推起来,推荐列表里会出现大量长尾噪音,用户点进去发现推荐的和自己买的东西毫无关联,体验直接崩。
3.3 余弦相似度计算的Java实现
不管用户-用户还是物品-物品,核心计算单元都是相似度函数。这个项目里相似度模块的代码我单独摘出来讲:
public class SimilarityCalculator { public static double cosineSimilarity(Map<String, Integer> vecA, Map<String, Integer> vecB) { double dotProduct = 0.0; double normA = 0.0; double normB = 0.0; // 合并两个向量的key集合,确保遍历到全特征空间 Set<String> unionKeys = new HashSet<>(vecA.keySet()); unionKeys.addAll(vecB.keySet()); for (String key : unionKeys) { int a = vecA.getOrDefault(key, 0); int b = vecB.getOrDefault(key, 0); dotProduct += a * b; normA += a * a; normB += b * b; } if (normA == 0.0 || normB == 0.0) { return 0.0; // 空向量相似度恒为0 } return dotProduct / (Math.sqrt(normA) * Math.sqrt(normB)); } }这个函数接受两个商品到评分的映射,计算余弦相似度。代码里用unionKeys而不是只遍历交集key,带了一个防御性好处:当某个向量为空时,遍历并集能让你看到key集合的异常,而不是用默认的0掩盖问题。如果你只遍历交集,空向量直接返回0,代码逻辑没毛病,但是问题就被藏住了。
评分映射方式对结果影响很大。直接用0/1表示是否购买,余弦相似度等价于Jaccard系数,只能反映共现关系,反映不了行为强度。项目里把浏览、加购、下单映射成1/3/5,相似度会偏向那些同样是重度购买的用户,效果明显更好。调参时先动这里的映射逻辑,不要一上来就换算法。
3.4 生成推荐列表:从相似度到Top-N输出
相似度算完,最后一步是为每个用户生成Top-N推荐。这一步通常用一个Reducer就能完成:Mapper读入相似度矩阵和用户评分向量,Reducer做加权求和后截取Top-N输出。
public class TopNReducer extends Reducer<Text, Text, Text, Text> { private int topN = 10; @Override protected void setup(Context context) { // Top-N数量从作业配置中读取,默认10 topN = context.getConfiguration().getInt("recommend.topn", 10); } @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 用PriorityQueue维护Top-N,避免全量排序 PriorityQueue<Candidate> queue = new PriorityQueue<>(topN); for (Text val : values) { String[] parts = val.toString().split(":"); String itemId = parts[0]; double score = Double.parseDouble(parts[1]); Candidate candidate = new Candidate(itemId, score); if (queue.size() < topN) { queue.offer(candidate); } else if (candidate.score > queue.peek().score) { queue.poll(); queue.offer(candidate); } } // 按score降序输出 List<Candidate> list = new ArrayList<>(queue); list.sort(Comparator.comparingDouble(c -> -c.score)); StringBuilder sb = new StringBuilder(); for (Candidate c : list) { if (sb.length() > 0) sb.append(","); sb.append(c.itemId).append(":").append(c.score); } context.write(key, new Text(sb.toString())); } static class Candidate { String itemId; double score; Candidate(String itemId, double score) { this.itemId = itemId; this.score = score; } } }这里有个工程习惯值得学习:不要在reduce里用List收集所有候选然后全量排序。一个热门用户可能有几千个候选商品,全排序是O(n log n),PriorityQueue固定容量为N,每次插入只维护前N个最大,复杂度O(log N),量级差一个档次。后续要把结果写进Redis或HBase,也是直接取这个有序队列,不需要再排一次。参数recommend.topn可以从作业提交命令里动态传,好处是调参不用重新打包代码,改个命令行参数就能比较Top-5和Top-20的线上效果差异。
4. 从零开始复现:Hadoop环境搭建与项目运行
4.1 伪分布式还是集群:先看数据规模再定环境
拿到代码后的第一个决策点是环境。数据量只有几十MB甚至几百MB,完全没必要搭三台机器的集群,那是给自己找麻烦。伪分布式模式让所有守护进程跑在同一台机器上,足够验证代码逻辑。只有当数据量到TB级别时,再考虑扩展到3个或更多节点的集群。
伪分布式搭建流程不复杂:装JDK 8、下载Hadoop发行版、配置SSH免密、配环境变量。网上的教程很多,但关键是版本匹配——JDK版本、Hadoop版本、pom.xml里的hadoop-client版本必须落在同一个兼容区间。我见过不少人用JDK 17跑Hadoop 2.7,作业一启动就报反射权限错误,这种问题排查起来比代码逻辑错误难多了。
4.2 核心配置:core-site.xml、hdfs-site.xml、mapred-site.xml
环境变量配好后,决定能不能跑起来的是三个XML文件。完整配置太占用篇幅,我把每个关键属性的作用说清楚。
core-site.xml:
<!-- 配置HDFS的入口地址,伪分布式写localhost --> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> </configuration>hdfs-site.xml:
<configuration> <!-- 关键参数一:副本数量 --> <property> <name>dfs.replication</name> <value>1</value> </property> <!-- 关键参数二:NameNode元数据目录 --> <property> <name>dfs.namenode.name.dir</name> <value>/usr/local/hadoop/data/namenode</value> </property> </configuration>mapred-site.xml:
<configuration> <!-- 指定MapReduce跑在YARN上,而不是默认的本地模式 --> <property> <name>mapreduce.framework.name</name> <value>yarn</value> </property> <!-- Reduce任务数量,伪分布式设置2~4即可 --> <property> <name>mapreduce.job.reduces</name> <value>2</value> </property> </configuration>几个参数的调整心得:副本数在伪分布式必须设为1,三副本会把单块磁盘容量几十秒内打满。如果你在2核4G的机器上跑,YARN容器内存要人为调小,默认的1G Container会直接让作业在ResourceManager分配阶段被拒。这些细节决定了你能否一次跑通,比代码本身更容易让人翻车。
4.3 数据上传、打包与作业提交命令
环境配好后按三个动作走:上传数据、Maven打包、提交作业。
# 1. 在HDFS上创建输入目录并上传本地日志 hdfs dfs -mkdir -p /input/behavior hdfs dfs -put user_behavior.log /input/behavior/ # 2. Maven构建,产出可提交的Jar mvn clean package -DskipTests # 3. 提交MapReduce作业,指定主类和参数路径 hadoop jar target/grms-1.0.jar com.grms.recommend.RecommendJob \ -Dmapreduce.job.name=grms-recommend \ -Dmapreduce.job.reduces=4 \ /input/behavior /output/recommend解释一下提交命令里的三个-D参数。mapreduce.job.name给作业起名,方便在YARN Web UI里定位;mapreduce.job.reduces控制Reduce任务数,伪分布式设成4,真实集群根据数据量来,后面第5.1节会讲数据倾斜时怎么调;最后两个路径参数是输入目录和输出目录。这里有个新手必踩的坑:输出目录不能提前存在,否则Hadoop直接抛FileAlreadyExistsException。
4.4 结果验证:查看HDFS输出和YARN日志
作业跑完后,先确认退出状态码,再用命令行查看输出:
# 查看作业退出状态 echo $? # 0表示运行成功 # 列出输出目录的part-r-xxxxx文件 hdfs dfs -ls /output/recommend/ # 打印前10行推荐结果 hdfs dfs -cat /output/recommend/part-r-00000 | head -10看到输出是“userId \t itemId1:score,itemId2:score”格式,说明链路已经通了。如果作业失败,不要瞎猜,直接去看YARN日志:
# 按作业ID查看日志 yarn logs -applicationId application_1690000000000_0001把堆栈里第一个非Caused by的报错拿去搜,大部分错误都能找到现成答案。真正难解决的是那些日志里不报错但结果明显不对的情况,这种往往要回到数据层面排查。
5. 避坑指南:MapReduce推荐系统的常见问题与排查
5.1 数据倾斜:Reduce端长尾任务的排查路径
现象:作业卡在Reduce阶段很久,99%的Reduce任务早跑完了,只剩一两个任务挂在那儿,作业迟迟不结束。
原因:商品推荐场景里,热门商品的共现频次远远超过普通商品。爆款A和爆款B的共现对,可能是普通商品对的几百倍,导致这一个Key的Reducer要处理的数据量远超其他Reducer,形成长尾。
解决:先看Counter确认哪个Key的记录数异常。方向有两个:如果是热门Key,用拆Key的方式把大Key拆成多个小Key,Reduce阶段先做局部聚合,再做全局聚合;如果只是数量分布不均匀,不要盲目调大Reduce数量,先把超过阈值(比如共现频次1000次)的Key单独抽出来做二次聚合,剩下的走正常Shuffle。跑这个项目时我的做法是:在Reduce入口根据Key额外加一个随机后缀,第一轮聚合完再去掉后缀做第二轮聚合,效果立竿见影。
5.2 小文件过多:HDFS NameNode内存告警
现象:跑完几轮清洗作业后,输出目录里全是一堆几KB甚至几百字节的小文件,NameNode内存涨得很快,集群响应变慢。
原因:默认情况下MapReduce每个Map任务生成一个输出文件。如果喂入的是几百个小日志文件,每份才几百KB,框架会启动几百个Map任务,输出几百个文件碎片。NameNode元数据要记录每个文件的信息,小文件一多内存就撑不住。
解决:清洗阶段先把输入合并成1GB左右的大文件;或者改用CombineFileInputFormat,它能把多个小文件打包成一个split,减少Map任务数。注意这个类不是默认的,需要在Driver里显式调用setInputFormatClass并指定参数。
5.3 本地能跑集群失败:依赖冲突与ClassNotFound
现象:在IDEA里单元测试全过,本地也能出结果,但打包上传到集群后,提交作业报ClassNotFoundException,找不到org.apache.hadoop.thirdparty.protobuf之类的类。
原因:编译期依赖和运行时依赖不一致。两种情况:pom.xml里写了hadoop-client依赖,但打包时没有把依赖打进去,集群上找不到类;另一种是hadoop-client依赖的protobuf版本和集群自带版本冲突,集群优先加载自带版本,抛版本冲突。
解决:不用assembly插件打fat jar,改用maven-shade-plugin,并在打包时对依赖做relocation,把项目自带的Hadoop类重命名到不冲突的包名。更简单粗暴的做法是:pom里把hadoop-common、hdfs、mapreduce-client-core等依赖的scope改成provided,这些运行时由集群提供,打包时自然排除掉。
5.4 中文编码问题:MapReduce读取GBK日志乱码
现象:清洗结果里的中文商品名变成一串问号,推荐结果没法看。
原因:日志文件是GBK编码,代码读取时默认UTF-8,中文被解释成乱码;另一种情况是HDFS上传时的编码转换已经出错。
解决:上传前用file -bi命令确认源文件charset,是GBK就先用iconv转成UTF-8再put。如果文件已经传到HDFS上,在Mapper的setup里显式指定编码:
@Override protected void setup(Context context) { String encoding = context.getConfiguration() .get("mapreduce.map.input.encoding", "UTF-8"); this.charset = "GBK".equalsIgnoreCase(encoding) ? Charset.forName("GBK") : StandardCharsets.UTF_8; }注意mapreduce.map.input.encoding参数在Hadoop 2.6之后才支持,老版本需要用InputStreamReader手动包装。
5.5 Windows下用IDEA调试Hadoop作业的三个坑
最近不少人想在Windows上直接用IDEA跑这个项目,我拆这个项目时也在Windows上折腾过一轮。
第一个坑是winutils缺失。运行时报Failed to locate the winutils binary in the hadoop home directory,解决方法是下载winutils.exe放进Hadoop home的bin目录,并在环境变量里配好HADOOP_HOME。
第二个坑是文件权限。Windows上JVM拿到的用户名和Linux完全不同,HDFS默认ACL会拒绝非superuser的写权限。在代码里设置:
System.setProperty("HADOOP_USER_NAME", "root");这行必须放在任何HDFS客户端API调用之前,比如FileSystem.get()之前。
第三个坑是网络代理。Hadoop 3.x客户端会优先读取环境变量里的HTTP_PROXY,Windows上尤其常见,导致连接NameNode被代理拦截,报Connection refused。排查方法是在IDEA的Run Configuration里把HTTP_PROXY和HTTPS_PROXY环境变量清掉。这个问题Linux上基本碰不到,Windows上几乎人人撞一次,我一开始还以为是防火墙问题,排查了一个下午。
6. 进阶:从离线批处理到近实时推荐,这套项目可以怎么改
6.1 全量重算的代价:用输出目录版本号控制迭代
这套MapReduce链路跑完一轮的时间,取决于数据量。TB级别下通常一到几个小时。显然不能每次用户刷新推荐列表都跑全量作业。工程上轻量的做法是每天凌晨用crontab触发全局重算,输出目录带时间戳:
hadoop jar target/grms-1.0.jar com.grms.recommend.RecommendJob \ /input/behavior/$(date +%Y%m%d) \ /output/recommend/$(date +%Y%m%d)这样每天产出独立目录,互不污染。应用层通过配置切换当天的目录,不需要重启服务。我个人的习惯是保留最近7天的输出,更早的用hdfs dfs -rm -r清掉,避免NameNode被历史版本拖垮。
6.2 增量更新的替代方案:Spark Streaming加状态合并
如果你确实要把延迟降到分钟级,MapReduce就兜不住了。常见做法是把这套逻辑用Spark Streaming重写,消费Kafka里的用户行为流,按窗口累积成微批次,每个微批次跑一次小规模的协同过滤增量更新,然后和前一天的全量结果做合并。复杂度主要在于状态保存:相似度矩阵要存在Redis或HBase里,窗口结束后只对新增行为计算增量相似度,再更新矩阵。
增量计算的时间复杂度远小于全量,但要注意尽量保持一致性。如果一致性要求低,直接把当天增量结果覆盖到全量结果上;要求严格,则保留历史版本做双写。这份代码结构很规整,Mapper和Reducer抽得很干净,改写Spark的transform算子花不了太多时间,我大概用一个周末完成了从MapReduce到Spark的迁移。
6.3 推荐质量验证:一份离线效果评估清单
改完之后,不要只盯着作业跑通,还要确认推荐效果没变差。我通常从三个维度验证:排序能力用离线AUC、召回率用Top-N命中率、多样性看候选列表里不同类型商品数量。具体做法是取最近30天行为数据,前29天训练,最后1天测试,跑完推荐后统计测试集商品在推荐列表里的比例。如果Top-10召回率低于5%,先回第2.3节检查评分映射,而不是急着换算法。
后来我把这套验证逻辑固化成三个自动检查:HDFS输出目录的part文件数量、推荐列表里是否有重复itemId、随机抽3个用户核对推荐结果。有一次就是part文件数量异常多,提前发现了一个数据倾斜隐患。从那以后我每次改完代码都强制走一遍这三个检查,再决定要不要上线。少踩了很多重复的坑,希望你也能用这套方法省下排查时间,希望帮到你。
本文还有配套的精品资源,点击获取