简介:在数据量暴涨的电商场景中,单机协同过滤算法往往难以应对千万级用户与百万级商品的计算压力。分布式计算框架将大规模数据分而治之,其中Hadoop天生适合此类离线批量任务。通过MapReduce多阶段作业,可以完成用户行为日志的存储、商品共现矩阵的构建以及相似度归一化计算,最终为每个用户生成TopN推荐列表。本文将从伪分布式环境搭建开始,逐步拆解协同过滤在HDFS与MapReduce中的代码实现,覆盖作业串联、DistributedCache应用和常见坑点,帮助读者掌握一个可运行的分布式推荐系统最小闭环。
1. 基于Hadoop的商品推荐系统:这是个课程设计,但别小看它
如果你在校园招聘或面试中聊过推荐系统,一定会发现一个尴尬的事实:大部分教学项目都跑在单机内存里,几十万条评分数据往DataFrame里一塞,用Python的字典算相似度,五分钟出结果。但等到面试官问一句“如果用户量涨到千万,商品涨到百万,你原来那套还能用吗”,基本就卡住了。基于Hadoop的商品推荐系统,就是用来回答这个问题的。它能解决的是:在分布式文件系统HDFS上存储用户行为日志,用MapReduce把“用户-商品”共现矩阵和物品相似度算出来,再通过推荐生成器给每个用户产出TopN商品列表。整套流程不依赖内存计算框架,不装Spark也能跑,非常适合拿来理解分布式推荐的最小闭环。
这个项目最适合两类人:一类是做Hadoop课程设计、毕业设计的学生,需要把“协同过滤”落到可运行的代码上;另一类是准备Hadoop相关岗位面试的从业者,想搞清楚MapReduce在真实业务问题里到底怎么拆分任务。我见过太多人把Hadoop学成了“会敲启动命令”,真正面对一个业务问题时不知道从哪下手。这个项目恰好把存储、计算、排序、多阶段作业串了起来,做完一遍,你对MapReduce的理解会超出背面试题的人一大截。下面我按自己复现这个项目时走过的完整路径来拆解,从环境搭建讲到代码实现,再讲到那些配置文件和工作日志里的玄学坑。
2. 伪分布式环境搭建:Hadoop安装与配置的四步走
2.1 环境选型:为什么课程设计首选伪分布式而非集群
很多人在搭建阶段就劝退了,因为一上来就想搭三台甚至五台机器的集群,结果光同步配置就折腾了两周。这个项目的数据量和计算量根本不需要真集群。Hadoop生态里伪分布式模式(Pseudo-Distributed)是官方支持的一种部署方式,所有守护进程都跑在同一台机器上,但每个进程是独立的Java进程,数据也真实地写入HDFS,完全能模拟分布式文件系统的行为。这是从零开始最平滑的路径。
另一种是单机模式(Local Mode),它不启动HDFS,直接在本地文件系统上跑MapReduce,跑起来很快,但体验不到HDFS的读写、数据块复制、SecondaryNameNode这些机制。面试官要是问“你HDFS的副本策略怎么验证的”,单机模式没法回答。所以我一般建议课程设计和面试准备都走伪分布式,这也是绝大多数“基于Hadoop的xx系统”毕业设计采用的模式。等你在伪分布式上把代码磨通了,再考虑集群化,这时你会发现只要把core-site.xml和hdfs-site.xml里的地址改成主机名,就能平滑迁移。
2.2 从零开始安装:JDK版本与SSH免密的血泪经验
先说JDK。Hadoop 3.x要求Java 8或Java 11,我踩过最大的坑是装了个Java 17,结果Hadoop 3.2.2的启动脚本在JVM参数解析上直接翻车,报错信息还特别隐晦,说什么Unrecognized option。所以先把JDK版本卡死:Hadoop 3.2.2配JDK 8最稳,Hadoop 3.3.x配JDK 8或11都行。安装步骤不多,但每一步都有讲究,下面是从零开始的标准流程:
# 1. 创建Hadoop专用用户,避免用root直接跑(root跑会有权限和脚本判断问题) sudo useradd -m hadoop sudo passwd hadoop # 2. 下载并解压Hadoop到指定目录(以3.2.2为例,解压后目录名带版本号) wget https://archive.apache.org/dist/hadoop/common/hadoop-3.2.2/hadoop-3.2.2.tar.gz sudo tar -zxvf hadoop-3.2.2.tar.gz -C /usr/local/ sudo mv /usr/local/hadoop-3.2.2 /usr/local/hadoop sudo chown -R hadoop:hadoop /usr/local/hadoop # 3. 配置环境变量,写到hadoop用户的~/.bashrc里 export HADOOP_HOME=/usr/local/hadoop export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64这里面的逻辑是:Hadoop的stop-all.sh和start-all.sh脚本会用ssh去连接localhost来管理守护进程,如果不配免密,每次启动都要输密码,伪分布式模式下体验极差。生成密钥的命令如下:
ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys如果你是在Windows下用IDEA搭Hadoop开发环境,本机没有原生ssh,一般是在Linux虚拟机或Docker容器里做环境,Windows这边只写代码和打包Jar。这个组合我后面专门讲。当你在终端执行ssh localhost不用输密码了,这一步就算过了。
提示:很多教程让你在/etc/profile里配环境变量,对单用户开发机来说没问题,但如果后续要跑Hadoop生态的其他组件(如Spark、Hive),不同组件之间可能有版本冲突,把Hadoop相关变量放在hadoop用户自己的~/.bashrc里,隔离性更好,排查问题也更快。
2.3 四个核心配置文件的含义与参数选择
伪分布式模式只需要改四个文件,都在$HADOOP_HOME/etc/hadoop目录下。很多新手习惯直接复制网上的配置,但从头理解参数含义,后面调优和排错都靠这个基础。
core-site.xml里指定NameNode的地址和临时目录:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/usr/local/hadoop/tmp</value> </property> </configuration>fs.defaultFS决定了你的HDFS访问入口,所有客户端的读写请求都走这个地址,端口9000是Hadoop RPC的默认端口。hadoop.tmp.dir如果不配置,HDFS的元数据会写到系统的/tmp目录下,系统一重启数据可能被清理,这是最常见的“格式化后重启丢数据”的根源之一。
hdfs-site.xml里指定副本数和NameNode的元数据目录(只改这个文件,伪分布式和集群的差别就在这里):
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>file:///usr/local/hadoop/tmp/dfs/name</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>file:///usr/local/hadoop/tmp/dfs/data</value> </property> </configuration>伪分布式只有一台机器,副本数必须设为1,否则DataNode启动后因为找不到足够的副本位置会一直报错。如果你把这个配置带到集群上,副本数就要改成3,所以这个参数是单机转集群时最容易忽略的一个坑。
mapred-site.xml这个文件在Hadoop 3.x里默认不存在,需要从模板复制,指定MapReduce的资源调度框架:
cp $HADOOP_HOME/etc/hadoop/mapred-site.xml.template $HADOOP_HOME/etc/hadoop/mapred-site.xmlyarn-site.xml指定YARN的ResourceManager和NodeManager地址。偏好设置里可以用<value>yarn</value>,强制走YARN调度。到这里配置阶段就结束了。用jps命令能看到NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNode这五个进程全部启动,环境就绪。我第一次做的时候漏了mapred-site.xml的复制,结果作业提交后一直卡在WAITING状态,一度以为是系统玄学问题,后来才发现是根本没切换到YARN模式,任务压根没被调度起来。
注意:修改任何配置文件后,如果改动涉及NameNode的元数据目录或core-site.xml,需要重新执行hdfs namenode -format。不需要格式化的场景是只改mapred-site.xml和yarn-site.xml。很多教程只说“改完要格式化”,实际上格式化会清空HDFS上的所有数据,如果你已经上传了数据,格式化等于全部报销。
2.4 Windows下用IDEA搭建开发环境:代码在本地,运行在虚拟机
这是热词里出现频率很高的问题,也是学生群体最常卡住的地方。常见的做法是:Windows本机用IDEA写Java代码和做单元测试,通过Maven的hadoop-client依赖在本地模式下测试逻辑;真正的Hadoop环境跑在Linux虚拟机里,代码编译打包成Jar后拷贝到虚拟机执行。
这背后的原因很直接:Hadoop的官方发行版没有原生Windows支持,跑在Windows上要额外装Winutils和Hadoop.dll,而且很容易遇到权限和路径兼容问题。与其在Windows上死磕,不如用这种“两地分离”的方案。
具体路径是:在IDEA的pom.xml里引入hadoop-client依赖,scope设为provided,这样本地编译和测试都能用,打包时不打入Jar;本地测试时设置Hadoop的fs.defaultFS为file:///,跑LocalJobRunner,这能验证你的Mapper和Reducer逻辑有没有低级错误;然后执行mvn clean package打出Jar包,通过scp推到虚拟机的/home/hadoop目录;最后在虚拟机上用hadoop jar命令提交到YARN执行。
这套流程我强烈推荐先做,因为它能极大缩小排查范围。你写了个推荐算法的Reducer,不确定是代码问题还是Hadoop环境问题,在本地跑一遍就知道。等本地输出正确了,再上虚拟机提交作业,这时候出错基本就是环境或路径问题。
3. 推荐链路的数据准备:从原始日志到“用户-商品-评分”三张表
3.1 推荐系统要什么数据:显式反馈与隐式反馈
商品推荐系统分为基于内容的推荐和基于协同过滤的推荐两大类。基于Hadoop实现的绝大多数是协同过滤。协同过滤的基础是用户行为数据,这个项目里用的通常是评分数据或购买记录。
数据分为两类:显式反馈(用户主动打分,比如1到5分)和隐式反馈(浏览、点击、加入购物车)。两者的处理逻辑差别很大。显式反馈直接用评分作为权重;隐式反馈没有负样本,通常把“有行为”记为1,“无行为”记为0或不做记录。
这个课题里最常见的数据集是MovieLens的评分数据,或者自建的电商订单数据。数据格式一般是三列:用户ID、商品ID、评分,用逗号或制表符分隔。我在复现时会把它转换成下面这种结构:
user_id,item_id,score 1001,3001,4 1001,3002,3 1001,3005,5 1002,3001,5 1002,3002,2 1003,3008,4你的HDFS上需要建立三个目录:/input存放原始评分数据,/output存放中间结果,/final存放最终推荐结果。上传命令是hdfs dfs -mkdir和hdfs dfs -put。
3.2 为什么要用两步MapReduce:相似度计算必须依赖共现矩阵
这是整个项目理解上的核心门槛:协同过滤不是用一个MapReduce就能算完的。它需要两个阶段,每个阶段各一个或多个MapReduce作业。
第一阶段计算“商品-商品”共现矩阵。这一步要统计的是:同时出现在同一个用户评分列表中的商品对有多少次。这个统计必须按用户分组,把每个用户评过分的商品两两配对,然后累加。
第二阶段基于共现矩阵计算相似度,再根据相似度和用户的历史评分生成推荐结果。为什么不能合并为一个作业?因为共现矩阵的计算结果是第二阶段的输入。MapReduce作业之间通过HDFS传递数据,前一个作业的输出目录是后一个作业的输入目录。如果你强行写到一个作业里,Mapper阶段的数据是用户-商品对,还没形成商品-商品的共现关系,Reducer拿不到完整的相似度信息,逻辑没法闭环。
所以项目的主流程是:作业1把评分数据转换成商品共现矩阵,写入HDFS;作业2读取共现矩阵和原始评分数据,计算相似度,然后对每个用户生成推荐列表。
3.3 代码实现:用户分组与商品共现的Mapper和Reducer
下面的代码是作业1的核心,完整可运行,我把关键逻辑都注释了:
public class CoOccurrenceMapper extends Mapper<LongWritable, Text, Text, Text> { private Text outKey = new Text(); private Text outValue = new Text(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 每行格式: user_id,item_id,score String[] fields = value.toString().split(","); if (fields.length < 3) { return; // 脏数据直接跳过,不打断作业 } String userId = fields[0]; String itemId = fields[1]; // 输出: key为用户ID,value为商品ID,让同一个用户的数据进入同一个Reducer outKey.set(userId); outValue.set(itemId); context.write(outKey, outValue); } }这里的关键设计是让“用户ID”作为Map输出的Key,这样同一个用户评过的所有商品就会被分到同一个Reducer的输入里,Reducer里拿着这个用户评过分的完整商品列表,两两组合输出共现对。
public class CoOccurrenceReducer extends Reducer<Text, Text, Text, IntWritable> { private IntWritable outValue = new IntWritable(); @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // 收集同一个用户的全部商品ID List<String> items = new ArrayList<>(); for (Text val : values) { items.add(val.toString()); } // 去重:同一个用户可能对同一商品有多条记录,只保留一条 Set<String> itemSet = new HashSet<>(items); List<String> itemList = new ArrayList<>(itemSet); // 两两组合,生成商品对,value固定为1,表示这一对共现一次 for (int i = 0; i < itemList.size(); i++) { for (int j = i + 1; j < itemList.size(); j++) { // 商品对之间要有一个分隔符,便于第二阶段的Mapper区分 outKey.set(itemList.get(i) + ":" + itemList.get(j)); outValue.set(1); context.write(outKey, outValue); } } } }这段代码有个容易被忽略的点:Reducer里我没有直接遍历values去两两组合,而是先放入Set去重。原因是同一个用户对同一商品可能因为数据问题有多条记录。如果不去重,商品对会被重复计数,共现矩阵的数值会被放大,最终影响相似度的准确性。这个坑在数据量小的时候不容易发现,一旦数据量上来,结果偏差非常明显。
接着还需要一个计算共现次数的Reducer,因为MapReducer输出后相同的商品对会有多个值,需要累加:
public class CountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private IntWritable outValue = new IntWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } outValue.set(sum); context.write(key, outValue); } }这三个类组合起来是标准的两阶段MapReduce:先按用户分组输出商品对,再累加计数。跑完后的输出格式是itemA:itemB \t count,比如3001:3002 \t 2。这个文件就是商品共现矩阵,是第二步相似度计算的唯一输入。
注意:这个实现里有个计算倾向,只统计了共现次数,没有除以商品的总出现次数。严格来说这不是标准余弦相似度,而是“共现次数”本身。但如果你的数据集是评分数据,评分的高低没有参与计算,后续还有优化的空间。课程设计阶段这样做已经能出合理结果,但面试被问“你的相似度公式是什么”时会露怯。所以我在第二阶段的Mapper里补充了归一化逻辑,见下一章。
4. 相似度计算与TopN推荐生成:MapReduce实现四个关键类
4.1 相似度矩阵的实现:基于共现次数的归一化处理
第二阶段的第一部分是把共现矩阵转换成语义明确的相似度矩阵。共现次数高的商品不一定相似,比如啤酒和尿布共现次数很高,但它们不是同类商品。所以业界通常用归一化方法,常见做法是计算Jaccard相似度或余弦相似度。
我在这个项目里实现的量化方式是:将共现次数除以两个商品各自“被出现的总次数”的乘积平方根,近似于余弦相似度。这就需要在同一个作业里同时读取两个输入:一是共现矩阵,二是每个商品被多少个用户买过的商品流行度统计。
这个场景正好用得上MapReduce的MultipleInputs机制,它允许一个作业读取多个输入目录:
Job job = Job.getInstance(conf, "SimilarityCalculation"); // 从共现矩阵目录读取 MultipleInputs.addInputPath(job, new Path("/output/cooccurrence"), TextInputFormat.class, SimilarityMapper.class); // 从商品流行度目录读取 MultipleInputs.addInputPath(job, new Path("/output/item-popularity"), TextInputFormat.class, PopularityMapper.class);为什么需要两个Mapper?因为两个输入的数据格式不同,共现矩阵是“itemA:itemB \t count”,流行度是“item \t count”,如果共用一个Mapper,就要在map方法里靠着判断字符串里有没有冒号来区分,逻辑混在一起还容易出错。分开写更清晰。
SimilarityMapper的做法是把共现矩阵的key拆成两个商品,value保留共现次数,转成自定义的可写对象,输出给Reducer。这个自定义对象编码在项目里通常叫PairWritable或者直接用Text拼接:
public static class SimilarityMapper extends Mapper<LongWritable, Text, Text, Text> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] parts = value.toString().trim().split("\t"); if (parts.length != 2) return; String itemPair = parts[0]; // 格式: itemA:itemB long count = Long.parseLong(parts[1]); String[] items = itemPair.split(":"); if (items.length != 2) return; // 把商品对和共现次数一起发给Reducer,Key用第一个商品 context.write(new Text(items[0]), new Text("PAIR:" + items[1] + ":" + count)); } }PopularityMapper的输出Key也是商品ID,value为“POP:商品总数”。这样Reducer里同一个商品ID会同时收到两类数据:PAIR开头的是它和哪些商品共现过,POP开头的是它自己的流行度。Reducer里把这两类数据组织好,就能算相似度。
Reducer端我用一个Map把所有信息暂存,遍历完成后统一计算。理论上数据规模大时不应该缓存全部数据到内存,但单机伪分布式跑课程设计的数据量,这个方案简单有效:
public static class SimilarityReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { Map<String, Long> coCountMap = new HashMap<>(); long itemPopularity = 0L; for (Text val : values) { String token = val.toString(); if (token.startsWith("PAIR:")) { String[] parts = token.substring(5).split(":"); coCountMap.put(parts[0], Long.parseLong(parts[1])); } else if (token.startsWith("POP:")) { itemPopularity = Long.parseLong(token.substring(4)); } } if (itemPopularity == 0L) return; // 计算商品key与其他商品的相似度 for (Map.Entry<String, Long> entry : coCountMap.entrySet()) { double sqrtProduct = Math.sqrt(itemPopularity); double similarity = entry.getValue() / (sqrtProduct + 1e-10); Text outKey = new Text(key.toString() + ":" + entry.getKey()); context.write(outKey, new Text(String.valueOf(similarity))); } } }注意这里我偷懒了,没有传入共现商品那一侧的流行度,精确的算法需要知道itemB左边商品的流行度,因此需要缓存共现矩阵并对每个商品对做两边匹配。正规做法是用一个带“点击流”的辅助数据结构,或者写两个MapReduce作业:一个算商品流行度,一个算相似度归约。为了篇幅不展开,我建议你在课程设计里按这种思路实现:先用一个简单的Job统计每个商品出现在多少个用户里,得到itemA:itemB:count和item流行度两张表,然后在Reducer里把相似度正式算出来。能用Spark或MapReduce连接的方式实现,保证结果可复现。
4.2 评分矩阵读取:为每个用户的推荐候选生成做准备
相似度算完后,还差一步:怎么把它变成“每个用户的TopN商品”。这一步需要的输入是“相似度矩阵”和“用户的历史评分”。MapReduce核心逻辑如下:
用户过去对商品A评过5分,商品A和商品B的相似度是0.8,那么商品B对这位用户的推荐得分就是 5 * 0.8 = 4.0。把用户所有历史商品对应的相似商品得分累加起来,排名靠前的就是候选推荐商品。
这个流程涉及两张表的关联,MapReduce里做关联最常见的方法是“Map侧缓存”。把相似度矩阵放进DistributedCache,让每个Mapper在读用户评分数据之前就把相似度矩阵加载到内存。然后对每个评分记录,从缓存中找出这个商品的所有相似商品,把相似度乘以评分,输出到Reducer。
下面是一个简化但可运行的实现:
public class RecommendMapper extends Mapper<LongWritable, Text, Text, Text> { // 缓存相似度矩阵: key为itemA:itemB,value为相似度 private Map<String, Double> similarityCache = new HashMap<>(); @Override protected void setup(Context context) { // 从缓存文件加载相似度矩阵,每行格式: itemA:itemB \t similarity try { URI[] cacheFiles = context.getCacheFiles(); if (cacheFiles != null && cacheFiles.length > 0) { Path path = new Path(cacheFiles[0]); FileSystem fs = FileSystem.get(context.getConfiguration()); BufferedReader reader = new BufferedReader(new InputStreamReader(fs.open(path))); String line; while ((line = reader.readLine()) != null) { String[] parts = line.split("\t"); if (parts.length == 2) { similarityCache.put(parts[0], Double.parseDouble(parts[1])); } } reader.close(); } } catch (Exception e) { e.printStackTrace(); } } @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 输入: user_id,item_id,score String[] fields = value.toString().split(","); if (fields.length < 3) return; String userId = fields[0]; String itemId = fields[1]; double score = Double.parseDouble(fields[2]); // 遍历相似度缓存,找出与当前商品相似的其他商品 for (String pair : similarityCache.keySet()) { String[] items = pair.split(":"); if (items[0].equals(itemId)) { double sim = similarityCache.get(pair); double recScore = score * sim; String outKey = userId + ":" + items[1]; context.write(new Text(outKey), new Text(String.valueOf(recScore))); } } } }这个Mapper的问题很明显:对每个评分记录都遍历一次相似度缓存,时间复杂度是O(评分条数 * 相似商品数),数据量稍大就慢。但课程设计的数据集通常只有几万条评分,伪分布式单机跑完全没压力。如果你想优化,可以在setup里建一个Map<商品ID, List<相似商品ID>>的倒排索引,这样map里的循环从“全量”降为“该商品只对应的相似商品”,性能提升明显。
Reducer端相对简单,按用户ID分组,累加候选商品的推荐得分,最后取TopN:
public class RecommendReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // key的格式: user_id:item_id String[] keyParts = key.toString().split(":"); String userId = keyParts[0]; String itemId = keyParts[1]; double totalScore = 0.0; for (Text val : values) { totalScore += Double.parseDouble(val.toString()); } // key改为用户,value为 商品ID:推荐得分 context.write(new Text(userId), new Text(itemId + ":" + totalScore)); } }到这里推荐生成的雏形已经有了,它会输入一个用户和这个用户的所有候选商品得分,每个得分一行,还没有做TopN排序。下一步你要么写一个“按得分排序的简单MapReduce作业”,要么把这个结果再加工成HTML片段。常见做法是再加一个排序作业,Map端以得分作为可比较的倒序Key,Reducer输出前N个。这块代码不复杂,但如果你在面试中被问“MapReduce怎么做全局排序”,就能引出这个部分:让Mapper输出Key为自增序号,或者用TotalOrderPartitioner做全局有序分区。
4.3 主类中的Job串联:MapReduce各作业之间的依赖关系
推荐系统的主类看起来像一串流水账,但作业依赖关系是它最值得讲清楚的部分。以下是项目主类中常见的Job串联方式:
public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Path inputPath = new Path(args[0]); Path coocPath = new Path("/output/cooccurrence"); // Job1: 生成共现矩阵 Job job1 = Job.getInstance(conf, "CoOccurrence"); job1.setJarByClass(RecommendDriver.class); job1.setMapperClass(CoOccurrenceMapper.class); job1.setReducerClass(CoOccurrenceReducer.class); job1.setOutputKeyClass(Text.class); job1.setOutputValueClass(Text.class); job1.setOutputFormatClass(SequenceFileOutputFormat.class); FileInputFormat.addInputPath(job1, inputPath); FileOutputFormat.setOutputPath(job1, coocPath); job1.waitForCompletion(true); // Job2: 计算相似度(依赖Job1输出) Job job2 = Job.getInstance(conf, "Similarity"); job2.setJarByClass(RecommendDriver.class); // 将Job1输出读取为共现矩阵,算出相似度,写入/distributedCache job2.waitForCompletion(true); // Job3: 生成推荐结果 Job job3 = Job.getInstance(conf, "Recommendation"); job3.setMapperClass(RecommendMapper.class); job3.setReducerClass(RecommendReducer.class); // 设置DistributedCache,传入相似度矩阵 job3.addCacheFile(new URI(coocPath + "#similarity")); job3.setNumReduceTasks(1); FileInputFormat.addInputPath(job3, inputPath); FileOutputFormat.setOutputPath(job3, new Path("/output/recommendation")); System.exit(job3.waitForCompletion(true) ? 0 : 1); }这段代码有一个很关键的配置:job2在job1完成之后才启动,job3又依赖job2的输出。MapReduce本身不会自动识别依赖链,Job.waitForCompletion返回true代表作业成功,再启动下一个。如果顺序写反,Job3会因为没有相似度文件而报FileNotFoundException。这是很多复现者容易翻车的地方,控制台报错信息不是“无法连接NameNode”,而是“Failed to locate the file”,一眼看上去和HDFS权限或网络有关,实际上就是作业顺序错了。
提示:上面这个主类的“Job2”没有给具体实现,完整的做法是在这个位置嵌入“统计商品流行度”和“计算相似度”两个作业。课程设计答辩时你可以重点讲这个环节,因为它体现了对多作业编排的理解,比单纯贴代码更能说明你掌握了MapReduce的流程控制。
5. 从零到一跑通时常见问题排查:DistributedCache、中文乱码、数据倾斜
5.1 NamedNode与DataNode进程都活着,但就是连不上9000端口
这个现象我见过太多次:jps一下五个进程都在,但报错Connection refused: 9000。原因是启动HDFS时写了Hostname,但core-site.xml里的fs.defaultFS用的是localhost,或者hadoop用户和root用户的配置不一致。Hadoop是Java进程,网络地址绑定和客户端解析的主机名必须完全匹配,否则你这边的客户端解析到127.0.0.1,服务端绑定到实际IP,根本对不上。
解决方法是保证三处一致:一是在core-site.xml里统一写localhost或具体主机名;二是在/etc/hosts里把主机名映射到127.0.0.1或实际IP;三是执行hdfs dfs -ls时指定hdfs://localhost:9000/来测试连通性。我曾经为这个问题折腾了整整一个下午,最后发现是/etc/hosts里有重复的映射记录,把原来正确的映射顶掉了。
注意:如果你改动了core-site.xml或hdfs-site.xml里的NameNode地址,一定要重新执行hdfs namenode -format。这个操作会清空NameNode上的元数据,虽然伪分布式模式数据量不大,但如果你忘了先备份上传到HDFS的测试数据,会损失惨重。格式化前先hdfs dfs -get把数据拉回本地,这是血泪经验。
5.2 中文字段全部变成问号:Hadoop平台默认不支持UTF-8
数据文件里有中文商品名,结果输出到终端全是???,这是字符集问题,不是程序bug。Linux虚拟机的locale环境默认可能是POSIX或en_US.UTF-8,但Hadoop内部如果没显式设置编码,解析时默认按UTF-8读取,可如果你的数据是用Windows记事本另存的,编码可能是GBK甚至带BOM的UTF-8,读进来就乱了。
先确认数据文件本身的编码,在Linux上用file -i命令查看:如果是charset=iso-8859-1或unknown,说明文件不是UTF-8编码。解决方法很简单,把原始数据在Linux下强制转成UTF-8再上传:
iconv -f GBK -t UTF-8 input.csv > input_utf8.csv hdfs dfs -put input_utf8.csv /input/如果数据里有BOM头,还需要去掉,否则第一列数据会带上一个不可见字符,导致Mapper里split(",")后第一行少一个字段。用sed去掉BOM就行:
sed -i '1s/^\xEF\xBB\xBF//' input_utf8.csv另外在代码中显式指定TextInputFormat的编码很容易踩到序列化坑,所以最省心的方式是保证数据文件本身是干净的UTF-8且无BOM。
5.3 Reduce阶段卡在某个进度条不动:常见的数据倾斜
现象是Map阶段跑得飞快,Reduce阶段99%卡了几十分钟。原因是数据分布不均衡:你的Reducer按用户ID分组,如果有一个超级用户买了上千件商品,他的商品对数量是平方级增长的,其他用户每人只有几十个商品对,这个Reducer要做的工作量远大于其他Reducer,整体进度被拖住。
方向有三个:一是改变分组策略,不要用用户ID给共现矩阵做Key,使用复合Key,把每个用户输出的商品对打散到多个Reducer;二是在生成商品对的地方增加采样或裁剪逻辑,对超长商品列表做随机抽取,限制最多配对数量;三是设置Combiner,在Mapper端先合并相同商品对,减少shuffle数据量。这个项目里第三个方向改动最小,效果也明显:
job.setCombinerClass(CountCombiner.class);Combiner的本质是局部Reducer,能在Map输出端先做一次求和,让网络传输和Reduce输入压力小一个量级。这个坑在课程设计答辩时经常被问到,你能说出“Combiner不能乱用,只适用于满足交换律和结合律的聚合函数”这句,面试基本就稳了。
5.4 用IDEA打包后提交作业报主类找不到
这个报错动不动就出现,原因很朴素。Maven里如果没配maven-jar-plugin或没有指定Main-Class,打出来的Jar包就是个“依赖都回挤到一起”的普通Jar,执行hadoop jar xxx.jar时它找不到入口。解决方式是在pom.xml里加一段打包插件配置:
<plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-jar-plugin</artifactId> <configuration> <archive> <manifest> <mainClass>com.example.recommend.RecommendDriver</mainClass> </manifest> </archive> </configuration> </plugin>如果是依赖了第三方库,光配Main-Class还不够,需要maven-shade-plugin把所有依赖打进一个转移包(Fat Jar)。Hadoop自身的客户端依赖应设置为provided,只用本地编译,不进入最终Jar,否则和集群环境里的Hadoop版本不一致会导致各种NoSuchMethodError。
这条坑特别符合“Windows下使用IDEA搭建Hadoop开发环境”的痛点:本机跑得通,打包到集群就跑不通,十有八九是依赖冲突。排查时可以先把shade插件去掉,打一个干净Jar跑试试,如果还有NoClassDefFoundError,再考虑合并依赖。
5.5 输出目录已存在导致的FileAlreadyExistsException
伪分布式模式下,HDFS上的目录删不掉,跑第二次作业时会直接报FileAlreadyExistsException。很多人误以为这是权限问题,反复chmod。实际上Hadoop的OutputFormat约定输出目录必须不存在,这是防止覆盖数据的保护机制。
解决办法很直接,每次跑新的作业前手动清掉旧输出:
hdfs dfs -rm -r /output/cooccurrence也可以在主类的main方法里加一段“存在即删除”的逻辑,这在工作目录里更省事,但要注意不要误删原始输入:
Path outputPath = new Path("/output/recommendation"); FileSystem fs = outputPath.getFileSystem(conf); if (fs.exists(outputPath)) { fs.delete(outputPath, true); }6. 进阶:从课程设计到能提供推荐效果的几个实用技巧
这个项目的天花板不在“能跑”,而在于“怎么让推荐结果看起来合理、能用”。很多人把代码跑通输出了几百行数字,但不知道如何验证推荐质量。我实践下来有两个最有效的手段。
第一个手段是加一个“去历史”步骤。上面实现的推荐有个明显毛病:用户买过的商品,和它相似的商品里,用户可能已经买过了,却还是会被推荐。对课程设计来说这可以接受,但面试官一追问就露怯。常规做法是加一个FilterMapper,在生成候选推荐时,查一下用户历史数据集,把已经交互过的商品剔除掉。实现逻辑不复杂,把用户历史行为也放进DistributedCache,Mapper里碰到候选商品先查一下是否在当前用户的历史集合里,在的话就跳过。这一步能让推荐结果看起来“聪明”很多。你可以在课堂上演示:用户A给《复仇者联盟》打了5分,推荐列表里如果还出现《复仇者联盟2》,显然不合理,加上过滤后结果干净了。
第二个手段是调关键参数。第一个参数是相似度阈值,很多相似的关联商品都是一根手指头的关系,共现次数只有1次,算出来的相似度虚高。建议在生成相似度矩阵时加一个过滤条件:共现次数小于2的商品对直接丢弃。第二个参数是推荐列表长度N,TopN里N设为多少合适,要看你的数据集中用户平均购买的商品数量。一般课程设计数据里N=10或N=20就行,N太大会把低得分的垃圾商品也带出来,N太小体现不出多样性。
再一个值得做的是可视化输出。MapReduce输出的文本文件在命令行里看非常吃力。我一般会写一个很简单的Python脚本,把HDFS上的输出拉回本地转成CSV格式,然后用Pandas筛选每个用户得分前10的商品列表。这个脚本不是项目的一部分,但能让答辩演示效果提升很多。你先用hdfs dfs -get把结果拉下来,再用pandas.read_csv读取(分隔符是\t),按用户分组排序,输出一个表格。这个方法比把输出分段打印到终端好得多。
最后建议你把Hadoop和Zookeeper整合的实操做一遍,这是目前开发岗面试中高频出现的场景。虽然商品推荐本身不依赖Zookeeper,但一旦你把集群从伪分布式扩展到三节点,NameNode的高可用就需要Zookeeper来选主。你可以把推荐系统部署成两节点集群,用Zookeeper做HA,这样简历上写“基于Hadoop的商品推荐系统,支持NameNode高可用”,含金量就不是普通课程设计能比的了。
我自己的习惯是每改一个参数,就把共现矩阵和最终推荐列表完整导出一份,跟上一版做对比。这个习惯帮我在一个数据倾斜问题上找到了真正原因:当时推荐结果里有三四个商品出现在几乎所有用户列表里,单看得分好像正常,但导出共现矩阵一看,这几个商品和所有其他商品的共现次数都异常高,问题出在“某个用户一次性购买了几百件商品”,把全局统计拉偏了。如果不导出中间结果,这类问题根本没法定位。
所以如果你照着这个项目复现,记得给自己留一个“后悔药”:主类里每个Job的输出目录都保留一套,不要每次覆写。排查问题的时候,中间结果就是你最好的调试工具。希望帮到你。
本文还有配套的精品资源,点击获取