
简介在机器学习与大数据处理深度融合的今天K近邻算法KNN作为经典的分类方法以其无需显式训练、逻辑直观的特点广泛应用于用户画像、推荐系统等场景。然而当数据规模膨胀海量样本间的距离计算成为性能瓶颈这正是分布式计算框架MapReduce的用武之地。MapReduce通过分而治之的思想将独立的距离计算任务并行化使KNN能够高效扩展至百万级用户。本文以电影网站用户性别预测为实例基于Hadoop平台从特征工程、距离度量、K近邻投票到分布式作业设计完整讲解如何用Java构建一个端到端的机器学习项目。项目采用MovieLens公开数据集涵盖数据预处理、向量构建、训练测试集划分及准确率调优等工程实践并展示特征对齐、参数选择与数据倾斜等关键问题的解决方案。该案例既适合理解算法原理也为大数据与机器学习结合的实战项目提供了可复用的参考模板。 最近好几个读者问我能不能把一个经典的机器学习算法和分布式计算框架结合起来做成一个完整的项目正好我之前做过一个“基于KNN算法和MapReduce实现电影网站用户性别预测”的项目今天把这个项目的完整思路、源码细节和踩坑记录都整理出来。这个项目的核心就是使用Java语言在Hadoop的MapReduce框架上实现KNN分类算法利用用户对电影的评分与观看偏好来预测用户的性别整个过程包含数据处理、特征构建、距离计算、K近邻投票等多个环节。对于正在学习大数据、准备面试、或者想做一个能写进简历的实训项目的朋友来说这个项目具备完整的业务闭环和清晰的技术栈非常适合作为实战参考。先说说这个项目能解决什么问题。电影网站想要做个性化推荐但新用户没有任何行为记录也就是冷启动问题。这时候如果能预测出用户性别就可以先按性别对应的偏好来做初级推荐。而预测性别的数据基础恰恰是很多网站都有的用户观影记录不需要额外采集成本极低。再加上KNN算法本身逻辑简单、无需显式训练配合MapReduce天然适合并行处理大规模距离计算所以这个项目在工程上非常落地。1. 项目整体设计与思路拆解1.1 业务场景为什么看过的电影能暴露性别先说结论不同性别的用户群体在电影类型偏好上存在明显的统计学差异。以经典的MovieLens数据集为例男性用户在动作、科幻、冒险类电影上的平均评分和观看占比通常高于女性用户而女性用户在爱情、剧情、动画类电影上的偏好则更明显。这不是说每个人都如此但从群体分布上看这种偏好差异是真实存在且可以被算法捕捉的。性别预测本质上是一个二分类问题输入是用户的行为特征向量输出是“男性”或“女性”。这个问题的训练数据非常容易获得注册时填写了性别的老用户就是天然的训练样本这批人的观影记录就是特征性别就是标签。而目标用户就是那些没有填写性别或需要交叉验证性别信息真实性的用户。这个项目的另一个价值在于它演示了如何把一个数学上很优雅的算法放到一个真实的分布式计算环境中去执行。KNN本身并不复杂但当用户量从几千扩展到几百万的时候单机计算距离矩阵就会变得不可行所以引入MapReduce是合理的工程决策。1.2 技术选型KNN和MapReduce为什么是绝配先解释KNN为什么适合这个场景。KNN全称K-Nearest Neighbors是一种基于实例的惰性学习算法。所谓惰性学习就是它没有显式的训练阶段所有计算都发生在预测阶段来一个新样本计算它与所有已知样本的距离找到距离最近的K个样本让这K个近邻投票决定新样本的类别。对比逻辑回归或决策树KNN的优点是不需要对数据分布做假设简单直观且天然支持多分类。在性别预测这个任务上特征维度不太高通常几十维数据量适中KNN完全够用。再解释为什么引入MapReduce。KNN有一个致命弱点预测一个样本需要遍历全部训练样本计算距离时间复杂度是O(N)N为训练集大小预测M个样本就是O(M×N)。当用户量到达百万级别这个计算量是恐怖的。而KNN的距离计算有一个非常好的特性每个样本与其他样本之间的距离是完全独立的这正好落在MapReduce擅长的数据并行范式里。Map阶段把待预测样本分发给多个计算节点每个节点并行计算它和部分训练样本的距离Reduce阶段汇总排序取前K个整个过程完美契合“分而治之”的思想。至于为什么用Java答案很简单Hadoop本身是Java写的MapReduce的原生编程接口就是Java。用Java实现不需要额外的中间件和进程通信开销调试也最方便。1.3 项目整体架构与数据流整个项目的实现分成两个依次依赖的MapReduce作业。第一个作业负责数据预处理和特征向量构建输入原始评分数据和电影元数据输出每个用户的特征向量。这个特征向量以电影类型为维度统计用户在每种类型下看过的电影数量也可以叠加评分信息最终形成一条“用户ID 特征向量”的记录。第二个作业是核心的KNN计算与预测把所有带性别标签的训练用户特征向量加载到DistributedCache中Map阶段读取待预测用户特征逐一计算与训练用户的距离输出待预测用户ID距离性别Reduce阶段对同一用户的所有距离排序取前K个按性别投票输出最终的预测结果。从数据流上看整个流程是原始数据 → 特征向量 → 距离矩阵分布式计算 → K近邻 → 投票结果。下面这张表可以清晰展示两个作业的输入输出作业Mapper输入Mapper输出Reducer输出Job1 特征构建ratings.dat、movies.dat(用户ID, 电影类型:评分)(用户ID, 特征向量)Job2 KNN预测待预测用户特征(用户ID, 训练用户ID:距离:性别)(用户ID, 预测性别)2. 核心数据与特征工程2.1 数据集准备MovieLens经典数据项目采用MovieLens 100K数据集这是推荐系统领域最经典的公开数据集之一来自明尼苏达大学的GroupLens研究组。数据包含三个核心文件结构如下users.dat的字段依次是用户ID、性别、年龄、职业编码、邮编。例如1::F::1::10::48067 2::M::56::16::70072这里性别字段就是我们要预测的目标标签。movies.dat的字段依次是电影ID、标题、类型列表。类型用竖线分隔例如1::Toy Story (1995)::Animation|Childrens|Comedy 2::Jumanji (1995)::Adventure|Childrens|Fantasy这个文件用来建立电影ID到类型的映射关系。ratings.dat的字段依次是用户ID、电影ID、评分、时间戳。例如1::1193::5::978300760 1::661::3::978302109评分范围是1到5的整数。数据预处理主要做三件事。第一清洗无效记录比如评分值不在1到5范围内、用户ID或电影ID为空、电影类型为空的记录直接过滤掉。第二处理用户维度同一用户的所有评分记录要归并到一条特征向量中。第三数据集划分把有性别标签的用户按一定比例划分为训练集和测试集训练集用于KNN的参考样本测试集用于评估预测准确率。这里有一个容易踩的坑MovieLens数据集的编码是ISO-8859-1不是UTF-8。直接用Java默认字符集读取中文或特殊字符时会乱码建议在解析文件时显式指定字符集。2.2 特征向量构建把观影行为变成数学向量特征工程是整个项目中影响准确率最大的环节。KNN算法依赖“距离”来衡量样本相似性距离计算又依赖向量表示所以向量怎么构建直接决定了算法上限。初始版本可以采用最简单的方案统计用户在每种电影类型下的观看数量。MovieLens 100K数据集一共有18种电影类型分别是Action、Adventure、Animation、Childrens、Comedy、Crime、Documentary、Drama、Fantasy、Film-Noir、Horror、Musical、Mystery、Romance、Sci-Fi、Thriller、War、Western。那么每个用户就可以被表示成一个18维的整数向量每一维是该类型下的观影次数。光有观影次数还不够因为只看次数会忽略用户的喜好强度。举个例子用户A看了10部爱情片但平均只给了2分用户B看了5部爱情片但平均给了4分显然B对爱情片的喜爱程度远高于A。所以在进阶版本中我把特征从“观看次数”升级为“类型加权评分”也就是对每一维分别统计观看次数和评分总和再用评分总和除以观看次数得到平均评分。最终每个用户被表示成一个36维的向量18个类型的次数维度 18个类型的平均评分维度。特征构建这一步需要写一个MapReduce作业来完成。我在实际项目里是用DistributedCache把movies.dat的映射关系加载到Mapper内存里然后逐条读入ratings.dat进行数据补全最后在Reducer中聚合所有属于同一用户的记录。下面给出Job1的核心代码先看主类框架public class FeatureJob { public static class FeatureMapper extends MapperObject, Text, Text, Text { private MapString, String movieTypeMap new HashMap(); private Text outKey new Text(); private Text outValue new Text(); Override protected void setup(Context context) throws IOException, InterruptedException { // 从DistributedCache中加载电影类型映射 URI[] cacheFiles context.getCacheFiles(); if (cacheFiles ! null cacheFiles.length 0) { Path path new Path(cacheFiles[0]); FileSystem fs FileSystem.get(context.getConfiguration()); try (BufferedReader reader new BufferedReader( new InputStreamReader(fs.open(path), ISO-8859-1))) { String line; while ((line reader.readLine()) ! null) { String[] fields line.split(::); if (fields.length 2) { movieTypeMap.put(fields[0], fields[1]); } } } } } Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(::); if (fields.length ! 4) { return; } String userId fields[0]; String movieId fields[1]; String rating fields[2]; String types movieTypeMap.get(movieId); if (types null) { return; } // 每个电影类型都输出一条携带评分方便Reducer统计 String[] typeArr types.split(\\|); for (String type : typeArr) { outKey.set(userId); outValue.set(type : rating); context.write(outKey, outValue); } } } public static class FeatureReducer extends ReducerText, Text, Text, Text { private Text outValue new Text(); Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { MapString, int[] typeStats new TreeMap(); for (Text val : values) { String[] parts val.toString().split(:); if (parts.length ! 2) { continue; } String type parts[0]; int rating Integer.parseInt(parts[1]); typeStats.computeIfAbsent(type, k - new int[2])[0] 1; // 次数 typeStats.computeIfAbsent(type, k - new int[2])[1] rating; // 总分 } StringBuilder sb new StringBuilder(); for (Map.EntryString, int[] entry : typeStats.entrySet()) { int count entry.getValue()[0]; double avgRating (double) entry.getValue()[1] / count; sb.append(entry.getKey()).append(:).append(count).append(:).append(avgRating).append(,); } outValue.set(sb.toString()); context.write(key, outValue); } } public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, Feature Job); job.setJarByClass(FeatureJob.class); job.setMapperClass(FeatureMapper.class); job.setReducerClass(FeatureReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // 加载movies.dat到DistributedCache job.addCacheFile(new URI(args[1])); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[2])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }这个版本的Reducer输出的是“类型1:次数1:平均分1,类型2:次数2:平均分2,...”这种字符串格式好处是灵活坏处是后续KNN计算时还需要再解析。如果你的项目对性能要求更高可以改为输出固定长度的特征向量每个维度对应一种类型提前把类型列表定义好。2.3 训练集与测试集的拆分策略有了特征向量之后下一步是拆分训练集和测试集。拆分时有个细节很容易被忽略要保证训练集和测试集中男女比例基本一致否则预测结果会偏向于样本量更多的那一类。最稳妥的做法是分层抽样即分别从男性用户和女性用户中按比例随机抽取。在MapReduce项目中还有一个数据组织细节训练集数据要放到DistributedCache中供所有Mapper共享。MapReduce的DistributedCache机制会在任务启动前把文件分发到每一个计算节点的本地磁盘Mapper在setup阶段就能读到避免了每次都从HDFS拉数据的网络开销。训练集通常有几百KB到几MB完全适合放DistributedCache。测试集则是第二个作业的正常输入路径。在评估环节测试集中的用户特征被当作“没有性别标签的新用户”来处理但真实性别被保留在另一份对照文件中用于最后计算准确率。3. KNN算法原理与MapReduce实现3.1 KNN分类原理与距离度量选择KNN算法的数学基础其实非常朴素。假设训练集中有N个已标记样本每个样本是一个d维向量x_i对应标签y_i。给定一个待预测样本x_q算法计算它与所有训练样本的距离选距离最小的K个然后在这K个样本中统计各类别的数量数量最多的类别就是预测结果。距离度量方式直接影响KNN的效果。我对比过三种常见度量方式欧氏距离也就是直线距离公式为d(x, y) sqrt(Σ(x_i - y_i)²)。它直观反映向量在特征空间中的绝对差距对数值大小敏感是KNN最常用的选择。曼哈顿距离公式为d(x, y) Σ|x_i - y_i|对异常值更鲁棒但会弱化多个维度上差距的累加效应。余弦相似度公式为cos(x, y) (x·y) / (|x|·|y|)衡量的是方向上的相似性而不是距离。如果用户特征向量的模长差异很大比如一个用户观影总量是另一个的10倍用欧氏距离会误判为不相似而余弦相似度能规避这个问题。在实际测试中如果特征向量只做次数统计欧氏距离效果尚可但如果加入了评分维度由于评分均值的取值范围是1到5观影次数可能高达几十两者量纲差异很大直接用欧氏距离会使得评分维度几乎不起作用。解决方案有两个一是特征归一化二是改用余弦相似度。我最终选择了先对特征向量做归一化再用欧氏距离这样保留距离的直观性同时让所有维度在同一尺度下参与计算。归一化的做法是每个维度减去该维度的均值再除以标准差也就是Z-score标准化。初次实现时只做简单的最大最小值缩放效果不够稳定因为观影次数呈长尾分布少数活跃用户的观影次数远超普通用户。换成Z-score后准确率提升了约5个百分点。3.2 算法流程与伪代码整个KNN预测的完整流程可以拆解为以下七个步骤读取训练集特征向量解析成内存中的对象列表。读取测试集待预测用户特征向量同样解析成对象。对待预测用户遍历训练集中所有用户计算特征向量的欧氏距离。输出(待预测用户, 距离, 训练用户性别)三元组。按待预测用户ID聚合所有三元组。按距离从小到大排序取前K个。对K个近邻的性别做投票输出票数多的性别。用伪代码可以这样表示对每个待预测用户 u: 初始化一个最小堆容量为K堆顶是当前最大距离 对训练集中每个用户 t: dist 欧氏距离(u.feature, t.feature) 如果堆未满直接插入(t.gender, dist) 否则如果dist小于堆顶距离弹出堆顶插入(t.gender, dist) 统计堆中男性的数量m和女性的数量f 如果m f预测u为男性否则预测u为女性这里用最小堆容量为K始终保持距离最小的K个元素而不是全量排序是为了控制内存消耗。在Reducer端每个用户可能对应上万条距离记录全量排序虽然也能在Reducer的内存里完成但堆结构明显更优雅。3.3 MapReduce作业设计与优化点第二个作业是整个项目的核心。我在设计Mapper时做了一个非常关键的优化把训练集放到DistributedCache中让每个Mapper在setup阶段一次性加载训练集到内存。这样Map阶段每读入一条待预测用户记录就直接在内存中遍历训练集计算距离输出一条聚合了该用户与所有训练用户距离信息的记录。这里要注意一个细节如果所有距离都输出到一个Reducer会出现严重的数据倾斜单个Reducer要处理全量数据完全丧失并行优势。实际项目中我采用了“用户ID取模分桶”的Partitioner策略把不同待预测用户分派到不同的Reducer。由于每个用户的KNN计算是独立的这样可以在不影响正确性的前提下把计算负载分散到多个节点。另外我在计划中没有使用Combiner原因是Reducer端需要在排序后取前K个再做投票而Combiner只在Map端本地聚合如果强行在Combiner阶段就筛选K个近邻会丢失全局信息导致结果不准确。这是一个典型的“为了优化而优化反而出错”的案例读者如果自己做这个项目一定要想清楚Combiner使用的边界。下面是Job2的核心代码先看Mapperpublic class KnnJob { public static class KnnMapper extends MapperObject, Text, Text, Text { private ListTrainUser trainUsers new ArrayList(); private Text outKey new Text(); private Text outValue new Text(); private int K; Override protected void setup(Context context) throws IOException, InterruptedException { Configuration conf context.getConfiguration(); K conf.getInt(knn.k, 5); URI[] cacheFiles context.getCacheFiles(); if (cacheFiles ! null cacheFiles.length 0) { Path path new Path(cacheFiles[0]); FileSystem fs FileSystem.get(conf); try (BufferedReader reader new BufferedReader( new InputStreamReader(fs.open(path), UTF-8))) { String line; while ((line reader.readLine()) ! null) { String[] fields line.split(\t); if (fields.length 2) { continue; } String userId fields[0]; String[] metaAndVector fields[1].split(\\|, 2); if (metaAndVector.length ! 2) { continue; } String gender metaAndVector[0]; double[] vector parseVector(metaAndVector[1]); trainUsers.add(new TrainUser(userId, gender, vector)); } } } } private double[] parseVector(String vectorStr) { String[] dims vectorStr.split(,); double[] vector new double[dims.length]; for (int i 0; i dims.length; i) { vector[i] Double.parseDouble(dims[i]); } return vector; } Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); String[] fields line.split(\t); if (fields.length 2) { return; } String userId fields[0]; // 特征向量部分格式为 gender|v1,v2,v3... 或者 unknown|v1,v2,v3... String[] metaAndVector fields[1].split(\\|, 2); if (metaAndVector.length ! 2) { return; } double[] testVector parseVector(metaAndVector[1]); // 维护一个大小为K的最大堆堆顶是当前最大距离 PriorityQueueNeighbor heap new PriorityQueue(K, (a, b) - Double.compare(b.distance, a.distance)); for (TrainUser trainUser : trainUsers) { double dist euclideanDistance(testVector, trainUser.vector); Neighbor neighbor new Neighbor(trainUser.userId, trainUser.gender, dist); if (heap.size() K) { heap.offer(neighbor); } else if (dist heap.peek().distance) { heap.poll(); heap.offer(neighbor); } } // 输出当前用户的K近邻 outKey.set(userId); StringBuilder sb new StringBuilder(); for (Neighbor neighbor : heap) { sb.append(neighbor.userId).append(:) .append(neighbor.gender).append(:) .append(String.format(%.4f, neighbor.distance)).append(,); } outValue.set(sb.toString()); context.write(outKey, outValue); } private double euclideanDistance(double[] v1, double[] v2) { int len Math.min(v1.length, v2.length); double sum 0.0; for (int i 0; i len; i) { double diff v1[i] - v2[i]; sum diff * diff; } return Math.sqrt(sum); } private static class TrainUser { String userId; String gender; double[] vector; TrainUser(String userId, String gender, double[] vector) { this.userId userId; this.gender gender; this.vector vector; } } private static class Neighbor { String userId; String gender; double distance; Neighbor(String userId, String gender, double distance) { this.userId userId; this.gender gender; this.distance distance; } } } }Mapper的输出已经对每个用户提前筛选出了K个近邻所以Reducer的逻辑就非常简单了。因为我们在测试集中把预测目标当成了“unknown”所以Reducer端只需要对这K个近邻的性别做投票public static class KnnReducer extends ReducerText, Text, Text, Text { private Text outValue new Text(); private String actualGender; Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { int maleCount 0; int femaleCount 0; StringBuilder neighborsInfo new StringBuilder(); for (Text val : values) { String[] neighbors val.toString().split(,); for (String neighbor : neighbors) { if (neighbor.isEmpty()) { continue; } String[] parts neighbor.split(:); if (parts.length 3) { continue; } String gender parts[1]; if (M.equalsIgnoreCase(gender)) { maleCount; } else if (F.equalsIgnoreCase(gender)) { femaleCount; } neighborsInfo.append(neighbor).append(;); } } String predictedGender maleCount femaleCount ? M : F; int total maleCount femaleCount; double confidence total 0 ? 0.0 : (Math.max(maleCount, femaleCount) * 100.0 / total); outValue.set(predictedGender \t String.format(%.2f, confidence) %\t neighborsInfo.toString()); context.write(key, outValue); } }Driver类的设置核心参数包括K值、缓存文件的路径、输入输出路径public static void main(String[] args) throws Exception { if (args.length ! 4) { System.err.println(Usage: KnnJob trainCache input output k); System.exit(-1); } Configuration conf new Configuration(); conf.setInt(knn.k, Integer.parseInt(args[3])); Job job Job.getInstance(conf, KNN Gender Prediction); job.setJarByClass(KnnJob.class); job.setMapperClass(KnnMapper.class); job.setReducerClass(KnnReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); job.addCacheFile(new URI(args[0])); FileInputFormat.addInputPath(job, new Path(args[1])); FileOutputFormat.setOutputPath(job, new Path(args[2])); System.exit(job.waitForCompletion(true) ? 0 : 1); }代码里还埋了一个后续评估准确率的接口测试集在构造特征向量时会额外保留真实性别到元信息中对应代码中的metaAndVector[0]字段虽然KNN作业里没有使用但输出时可以通过对照文件来评估预测准确率。这是一种工程上很实用的技巧把“特征”和“标签”分开存储防止算法在评估时无意中偷看答案。3.4 特征对齐问题维度不一致怎么解决在KNN的距离计算中特征向量的维度必须一致。但实际开发中我发现一个很容易踩的坑Job1的Reducer输出是按类型动态构造的字符串如果某个用户没有看某类电影输出中就直接缺失了这个类型对应的维度。比如用户A的输出是“Action:5:3.2,Romance:2:4.0”用户B的输出是“Action:3:2.8,Comedy:4:3.5”这两个字符串解析出来的维度数量和顺序都不一样直接做距离计算会错位。解决方案是在Job1和Job2之间加一个整理环节把所有的字符串统一转换成固定长度的稠密向量。具体做法是预先定义好18种类型的顺序循环遍历这个类型列表查询当前用户的统计结果没有记录的类型就填充为0。这样每个用户输出的特征向量维度都相同且顺序一致距离计算才有意义。如果你在实现时偷懒用HashMap直接存特征然后计算两个Map的距离结果一定是错的。我第一次跑的时候准确率只有百分之五十几排查了很久才发现是特征对齐的问题。把这个问题修掉之后准确率马上提高了十多个百分点这个坑值得单独记一笔。4. 实操过程与运行效果4.1 环境准备与版本选型开始实操前先把环境说清楚。我的项目在以下环境中完整跑通过读者可以参考不一定要完全一致组件版本JDK1.8Hadoop2.10.2伪分布式模式Maven3.6.3操作系统CentOS 7数据集MovieLens 100K这里提醒一下Hadoop 3.x的API和2.x略有差异比如addCacheFile的用法、部分类所在的包路径如果读者用的是Hadoop 3.3可以参考官方文档做适配。另外JDK版本不建议高于1.8因为Hadoop 2.x对更高版本的JDK兼容性不佳实际运行可能出现一些莫名其妙的反射异常。Maven的pom.xml核心依赖如下dependencies dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-client/artifactId version2.10.2/version /dependency dependency groupIdorg.apache.hadoop/groupId artifactIdhadoop-common/artifactId version2.10.2/version /dependency /dependencies打包插件建议使用maven-shade-plugin它会打出一个包含所有依赖的fat jar省去运行时找依赖的麻烦build plugins plugin groupIdorg.apache.maven.plugins/groupId artifactIdmaven-shade-plugin/artifactId version3.2.4/version executions execution phasepackage/phase goals goalshade/goal /goals /execution /executions /plugin /plugins /build4.2 数据上传与集群准备在运行作业之前需要把数据上传到HDFS。先用命令行创建目录再把本地的数据文件传上去hdfs dfs -mkdir -p /user/hadoop/movie/input hdfs dfs -put ratings.dat /user/hadoop/movie/input/ hdfs dfs -put movies.dat /user/hadoop/movie/input/ hdfs dfs -put users.dat /user/hadoop/movie/input/这里有一个顺序问题第一个作业只需要ratings和movies不需要users。第二个作业需要的是第一个作业的输出以及训练集对应的特征文件。训练集特征文件不需要从外部导入因为它是第一个作业产出的子集。如果你在本地Windows环境跑要额外注意Hadoop的winutils.exe和相关依赖否则会报“Failed to locate the winutils binary”错误。最好还是放到Linux环境里操作省心很多。4.3 运行Job1特征向量构建提交第一个作业的命令如下hadoop jar movie-gender.jar com.example.FeatureJob \ /user/hadoop/movie/input/ratings.dat \ /user/hadoop/movie/input/movies.dat \ /user/hadoop/movie/feature-output运行成功后查看特征输出hdfs dfs -cat /user/hadoop/movie/feature-output/part-r-00000 | head -5输出的每一行大致是这样1 Action:4:3.5,Adventure:3:3.67,Animation:2:4.0,Childrens:2:4.5,Comedy:5:3.8,... 2 Action:2:2.5,Crime:3:3.0,Drama:4:4.25,Romance:6:3.83,...第一列是用户ID第二列是“类型:观看次数:平均评分”的逗号分隔列表。这一阶段输出的文件就是后续KNN计算的基础。4.4 拆分训练集与测试集由于MovieLens的users.dat中有性别标签我可以把特征输出中的用户与users.dat做一次Join把用户分成训练集和测试集。一个简化方案是把特征输出中80%的用户作为训练集剩下20%作为测试集。测试集输入给KNN作业时性别字段标记为unknown。这个拆分既可以在Hive中做也可以用MapReduce或Shell脚本实现甚至本地处理也行因为这里的数据规模很小。训练集需要转换成KNN作业要求的格式每一行是用户ID\t性别|v1,v2,v3,...,v36测试集同样转换只是性别字段写成unknown用户ID\tunknown|v1,v2,v3,...,v36这个转换过程建议放在Job1的输出之后单独写一个小工具完成不要硬塞进Job1里因为训练集和测试集的划分比例会影响评估结果把逻辑拆分出来更容易调整。4.5 运行Job2KNN预测训练集文件上传到HDFS后提交第二个作业hadoop fs -mkdir -p /user/hadoop/movie/train hadoop fs -put train.txt /user/hadoop/movie/train/ hadoop jar movie-gender.jar com.example.KnnJob \ /user/hadoop/movie/train/train.txt \ /user/hadoop/movie/input/test.txt \ /user/hadoop/movie/knn-output \ 7最后一个参数7是K值。运行结束后查看预测结果hdfs dfs -cat /user/hadoop/movie/knn-output/part-r-00000 | head -20输出的每一行格式如下198 F 71.43% 335:M:5.2915;287:F:5.3852;563:M:5.5678;...第一列是待预测用户ID第二列是预测性别第三列是置信度K个近邻中多数性别所占的比例后面就是具体的近邻列表。K7时如果4个近邻是女性3个是男性置信度是57.14%预测为F。4.6 准确率评估与结果分析预测完成之后还需要对照test.txt中的真实性别来计算准确率。可以用下面的ShellAWK方式快速统计hdfs dfs -cat /user/hadoop/movie/knn-output/part-* predict.txt cut -f1,2 predict.txt predict_gender.txt # 将真实标签和预测标签做Join统计在我的测试中使用K7、欧氏距离、36维特征18个次数18个平均评分在MovieLens 100K数据集上准确率可以达到72%左右。这个准确率看起来不是特别高但对于一个纯行为特征的二分类预测来说已经不错了。如果完全随机猜测准确率只有50%。影响准确率的因素有很多训练和测试的用户分布是否一致、特征表达是否充分、距离度量是否合适、K值是否恰当。把这些因素都调好之后准确率还有上升空间但很难超过80%。毕竟电影偏好只是性别的弱关联信号总有一些用户的行为模式和异性群体更相近这是数据本身的局限不是算法的锅。5. 常见问题与排查技巧实录5.1 高频问题排查速查表做这个项目的过程中我整理了以下高频问题的排查方法基本覆盖了从环境配置到结果分析的大部分坑现象可能原因解决方案运行时报ClassNotFoundException没有打fat jar使用maven-shade-plugin打包中文乱码数据文件是ISO-8859-1编码读取时显式指定字符集特征维度对不上Job1输出是动态格式统一按18种类型顺序转成固定维度准确率一直在50%左右特征未对齐或标签泄露检查训练集和测试集的特征构造逻辑DistributedCache文件读不到路径写错或文件权限问题确认HDFS路径存在用hdfs dfs -ls检查Reducer数据倾斜所有用户都分到同一个Reducer增加Partitioner按用户ID取模分桶Mapper内存溢出训练集过大全部加载到内存增加Mapper堆内存或改用子采样测试集性别偷看训练集特征元信息未正确剥离评估阶段把测试集的gender设为unknown多个Reducer输出文件正常现象part-r-xxxxx多个用通配符part-*合并查看KNN结果全是一个性别训练集男女比例失衡做分层采样确保训练集性别平衡5.2 参数调优K值、距离公式与特征组合的实战对比我在项目中做了一组对照实验直观展示各因素对准确率的影响。固定训练集和测试集不变分别调整K值K值准确率369.2%570.8%772.1%970.5%1168.9%K值太小时模型对噪声敏感K值太大会让远处的样本稀释近邻的影响力。在这个数据集上K7是甜点值。这也是为什么一般推荐K取奇数可以避免平票。再看距离度量的影响同样是K7距离度量准确率欧氏距离未归一化63.5%欧氏距离Z-score归一化72.1%曼哈顿距离归一化68.4%余弦相似度用相似度取最大K70.2%这个结果说明数据归一化比距离公式本身的影响更明显。未归一化的欧氏距离会被观影总数这类大数值维度主导归一化后各维度公平参与准确率提升接近9个百分点这个提升幅度是非常可观的。最后看特征组合的影响特征组合准确率仅类型观看次数18维66.8%仅类型平均评分18维64.3%次数 平均评分36维72.1%这里反映出一个核心经验单一维度的表达能力有限。观影次数描述“量”平均评分描述“质”两者是互补关系组合之后准确率明显上升。5.3 数据倾斜与内存优化经验在实际运行中第二个作业的Reducer端可能会出现数据倾斜。这是因为用户ID并不是均匀分布的某些活跃用户的近邻数据量很大。调整Partitioner可以让不同用户ID均匀分布到不同Reducer但如果单纯使用默认的HashPartitioner某个Reducer接收多个大用户时依然可能成为瓶颈。我的优化方案有两层。第一层是在Mapper端使用大小为K的最小堆优先筛选距离最小的K个近邻只输出这K个而不是输出所有距离这样每个用户的数据量从“训练集大小”降到KShuffle的数据量大幅下降。第二层是自定义Partitioner按用户ID哈希后取模把负载分散到多个Reducer。如果你面临更大的数据量可以考虑两阶段KNN先用一组采样训练用户粗选候选集再在候选集上精确计算距离。有点类似于“先粗筛再精排”的思路在推荐系统里很常用。5.4 从实训项目到生产系统的扩展思考最后聊聊这个项目可以怎么延伸。性别预测只是KNN和MapReduce组合的一个演示场景同样的框架完全可以迁移到其他用户画像预测任务上。比如预测用户年龄段、预测用户职业类型、甚至预测用户是否会流失。只要把标签字段替换掉特征工程重新设计一下整体架构无需改动。如果你的数据规模真的到了单机Hadoop也跑不动的程度可以把MapReduce升级为Spark利用RDD的map和reduceByKey等算子实现同样的逻辑。这个项目的MapReduce思想在Spark中完全可以平移理解了MapReduce的“先并行计算再聚合”的思路写Spark版本会非常有条理。从学习价值上看这个项目训练的是“算法、工程、业务”三者结合的思维方式。KNN算法本身很简单但把它放进MapReduce框架、处理分布式环境中的数据对齐问题、设计高效的特征向量格式这些才是真正值钱的经验。我在实际开发里的一个核心体会是机器学习项目的成败往往不在算法本身而在于数据的组织和特征的设计。第一次用KNN做性别预测时我以为重点在调K值、选距离公式结果大量时间花在了数据清洗和特征对齐上。建议所有做这个项目的读者把时间分配向数据倾斜多花一些时间研究特征收益会远超你的预期。本文还有配套的精品资源点击获取