
简介基于Hadoop的好友推荐系统设计与实现项目面向大数据、人工智能、物联网等计算机相关专业的在校学生、教师及企业开发者解决推荐系统从数据预处理、相似度计算到结果展示的完整实现问题。压缩包共2000个文件大小79.5MB其中包含1260个PNG图片、403个CSS样式资源、73个Java源文件、55个class文件以及jar依赖库、JSP页面、XML配置等前端资源与后端代码分区明确便于查看和修改。目前已有162人学习下载。项目附带完整部署文档代码均经过运行验证核心模块涵盖数据访问、聚类分析、距离计算、结果绘图等能够帮助读者理清基于MapReduce的推荐算法实现思路。该资源评审分数达95分适合作为毕业设计、课程设计或项目初期演示也可供有一定基础的开发者在此基础上扩展功能深入学习Hadoop生态。1. 基于 Hadoop 的好友推荐系统为什么把协同过滤拆成 4 个 MapReduce 作业在一张千万节点、上亿条关注关系的社交网络图上如果直接在 MySQL 里执行“算两两用户相似度”的 SQL一次全表笛卡尔积就能把数据库连接池打满。这套基于 Hadoop 的好友推荐系统源码没有引入 Spark 或图计算组件而是用四个职责明确的 MapReduce 作业完成“初始化聚类中心—计算用户到中心距离—增量统计簇内向量—簇内生成推荐结果”的完整链路代码里可以看到FindInitDCMapper、CalDistanceMapper、DeltaDistanceMapper、ClusterDataMapper这四个 Mapper 类配合DBService落库和DrawPic出图。对于正在做课程设计、毕业设计的人以及想搞清“协同过滤在离线批处理里怎么工程化”的开发者这套资料里的部署文档比单纯讲算法的博客更适合照着复现。该方案既能在 Hadoop 伪分布式环境里单机跑通也能平滑迁移到三节点集群。2. 推荐模型与距离计算K-Means 与协同过滤的 MapReduce 化2.1 为什么先聚类再算相似度好友推荐的离线阶段通常要做两类事情第一类是基于共同好友统计 Jaccard 相似度第二类是把用户画像、关注标签、点击行为编码成向量然后计算欧氏距离或余弦相似度。单纯用第一类方案在只有几万用户时就需要进行数亿次交集计算而且用户关系表里重度用户和普通用户的行为分布差异很大shuffle 阶段数据量会被头部用户放大。项目把第二类方案和 K-Means 聚类结合起来先用用户行为向量做聚类把全量用户划分到 K 个簇再只在簇内计算用户间相似度最终从簇内挑选“不是当前用户好友但向量距离最近”的用户作为推荐结果。这个思路的关键收益是在需要 Top-N 推荐时不需要为所有用户对计算分数只需要保证每个簇的规模可控。阶段作业说明Mapper 类输出内容1从样本中挑 K 个初始中心点FindInitDCMapper中心点文件2计算每个用户到所有中心的距离CalDistanceMapper用户 ID - 最近簇 ID3增量统计每个簇的向量和与样本数DeltaDistanceMapper簇 ID - 局部统计量4在稳定簇内生成相似用户对ClusterDataMapper候选好友及相似度源码里CloudAction是作业调度入口HUtils负责把-D参数传递给ToolRunnerDBService只做推荐结果写库DrawPic负责从 HDFS 拉结果画图。也就是说这些类并不是装饰而是按作业边界划分的。2.2 CalDistanceMapper 实现在 Mapper 里加载聚类中心这一节给出第二阶段 Mapper 的可复现代码。要注意启动作业时需要用-files把中心点文件分发到每个节点然后在setup()里读取本地缓存而不是在 Mapper 里直接访问 HDFS否则每个 map task 都会误解一个远程路径网络开销会变大。public class CalDistanceMapper extends MapperLongWritable, Text, Text, Text { private final Listdouble[] centers new ArrayList(); Override protected void setup(Context context) { URI[] cacheFiles context.getCacheFiles(); if (cacheFiles ! null) { for (URI uri : cacheFiles) { Path path new Path(uri); try (BufferedReader br new BufferedReader( new InputStreamReader(new FileInputStream(path.toString())))) { String line; while ((line br.readLine()) ! null) { line line.replace(\uFEFF, ); if (line.trim().isEmpty()) continue; String[] parts line.split(,); double[] center new double[parts.length - 1]; for (int i 1; i parts.length; i) { center[i - 1] Double.parseDouble(parts[i]); } centers.add(center); } } catch (Exception e) { throw new RuntimeException(读取聚类中心失败, e); } } } } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(,); if (fields.length 2) { return; } String userId fields[0]; double[] vector new double[fields.length - 1]; for (int i 1; i fields.length; i) { vector[i - 1] Double.parseDouble(fields[i]); } int nearestIndex 0; double minDistance Double.MAX_VALUE; for (int i 0; i centers.size(); i) { double dist euclideanDistance(vector, centers.get(i)); if (dist minDistance) { minDistance dist; nearestIndex i; } } context.write(new Text(cluster_ nearestIndex), new Text(userId \t minDistance)); } }逻辑说明这个 Mapper 的输入是 HDFS 上/user/hadoop/input/behavior.txt每一行是一个用户 ID 和它的行为向量输出 key 是cluster_0、cluster_1这样的簇 IDvalue 是userId|距离值。后续 Reducer 只需要按簇 ID 做 group不需要再解析用户的全部特征。参数方面中心点文件第一列必须是中心点 ID后面的列必须和用户行为向量的维度一致否则ArrayIndexOutOfBoundsException会在某个 map task 里随机出现而且不会直接在 Driver 日志里显示要打开 container 日志才能看到。如果特征向量是 0/1 稀疏向量建议把欧氏距离换成余弦距离但那样簇中心更新就不能简单用均值需要归一化后求平均。常见做法是在第四阶段再把距离换成相似度排序。2.3 中心点迭代与收敛FindInitDCMapper 与 DeltaDistanceMapper 的配合FindInitDCMapper做的事情是从原始数据里抽 K 行作为初始中心而不是随机在内存里生成 K 个点因为 HDFS 上文件的某个 block 可能在集群任意节点直接使用Random选点无法保证每个 Mapper 都看到全量数据。常见做法是用TableSampler或者干脆用head取前 K 行。DeltaDistanceMapper在第三阶段输出的是簇内局部统计量中心点的真正均值计算由配套 Reducer 完成这里可以用一段 driver 片段说明迭代逻辑double previousCost Double.MAX_VALUE; for (int iter 0; iter maxIterations; iter) { Job job Job.getInstance(conf, KMeans- iter); job.setJarByClass(CloudAction.class); job.setMapperClass(CalDistanceMapper.class); job.setReducerClass(DeltaUpdateReducer.class); job.setNumReduceTasks(k); // 输入是上一次的中心点文件输出是本次新中心 FileInputFormat.setInputPaths(job, lastCentroidPath); FileOutputFormat.setOutputPath(job, newCentroidPath); if (!job.waitForCompletion(true)) { System.exit(1); } // 比较前后两次中心点移动距离小于 epsilon 就停止 double cost computeDelta(lastCentroidPath, newCentroidPath); if (Math.abs(cost - previousCost) epsilon) { break; } previousCost cost; }逻辑说明每次迭代输出一个新中心文件下一次的CalDistanceMapper通过-files或addCacheFile读取这个文件。computeDelta是一个辅助方法读取两个路径下中心文件的坐标并计算差的平方和。参数说明maxIterations一般设 5 到 8 就可以因为伪分布式环境下每次迭代都要经历 shuffle、merge、sort迭代次数越多调试成本越高epsilon设为 0.01 时通常第 3 次迭代后连续两次变化的绝对值就已经小于这个值。3. Hadoop 开发环境搭建与项目部署实战3.1 先跑通伪分布式再上集群很多人在拿到源码后第一反应是开三个虚拟机搭集群结果光格式化 NameNode 和配置免密就花了半天。对这个项目来说伪分布式足够验证算法和结果也更容易通过 IDEA 远程调试。部署文档里如果指定了 JDK 和 Hadoop 版本尽量保持一致没有指定时建议使用 JDK 1.8 搭配 Hadoop 2.7.x 或 3.2.x这两个系列的mapred-site.xml模板兼容性最好。sudo tar -zxvf hadoop-${HADOOP_VERSION}.tar.gz -C /opt sudo useradd hadoop sudo chown -R hadoop:hadoop /opt/hadoop sudo su - hadoop echo export JAVA_HOME/usr/lib/jvm/java-8-openjdk-amd64 ~/.bashrc echo export HADOOP_HOME/opt/hadoop ~/.bashrc echo export PATH\$PATH:\$HADOOP_HOME/bin:\$HADOOP_HOME/sbin ~/.bashrc source ~/.bashrc注意变量${HADOOP_VERSION}要替换成实际下载的版本号不能原样拷贝。JAVA_HOME写错时hadoop version会直接报“找不到 Java”而不是给出版本号另一个常见问题是使用 OpenJDK 11 以上版本编译的 jar 包放到 Hadoop 3.2 上会出现ClassNotFound和模块化访问错误所以不要在这个项目里强行升级 JDK。3.2 四个 XML 文件该怎么改伪分布式的核心配置只需要四个文件。core-site.xml指定文件系统入口hdfs-site.xml控制副本数和元数据路径mapred-site.xml决定 MapReduce 跑在 YARN 上yarn-site.xml负责资源调度。先看前两个!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/opt/hadoop/tmp/dfs/name/value /property property namedfs.datanode.data.dir/name value/opt/hadoop/tmp/dfs/data/value /property /configuration参数说明dfs.replication1是伪分布式必须做的一步如果保持默认 3上传小文件后会一直处于Under replicated状态做 MR 虽然能跑但清理输出目录时容易卡在Waiting to replicate。dfs.namenode.name.dir和dfs.datanode.data.dir不要放到/tmp否则重启机器后元数据会被清空。然后配置mapred-site.xml。在 Hadoop 2.x/3.x 中这个文件默认不存在需要从模板复制cp $HADOOP_HOME/etc/hadoop/mapred-site.xml.template $HADOOP_HOME/etc/hadoop/mapred-site.xmlproperty namemapreduce.framework.name/name valueyarn/value /property如果不设置这个属性MapReduce 会默认跑在本地进程里hadoop jar无法看到集群进度。启动服务hdfs namenode -format start-dfs.sh start-yarn.sh jps正常情况下会有 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 五个进程。如果 DataNode 没起来最典型的原因是/opt/hadoop/tmp目录残留了上次的元数据。解决方式是备份数据后清空 tmp 目录重新格式化。这个操作在生产环境是禁忌但纯学习环境想快速恢复是最常用的手段。3.3 上传数据、提交作业与 MySQL 初始化Hadoop 服务起来后把本地行为数据放到 HDFShdfs dfs -mkdir -p /user/hadoop/input hdfs dfs -put behavior.txt /user/hadoop/input/ hdfs dfs -ls /user/hadoop/input/然后初始化数据库。DBService类里通过BaseDAOImpl建立连接落库表结构可以参考CREATE TABLE tb_user_feature ( user_id VARCHAR(32) PRIMARY KEY, feature_vector VARCHAR(2048) ); CREATE TABLE tb_recommend_result ( user_id VARCHAR(32), friend_id VARCHAR(32), score DOUBLE, is_friend TINYINT DEFAULT 0, PRIMARY KEY (user_id, friend_id) );说明feature_vector用来做人工校验线上跑批时不应该把完整向量存 MySQL否则ClusterDataMapper读取用户数据时会和 MySQL 的 JDBC 连接池互相争抢资源。常见做法是HDFS 上行为向量只保留 user_id 和稀疏字段下标落库时只写推荐结果。DBService.batchInsert()用 PreparedStatement 的addBatch()批量写入避免每条结果各打开一次连接。提交作业hadoop jar your-project.jar com.hadoop.recommend.CloudAction \ -D kmeans.k15 \ -D kmeans.max.iterations5 \ -D mapreduce.job.namefriend-recommend \ /user/hadoop/input/behavior.txt \ /user/hadoop/output/recommend命令中的com.hadoop.recommend.CloudAction是主类kmeans.k和kmeans.max.iterations是自定义参数最后一个输出路径必须不存在。如果希望多次运行同一输出路径可以在CloudAction的 main 里显式判断并删除Path outPath new Path(args[args.length - 1]); if (outPath.getFileSystem(conf).exists(outPath)) { outPath.getFileSystem(conf).delete(outPath, true); }逻辑说明ToolRunner.run会先解析-D开头的参数到 Configuration再调用run()方法提前删除输出目录避免频繁手动hdfs dfs -rm -r。3.4 部署资料里最容易被忽略的配置项部署文档里除四个 XML 外还要检查yarn-site.xml的yarn.nodemanager.vmem-check-enabled。伪分布式内存紧张时这个参数会误杀任务。学习环境里建议设置为 falseproperty nameyarn.nodemanager.vmem-check-enabled/name valuefalse/value /property参数说明这个参数控制 NodeManager 是否检查虚拟内存超限。默认开启时即使 map 任务本身没 OOM只要虚拟内存超过yarn.nodemanager.vmem-pmem-ratio的比值Container 也会被强杀。如果日志里有Killing container due to usage of virtual memory第一件要做的事不是加大物理内存而是确认这个参数是否关闭。进入生产环境前要把它改回来避免失控任务拖垮整个 NodeManager。4. MapReduce 调优与集群部署中的边界问题4.1 作业卡在 map 100% reduce 0% 时先看日志再改参数MapReduce 作业在伪分布式里最常见的故障表现是Map 跑到 100%Reduce 一直停在 0%但 job 没有失败。原因是 Reducer 从 map 输出拉取数据时如果输入路径里有空文件或part-r-*输出格式被误当成输入FileInputFormat会产生无效分片导致 Reducer 等不到有效数据。另一个可能的原因是 Reduce 端内存不够频繁 spill 到磁盘。先用以下命令确认yarn application -appStates RUNNING -list yarn logs -applicationId application_1702323456789_0001日志里如果出现Container killed by ApplicationMaster就去调整mapreduce.reduce.memory.mb。好友推荐这种需要排序的作业Reducer 端往往比 Map 端更吃内存因为每个 Reducer 要维护当前簇的局部聚合结果。经验值是 Map 1GB、Reduce 2GBJava opts 分别对应 819m 和 1638mproperty namemapreduce.reduce.memory.mb/name value2048/value /property property namemapreduce.reduce.java.opts/name value-Xmx1638m/value /property注意mapreduce.reduce.java.opts的堆上限一般要比 container 内存少 20%否则 JVM 本身的内存开销会触发vmem-check。4.2 数据倾斜大V用户把单个 Reducer 打满计算好友相似度时如果只按friend_id做 key几百万粉丝的大 V 会把某个 Reducer 的任务量拉高到其他 Reducer 的几十倍。项目里的ClusterDataMapper在生成候选好友时也会面临同样的倾斜问题。常见缓解做法是在聚合前对 key 加盐例如将(friend_id, salt)作为中间 keyReduce 后再对同一个 friend_id 的结果二次聚合。代价是增加一轮排序或第二步作业但在数据量可控时能明显拉平负载。另一个做法是在生成候选集时过滤掉明显膨胀的数据一个用户如果与超过 1000 人关联则在候选集里只保留相似度最高的 30 个。可以在 Mapper 里直接判断if (userFriendCount MAX_CANDIDATE_COUNT) { context.write(new Text(userId), new Text(join(similarityTopN))); }这里的MAX_CANDIDATE_COUNT可以从-D recommend.max.candidate30传入答辩时能说明你考虑了推荐结果的“头部长尾”问题。参数说明recommend.max.candidate最好设为最终 Top-N 的 1.5 到 3 倍比如最终展示 20 个好友候选人可以保留 30 到 60 个后续再做过滤和排序。4.3 伪分布式和集群环境的资源参数差异如果要把项目迁移到三节点集群YARN 的资源设置要按节点内存重新计算。假设每个节点 8GB给操作系统和 HDFS 保留 2GBNodeManager 可用内存设为 6GBproperty nameyarn.nodemanager.resource.memory-mb/name value6144/value /property property nameyarn.scheduler.maximum-allocation-mb/name value2048/value /property property nameyarn.scheduler.minimum-allocation-mb/name value512/value /property参数说明yarn.nodemanager.resource.memory-mb是单节点可分配给所有 Container 的物理总内存yarn.scheduler.maximum-allocation-mb决定单个 Map/Reduce 任务最多申请多少。在 6GB 节点上设置最大 2GB表示最多同时跑 3 个任务容器。如果把最大值设成 6GB单个任务很爽但两个任务同时到达时会互相等资源整个作业反而变慢。跨节点部署时还要注意mapreduce.map.memory.mb不能超过yarn.scheduler.maximum-allocation-mb否则 ResourceManager 一直无法满足资源请求作业卡在 ACCEPTED 状态。遇到这种“任务不跑”的情况用yarn application -status看ResourceRequest里的Memory再对比yarn-site.xml里的上限。4.4 当部署文档说要整合 Zookeeper 时到底改哪里所谓“hadoop 和 zookeeper 整合实战”通常指启用了 NameNode HA 或 ResourceManager HA。伪分布式用不到 ZK但在三个节点上做高可用时core-site.xml需要加ha.zookeeper.quorum配置hdfs-site.xml需要声明dfs.nameservices、dfs.ha.namenodes.ns以及两个 NameNode 的地址同时让 JournalNode 进程负责同步 edit log。这里的实操关键不是背配置项而是记住顺序先启动三台机器的 JournalNode再格式化 NameNode最后初始化 Standby NameNode。如果顺序反了HA 状态会一直停在initializinghdfs haadmin -getServiceState nn1看不到 active。可以通过脚本验证hdfs zkfc -formatZK hdfs haadmin -getAllServiceState注意zkfc命令只负责把故障转移状态写入 Zookeeper不能和 NameNode 格式化混在一起。在高可用环境里集群高可用和推荐业务没有直接关系只是保证离线任务不会因为 NameNode 单点挂掉而中断。5. 用 DrawPic 快速验证推荐结果从 HDFS 拉数据到本地画图在伪分布式环境里调试推荐结果时每次用hdfs dfs -cat查看一长串记录很不直观。项目里的DrawPic类就是干这个的把 HDFS 上的推荐结果文件拉到本地解析userId \t candidateId \t score行然后按 score 排序并输出 Graphviz 点文件。这样可以用一张图快速判断结果是稀疏还是稠密是不是存在某个用户被大量推荐的情况。public class DrawPic { public static void drawTopN(String hdfsPath, String localPath) throws IOException { Configuration conf HUtils.getConf(); FileSystem fs FileSystem.get(conf); fs.copyToLocalFile(new Path(hdfsPath), new Path(localPath)); ListString[] rows new ArrayList(); try (BufferedReader br new BufferedReader( new InputStreamReader(new FileInputStream(localPath), StandardCharsets.UTF_8))) { String line; while ((line br.readLine()) ! null) { String[] parts line.split(\t); if (parts.length 3 !parts[2].equals(NaN)) { rows.add(parts); } } } rows.sort((a, b) - Double.compare(Double.parseDouble(b[2]), Double.parseDouble(a[2]))); StringBuilder dot new StringBuilder(digraph recom {\n); for (int i 0; i Math.min(20, rows.size()); i) { String[] r rows.get(i); dot.append( \).append(r[0]).append(\ - \) .append(r[1]).append(\ [label\).append(r[2]).append(\];\n); } dot.append(}\n); Files.write(Paths.get(recommend.dot), dot.toString().getBytes(StandardCharsets.UTF_8)); } }逻辑说明copyToLocalFile会把 HDFS 上整个输出文件拉到本地如果输出是多个part-r-00000就不要指定文件路径而是指定输出目录再在本地循环处理part-*文件。排序只保留 Top 20是为了避免 1000 个候选好友画出一团乱线。.dot文件可以直接用dot -Tpng recommend.dot -o recommend.png生成图片放进课程设计报告里。要注意的是HUtils.getConf()里不要硬编码hdfs://localhost:9000而应该从 XML 中读取fs.defaultFS。因为在虚拟机上调试时每台机器的 hostname 可能被改来改去硬编码 IP 会导致DrawPic在集群模式下连不上 NameNode。一般 HUtils 的做法是public static Configuration getConf() { Configuration conf new Configuration(); conf.addResource(core-site.xml); conf.addResource(hdfs-site.xml); return conf; }提示画图前最好先检查 HDFS 结果文件的行数和字段完整性。如果某条记录的 score 是NaN通常是该用户的特征向量全为 0余弦相似度分母为零。可以在ClusterDataMapper里过滤vectorNorm 0的用户避免脏数据进入最终结果。这个技巧在答辩演示时很实用把DrawPic跑出来的图放在主流程展示页旁边比贴一段日志更能说明系统已经跑通。如果换一批输入数据只需要改DrawPic的入参和过滤条件不需要重新提交整个 Hadoop 作业。本文还有配套的精品资源点击获取