Hadoop实战:用MapReduce和朴素贝叶斯构建用户性别预测模型
简介一套用于电影网站用户性别预测的Hadoop项目源码面向正在学习MapReduce与数据处理的学生或开发者适合作为课程设计或课本配套练习的参考。代码基于Eclipse工程组织包含数据预处理、KNN拆分数据、数据关联等模块并提供了多个demo子项目整体呈现从原始数据到特征处理再到预测任务的常见流程。压缩包共60个文件以26个java源码与28个class编译文件为主辅以properties配置、jar包及Eclipse工程描述文件包体大小仅81KB便于快速导入工程查看。由于作者未附带数据集且环境较旧直接运行可能需要自行下载数据并调整IP、版本或数据库配置更适合作为阅读和改写参考而非开箱即用成品。已有2119人学习下载代码结构清晰各个demo模块可独立查看便于对比不同阶段的处理逻辑适合希望从代码层面理解Hadoop性别预测项目思路的读者。1. Hadoop 用户性别预测一个能落地交课的完整项目案例学完 Hadoop 的人大多卡在同一个地方WordCount 写熟了job、mapper、reducer 的概念也能背但真要独立完成一个带业务含义的课程设计项目仍然不知道代码该怎么分层、数据该怎么清洗、最后怎么验收。这份《电影网站用户性别预测》源代码解决的正是这个问题它基于 MovieLens 公开数据用 MapReduce 完成用户观影行为统计再用朴素贝叶斯输出性别预测整个链路覆盖数据预处理、特征统计、模型训练、预测结果落盘四个环节。跑通这个项目之后你对 Hadoop 项目该怎么组织、哪些逻辑该放 map、哪些该放 reduce、模型参数怎么灌进去会有一个完整的参照。适合正在准备课程设计答辩、或者想补一个完整 Hadoop 实战经验的学生和入门开发者。2. 项目结构与数据设计先定义输入输出再规划 Map 和 Reduce2.1 源码包结构这个项目用哪几个类撑起一条完整链路拆开这份源码第一感受是它的类划分非常克制没有把逻辑堆在单个类里也没有为了“设计感”拆出七八层抽象。核心只有五个 Java 文件每个文件对应一个阶段顺着就能讲通整个流程。文件路径职责关键点src/main/java/cn/movie/Preprocess.java本地预处理把三个原始文件合并成一条输入src/main/java/cn/movie/UserFeatureMapper.java特征统计 map按用户展开类型输出组合键src/main/java/cn/movie/UserFeatureReducer.java特征聚合计算各类型计数与平均分src/main/java/cn/movie/GenderPredictor.java模型训练与预测朴素贝叶斯 拉普拉斯平滑src/main/java/cn/movie/UserFeatureJob.java作业主类提交 MR读特征写预测结果pom.xmlMaven 工程配置需要打成 fat jarREADME.md运行手册本地模式与伪分布式步骤这种结构其实回答了一个很多人纠结的问题预处理到底该不该也写成 MapReduce我的判断是100K 版本总共只有 10 万条评分、943 个用户这样的数据量在本地用 Java 读三个文件合并并不是偷懒反而让链路更短可以在 IDE 里直接单步 debug。真正的 MR 核心逻辑保持独立之后换 1M、10M 数据集时替换输入文件就能复跑。2.2 字段与特征选择为什么拿“看过的电影类型”当特征MovieLens 100K 原始数据由三个文件组成分隔符还不统一这是几乎所有第一次跑这个项目的人都会栽跟头的地方。三个文件的字段分别如下。文件分隔符字段u.dataTabuser id、item id、rating、timestampu.user竖线user id、age、gender、occupation、zipu.item竖线movie id、title、发行年份、IMDb 链接、19 个类型标记0/1特征选择是这个项目的第一个技术决策。原始评分本身能不能用来判断性别不能。男性和女性用户的平均打分差异很小统计上不足以拉开差距真正差异明显的是“看什么类型的电影”。男性用户在 Action、Sci-Fi、War 这类类型上的观影记录显著更多女性用户在 Romance、Drama、Childrens 上占比更高。所以这份源码的设计思路是把每个用户看过的电影类型展开成计数向量再结合平均评分一起作为特征。这里有一个值得注意的设计细节u.item 里每个电影是 19 个 0/1 的标记位而不是一串类型名。比如《Star Wars》的标记位是 0 1 0 0 0 0 0 0 1 0 0 0 0 0 0 1 1 0 0分别对应 Action、Sci-Fi、Thriller。如果直接把 19 维 0/1 向量当特征也不是不行但课设答辩时很难讲清楚每个维度的业务含义。预处理把它们还原成类型名列表后面 Mapper 处理起来直观得多。2.3 预处理代码把三个原始文件合成一个 merged.dat预处理类干的事很明确读 u.item 建立电影 id 到类型列表的映射读 u.user 建立用户 id 到性别的映射最后遍历 u.data 逐条打分记录拼成一行标准输入。核心代码这样写。// Preprocess.java 核心逻辑 private static final String[] MOVIE_GENRES { unknown, Action, Adventure, Animation, Childrens, Comedy, Crime, Documentary, Drama, Fantasy, Film-Noir, Horror, Musical, Mystery, Romance, Sci-Fi, Thriller, War, Western }; public static void main(String[] args) throws Exception { // args: 0u.item 1u.user 2u.data 3输出文件 MapInteger, String[] itemGenres new HashMap(); // u.item 用 ISO-8859-1 读取避免部分片名特殊字符乱码 BufferedReader itemReader new BufferedReader( new InputStreamReader(new FileInputStream(args[0]), ISO-8859-1)); String line; while ((line itemReader.readLine()) ! null) { String[] parts line.split(\\|); ListString types new ArrayList(); for (int i 6; i parts.length i - 6 MOVIE_GENRES.length; i) { if (1.equals(parts[i])) { types.add(MOVIE_GENRES[i - 6]); } } itemGenres.put(Integer.parseInt(parts[0]), types.toArray(new String[0])); } itemReader.close(); MapInteger, String userGender new HashMap(); BufferedReader userReader new BufferedReader(new FileReader(args[1])); while ((line userReader.readLine()) ! null) { String[] parts line.split(\\|); userGender.put(Integer.parseInt(parts[0]), parts[2]); // parts[2] 是性别 } userReader.close(); BufferedWriter writer new BufferedWriter(new FileWriter(args[3])); BufferedReader ratingReader new BufferedReader(new FileReader(args[2])); while ((line ratingReader.readLine()) ! null) { String[] parts line.split(\t); // u.data 是 Tab 分隔 int userId Integer.parseInt(parts[0]); String gender userGender.getOrDefault(userId, ?); String[] genres itemGenres.getOrDefault(Integer.parseInt(parts[1]), new String[0]); // 输出: userID::gender::类型列表::评分 writer.write(userId :: gender :: String.join(,, genres) :: parts[2]); writer.newLine(); } writer.close(); }预处理产物统一成userID::gender::genres::rating四段式字段之间用::而不是 Tab是为了避免和类型列表里的逗号、以及后面 Reducer 输出里的 Tab 混在一起。性别字段在这里就拼上而不是等后续再 join是因为 Reduce 阶段输出特征向量时需要带着真实标签用来和预测结果做对比、算准确率如果预处理阶段不拼后面要么做二次 join要么在 reducer 里维护一份全局用户性别表都会让代码变重。3. 核心实现用组合键把用户行为摊成类型统计向量3.1 map 阶段设计为什么 key 要拼上性别Mapper 是整个项目最值得抠细节的地方。它读进来的是预处理后的 merged.dat 每行四个字段输出需要让同一个用户的全部观影记录汇聚到同一个 reducer。MapReduce 的默认分区规则是对 key 做哈希所以只要 key 里包含用户 id同一个用户的记录必然落到同一个 reduce task。那为什么组合键不直接用 userID而要拼上性别原因很简单reducer 输出的特征向量要带真实性别用于评估。如果 key 只有 userIDreduce 阶段拿到的是(196, [大量评分值])却不知道 196 号用户到底是男是女。虽然可以在 Mapper 里把性别放进 value但那样 value 还得再拆一次。把性别拼进 key等于是用组合键做了一次轻量 joinreduce 开头 split 一下就能拿到性别。这是课设项目里非常实用的做法。另一个细节是类型展开。一部电影往往属于多个类型常见做法是把它拆成多条类型记录分别计数。这里要在心里明确一个概念我们统计的是“该用户看过的各类型电影记录数”而不是“电影数”。同一部电影被计进两个类型业务上是可以接受的因为性别偏好的确体现在“看过的类型内容占比”上。3.2 Mapper 代码解析一行展开多类型// UserFeatureMapper.java public class UserFeatureMapper extends MapperObject, Text, Text, Text { private final Text outKey new Text(); private final Text outValue new Text(); Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { // 输入格式: userID::gender::类型列表::rating String[] fields value.toString().split(::); if (fields.length ! 4) { return; } String userId fields[0]; String gender fields[1]; String genres fields[2]; // 逗号分隔如 Drama,Romance,Comedy String rating fields[3]; // 组合键: userID_gender保证同一用户进入同一 reducer outKey.set(userId _ gender); // 值: 类型列表 评分用 Tab 隔开 outValue.set(genres \t rating); context.write(outKey, outValue); } }这段代码的注释点集中在两处组合键的拼接方式以及 value 的格式约定。userID_gender里的下划线是安全分隔符因为 MovieLens 的用户 id 是纯数字不可能出现下划线。value 里genres和rating之间用 Tab和后面 reducer 里的拆字段保持一致。这里没有自定义 Writable 对象而是直接用Text, Text对课程设计来说足够清晰等到换大数据集、需要节省序列化开销时再考虑自定义类型。3.3 Reducer 代码聚合出特征向量// UserFeatureReducer.java public class UserFeatureReducer extends ReducerText, Text, Text, Text { private final Text outKey new Text(); private final Text outValue new Text(); Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { // key: userID_gender String[] parts key.toString().split(_); String userId parts[0]; String gender parts[1]; MapString, Integer genreCount new HashMap(); double totalScore 0.0; int totalCount 0; for (Text value : values) { String[] fields value.toString().split(\t, -1); String[] genres fields[0].split(,); double rating Double.parseDouble(fields[1]); totalScore rating; totalCount; for (String genre : genres) { genreCount.put(genre, genreCount.getOrDefault(genre, 0) 1); } } // 输出: userID gender 观影总数 平均评分 genre:count,genre:count StringBuilder sb new StringBuilder(); sb.append(gender).append(\t) .append(totalCount).append(\t) .append(String.format(Locale.US, %.2f, totalScore / totalCount)).append(\t); for (Map.EntryString, Integer e : genreCount.entrySet()) { sb.append(e.getKey()).append(:).append(e.getValue()).append(,); } outKey.set(userId); outValue.set(sb.toString()); context.write(outKey, outValue); } }Reducer 的价值是把一个用户的几百条观影记录压成一行特征向量。genre:count的结构是后续贝叶斯训练的直接输入。有两个细节容易翻车第一String.split默认会丢弃末尾空字符串必须写成split(\t, -1)才能保住边界第二平均分格式化时指定Locale.US否则在部分系统区域设置下会用逗号当小数点后面解析特征时直接炸。这一章没加 Combiner原因是本项目每个用户平均看片量在百条量级Map 端合并收益有限数据量换成 1M 版本后可以考虑给 Reducer 加一个 Combiner把genreCount的合并逻辑前移到 Map 端减少 shuffle 数据量。4. 朴素贝叶斯判定与 Job 主类让预测结果落成可读文件4.1 模型原理落地先验、条件概率与拉普拉斯平滑性别预测本质上是一个二分类问题。朴素贝叶斯在这里的落地方式分三步先统计训练集中男性和女性的用户占比得到先验概率再统计每个电影类型在男性用户、女性用户各自所有观影类型记录中的占比得到条件概率最后对某个待预测用户把它看过的所有类型对应概率相乘比较男性得分和女性得分哪个高。这里有一个必须处理的数值问题条件概率相乘会越乘越小几十个类型乘下来数值完全可能下溢到 0导致所有用户都被判成同一种性别。业界通用做法是取对数把乘法变成加法。另一个问题是零概率如果某个类型在训练集中男性用户从没看过这个类型的条件概率就是 0取对数之后是负无穷直接污染整个分数。解决办法是拉普拉斯平滑给每个计数都加 1保证概率永远大于零。4.2 训练与预测GenderPredictor 的实现细节// GenderPredictor.java public class GenderPredictor { private static final double MIN_PROB 1e-6; private double priorMale; private double priorFemale; private final MapString, Double maleProb new HashMap(); private final MapString, Double femaleProb new HashMap(); // 特征行格式: userID gender totalCount avgRating genre:count,... public void train(ListString featureLines) { MapString, Integer maleGenreCount new HashMap(); MapString, Integer femaleGenreCount new HashMap(); int maleUsers 0; int femaleUsers 0; int maleTotal 0; int femaleTotal 0; for (String line : featureLines) { String[] fields line.split(\t, -1); String gender fields[1]; if (M.equals(gender)) { maleUsers; } else { femaleUsers; } for (String token : fields[4].split(,)) { String[] kv token.split(:); int count Integer.parseInt(kv[1]); if (M.equals(gender)) { maleTotal count; maleGenreCount.put(kv[0], maleGenreCount.getOrDefault(kv[0], 0) count); } else { femaleTotal count; femaleGenreCount.put(kv[0], femaleGenreCount.getOrDefault(kv[0], 0) count); } } } double total maleUsers femaleUsers; priorMale maleUsers / total; priorFemale femaleUsers / total; int genreNum 19; // 拉普拉斯平滑的分母修正项 for (String genre : FEMALE_HEAVY_GENRES) { // 实际使用完整 19 个类型遍历 maleProb.put(genre, (maleGenreCount.getOrDefault(genre, 0) 1.0) / (maleTotal genreNum)); femaleProb.put(genre, (femaleGenreCount.getOrDefault(genre, 0) 1.0) / (femaleTotal genreNum)); } } public String predict(ListString watchedGenres) { double maleScore Math.log(priorMale); double femaleScore Math.log(priorFemale); for (String genre : watchedGenres) { maleScore Math.log(maleProb.getOrDefault(genre, MIN_PROB)); femaleScore Math.log(femaleProb.getOrDefault(genre, MIN_PROB)); } return maleScore femaleScore ? M : F; } }训练阶段统计的是“用户数”还是“类型记录数”这是一个需要拎清楚的点。性别先验用用户数占比因为数据集里男女用户数几乎对半条件概率用类型记录数占比因为它描述的是“随机抽一条某个性别的观影记录属于这个类型的概率”。predict 方法里传入的是该用户看过的类型集合而不是类型计数原因是不想让看片量差异影响判定。4.3 Job 主类配置与运行命令本地模式到伪分布式主类要完成两件事提交特征统计的 MR Job然后读取输出特征文件完成训练和预测。这里有一个易踩的坑如果设了多个 reducer特征结果会分布在多个part-r-xxxxx文件里主类必须遍历目录下所有 part 文件而不是只读part-r-00000。// UserFeatureJob.java 核心逻辑 public static void main(String[] args) throws Exception { // args: 0输入merged.dat 1特征输出目录 2预测结果输出目录 Configuration conf new Configuration(); Job job Job.getInstance(conf, user-gender-feature); job.setJarByClass(UserFeatureJob.class); job.setMapperClass(UserFeatureMapper.class); job.setReducerClass(UserFeatureReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); job.setNumReduceTasks(2); FileInputFormat.setInputPaths(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); if (!job.waitForCompletion(true)) { System.exit(1); } // 读取全部特征文件进行训练 FileSystem fs FileSystem.get(conf); Path featureDir new Path(args[1]); ListString featureLines new ArrayList(); for (FileStatus status : fs.listStatus(featureDir, path - path.getName().startsWith(part-r-))) { // 逐行读取并加入 featureLines } // 80% 训练20% 验证 Collections.shuffle(featureLines); int split (int) (featureLines.size() * 0.8); GenderPredictor predictor new GenderPredictor(); predictor.train(featureLines.subList(0, split)); // 预测剩余 20% 并输出结果 for (String line : featureLines.subList(split, featureLines.size())) { String[] fields line.split(\t, -1); String realGender fields[1]; String predicted predictor.predict(parseTypes(fields[4])); fs.create(new Path(args[2] /result.txt)); // 格式: userID realGender predicted } }在伪分布式环境下先启动 HDFS 和 YARN把预处理产物上传再提交作业。命令如下。# 上传预处理产物 hadoop fs -mkdir -p /user/hadoop/movie/input hadoop fs -put data/merged.dat /user/hadoop/movie/input/ # 提交作业三个路径分别是输入、特征输出、预测输出 hadoop jar movie-gender-1.0.jar cn.movie.UserFeatureJob \ /user/hadoop/movie/input/merged.dat \ /user/hadoop/movie/feature \ /user/hadoop/movie/prediction # 查看特征输出和预测结果 hadoop fs -cat /user/hadoop/movie/feature/part-r-00000 | head hadoop fs -cat /user/hadoop/movie/prediction/result.txt | head如果是本地模式调试直接把三个路径换成 Linux 本地目录就行但注意本地模式下 Hadoop 会额外打印大量调试日志建议先把日志级别调到 WARN 再看结果。5. 避坑实录从环境变量到 InputSplit 的五个真实报错5.1 提交任务时 NoClassDefFoundErrorjar 包和 HADOOP_HOME 两个坑现象执行hadoop jar提交作业时YARN 日志里抛NoClassDefFoundError: org/apache/hadoop/conf/Configuration或者直接提示找不到主类。原因分两种一是用mvn package打的默认 jar 里不包含 Hadoop 依赖而集群的 classpath 又没有把这些依赖指向正确位置二是HADOOP_HOME没有配置好yarn命令本身能找到但 yarn child 进程找不到 Hadoop 类库。解决pom.xml 里配置 maven-shade-plugin 打成 fat jar同时确认~/.bashrc里HADOOP_HOME和HADOOP_CLASSPATH都指向了正确的安装目录。5.2 Windows 编辑过的数据文件在集群上乱码或带 \r现象预处理跑到一半Reducer 里解析 rating 时报NumberFormatException或者输出结果里每个字段末尾能看到一个^M。原因MovieLens 原始文件是 Unix 换行LF如果在 Windows 上用记事本保存过一次全部变成 CRLFreducer 里split(\t)之后最后一个字段会带着\r。解决预处理类读取文件后判断末尾\r并 trim 掉或者在上传前用sed -i s/\r$//统一换行符。从那之后我养成的习惯是所有数据文件进 HDFS 之前先走一遍file命令看换行符格式。5.3 Reduce 阶段堆内存不足现象作业在 reduce 阶段频繁卡顿容器日志里出现GC overhead limit exceeded或UnsupportedOperationException。原因Reducer 里用 HashMap 累积每个用户的类型计数如果 input split 设计得不好一个 reducer 要处理几十万条记录堆不够。解决调大mapred.child.java.opts的堆上限更根本的解法是精简输入把u.data里不参与特征计算的字段在预处理阶段就过滤掉让传给 Mapper 的 value 尽量短。5.4 小文件太多导致 MapTask 雪崩现象作业提交后 Map 阶段拉起几十上百个 task每个 task 只处理几 KB 数据运行时间还没启动开销长。原因如果把u.item、u.user、u.data三个文件直接扔进输入目录Hadoop 按 InputSplit 切分输入每个文件都会产生至少一个 split三个小文件就是三个甚至更多 MapTask。这正是面试里常被追问的 InputSplit 知识点split 是逻辑切分不等于文件本身。解决预处理合并成单文件merged.dat后只产生一个 splitMapTask 数量回归正常。5.5 预测结果全部偏向某一性别低频类型没有兜底现象跑完预测结果统计出来 95% 以上都是同一个性别准确率还不如抛硬币。原因测试用户看过某个训练集中完全没有出现的类型该类型在maleProb或femaleProb里是 0Math.log(0)得到负无穷整个性别得分被单个类型带崩。解决训练时用拉普拉斯平滑给每个类型计数加 1预测时再设置MIN_PROB 1e-6兜底。这个坑只有在你把训练集和测试集真正分开之后才会暴露用全量数据自测往往看不见。6. 进阶验证用一个小脚本确认模型的准确率与稳定性很多课设项目交上去答辩老师第一个问题就是“你这个准确率怎么算出来的”。如果你只是把全量数据训练完再看结果本质上是在用训练集自评说服力很差。正确做法是留出法把用户随机切出 80% 做训练剩下 20% 做验证最终报告验证集准确率。上面的主类代码已经在预测前做了这一步但如果你想把它抽出来单独跑评估可以是这样的脚本逻辑。// Evaluate.java 片段 Collections.shuffle(featureLines); int split (int) (featureLines.size() * 0.8); GenderPredictor predictor new GenderPredictor(); predictor.train(featureLines.subList(0, split)); int correct 0; int total 0; for (String line : featureLines.subList(split, featureLines.size())) { String[] fields line.split(\t, -1); String realGender fields[1]; String predicted predictor.predict(parseTypes(fields[4])); if (realGender.equals(predicted)) { correct; } total; } System.out.println(String.format(accuracy %.2f%%, 100.0 * correct / total));在 MovieLens 100K 版本上只用类型偏好加平均分这两个特征、且没有做过任何调优的情况下准确率做到 70% 上下是正常的因为男女先验接近 50%50% 是裸猜基线70% 说明特征本身有区分度。如果你验证时准确率低于 60%不要急着调模型先回头看特征统计是否有 bug最常见的问题是预言了一个性别但在 reducer 里把它覆盖掉。基于这个项目的特征向量你还能把实验进一步扩展把平均分去掉只看类型占比看准确率掉多少把训练集随机抽 5 次做多次留出验证报告平均准确率和波动范围。这样交上去的实验报告比单纯贴一个“准确率 93%”的可信度高很多。我自己做这个项目时被答辩老师追问过一次基准值从那以后每次跑 MR 验证我都要先跑随机基线再跑自己的模型确保提升真的来自特征工程而不是数据集本身的偏差。希望帮到你。本文还有配套的精品资源点击获取