基于MapReduce和朴素贝叶斯的电影用户性别预测实战

发布时间:2026/8/29 3:39:27
基于MapReduce和朴素贝叶斯的电影用户性别预测实战 简介在大数据场景下海量用户行为日志的存储与计算是构建用户画像的关键。分布式计算框架Hadoop及其核心编程模型MapReduce为处理此类数据提供了可扩展的解决方案。面对性别预测等分类问题朴素贝叶斯算法因其原理简单、计算高效且易于并行化常被用于工程实践。本文基于电影网站用户浏览日志与注册信息采用MapReduce实现朴素贝叶斯分类器的完整训练与预测流程涵盖特征工程、先验概率与条件概率统计、模型部署及数据倾斜等优化细节。该方法不仅适用于课程设计也为真实场景下基于分布式的用户属性预测提供了可参考的实践路径。 前阵子帮人看了一个Hadoop课程设计的项目电影网站用户性别预测。需求很常见训练数据是用户的浏览日志和基本注册信息目标是预测那些没填性别的用户是男是女。很多人拿这类题目第一反应是用Pandas跑逻辑回归但对于真正的课程设计和入门Hadoop的项目来说用MapReduce自己实现一个可解释的分类器反而更能体现对分布式计算的理解。这篇就把我当时拆分任务、设计特征、写Mapper/Reducer、以及部署过程中踩过的坑完整复盘一遍代码基于Hadoop 3.x环境给出的是能直接改路径跑起来的核心结构。1. 需求拆解电影网站的用户性别预测到底在算什么1.1 项目背景与数据形态题目里提到的“电影网站用户性别预测”本质上是一个二分类问题。样本是用户标签是性别特征是用户在这个网站上的行为痕迹。实战项目里训练数据一般会包含这样几类信息注册信息user_id、gender训练集有预测集没有、age、注册渠道等浏览行为看过的电影id、浏览的影片类型、浏览时间、停留时长评分行为给多少部电影打过评分、平均分、评分分布时段特征活跃时段集中在白天还是后半夜。我在项目里用的训练数据格式是TSV一行一个用户字段用\t分隔大致长这样user_id gender age view_genres avg_rating active_hours 10001 M 25 动作,科幻,犯罪 3.8 23,0,1 10002 F 28 爱情,剧情,喜剧 4.2 20,21,22 10003 31 动作,战争,历史 3.1 22,23注意预测集里gender字段是空的这些用户可能注册时没有填写性别也可能是在后续使用中被判定为未知性别。我们做的事情就是给这些行填上一个M或F的预测值。1.2 为什么答案不是一台电脑跑Python很多第一次接触这个题目的朋友会问数据量就几十MB甚至课程设计给的数据集可能只有几万行为什么非要上Hadoop这背后有两个层面的理由。第一个层面是课程设计的考察点。既然题目明确写了“基于Hadoop”那么核心就是要你用分布式计算的思路来解决问题而不是教你用Excel。第二个层面是真实场景的考量。在真实业务中一个视频网站的用户行为日志一天就能产生几个GB而且日志分散在多台服务器的HDFS上单机Python脚本要么加载不了全部数据要么只能抽样处理。用Hadoop的好处是计算跟着数据走处理几GB的文件和几十MB的文件代码逻辑几乎不用变调整资源就能扩展。这里还有一个实际原因预测任务本身需要进行大量“统计频次”的操作比如统计男性用户看动作片的比例、女性用户看爱情片的比例。这种按key聚合的计数操作恰好是MapReduce最擅长的事情。1.3 选型Hadoop生态里谁负责干什么这个项目不用引入Spark、Flink等复杂框架也不用写Hive SQL来做完整流程。我的选型思路是HDFS负责存储原始日志和模型文件MapReduce负责三个核心计算任务统计性别先验概率、统计特征条件概率、批量预测可选Hive用来做数据清洗和格式转换如果原始数据不是规整的TSV用Hive的LATERAL VIEW配合UDTF拆解影片类型会很方便模型最终是一个文本文件每行是“特征名\t性别\t概率”不需要像深度学习那样保存大文件。简单说这个项目用到的Hadoop知识集中在MapReduce编程模型本身适合用来理解分布式计算的核心思想。2. 特征设计从浏览日志和评分记录里提炼判断依据2.1 原始字段长什么样这个项目的预测效果好不好特征设计占了80%的权重。如果直接把user_id送进模型没有任何意义。如果只用性别先验概率预测那所有人的预测结果都一样毫无区分度。所以要找那些和性别相关性强的行为特征。我做特征工程时先把每个用户的行为日志聚合到一行上。这部分如果数据量大用Hive或Pig来聚合比较合适如果数据量小一个Python脚本先做预处理也行。核心是聚合出以下几类原始特征浏览过的电影题材列表可能出现多个题材比如“动作,科幻”浏览电影的总数量、平均评分、评分标准差活跃时段按小时统计把用户最常访问的时段档位提取出来浏览设备或访问来源有些数据源里会记录UA信息移动端和PC端的比例也能作为特征。我在实践里发现电影题材偏好是最有区分度的特征之一但也最容易导致数据倾斜。因为用户的题材特征是列表不是单一值后面处理时要拆开统计。2.2 构造特征题材偏好、活跃时段、评分倾向以我用的训练数据为例每个用户一行字段已经聚合好。我从里面构造了四组特征第一组是题材偏好特征。把view_genres字段用逗号拆开变成多个特征项。比如“动作,科幻,犯罪”就会拆成view_genres动作、view_genres科幻、view_genres犯罪三个特征。在朴素贝叶斯里假设这些特征相互独立分别计算“男性中动作片比例有多高女性中动作片比例有多高”。第二组是活跃时段特征。把active_hours字段按小时档位划分成凌晨0-5、上午6-11、下午12-17、晚间18-23四类然后取用户活跃最多的一个档位作为特征。比如“23,0,1”会归类为凌晨这是一个很典型的男生特征。第三组是评分倾向特征。avg_rating离散化成几个区间3分以下、3到4分、4分以上。这个特征的作用是区分用户是“重度评论型用户”还是“随便逛逛型用户”和性别的相关度虽然不如题材偏好但能让模型多一个维度。第四组如果数据里有age也可以离散化成年龄段。不过注意age在真实场景中经常缺失后面会单独处理。特征构造完成后每一条训练样本变成类似这样的结构10001 M gender_count feature_view_genres动作 feature_view_genres科幻 feature_active凌晨 feature_rating3.8这里为了输入给MapReduce我把一行样本拆成多行每行是“user_id, gender, feature_name”后续Map的输入就可以直接处理这种扁平化格式。当然也可以保留一行多特征在Mapper里自己拆分两种方式都行。2.3 特征离散化与缺失值处理MapReduce上的朴素贝叶斯模型要求特征值尽量是离散的因为连续值的概率密度估计在分布式环境里比较麻烦。所以我们把连续值做离散化比如avg_rating用区间分桶age用年龄段分桶。这样做还有一个额外好处离散化后的特征对异常值不敏感。比如某个用户的评分是9999显然不合理但分桶后只会被放到最高区间不会被当作极端值干扰计算。缺失值处理有两种思路。第一种是直接丢弃缺失特征的用户但这样会损失训练样本。第二种是给缺失特征单独分配一个枚举值比如“ageunknown”把它当成一个普通特征来统计概率。我推荐第二种因为在实际数据中性别已知的用户中也有不少缺失年龄直接丢掉太可惜。另外在预测阶段如果某个测试用户缺失了某个特征那么朴素贝叶斯公式里就跳过这一项只计算已有特征的概率不把缺失特征当作0概率处理。3. 朴素贝叶斯分类器与MapReduce的结合点3.1 朴素贝叶斯的一页纸数学项目里用的分类器是朴素贝叶斯数学上非常简单适合在MapReduce里并行实现。我们要计算的是P(性别|特征1, 特征2, ..., 特征n) ∝ P(性别) × P(特征1|性别) × P(特征2|性别) × ... × P(特征n|性别)P(性别) 是先验概率直接从训练集里统计男性和女性的人数比例得到。P(特征i|性别)是条件概率表示在男性或女性样本中出现某个特征比如“喜欢动作片”的比例。预测时对每个候选性别把这个性别下所有特征的条件概率连乘起来再乘上先验概率最后选概率更大的那个性别作为预测结果。因为特征值是离散的条件概率的估计公式是P(特征i某个值|性别) (该性别下出现这个特征值的样本数 1) / (该性别总样本数 特征取值个数)这里加1的平滑处理叫拉普拉斯平滑目的是防止训练集中没有出现过的特征组合在预测时概率直接变成0。3.2 训练过程如何拆成MR作业整个训练过程拆成了三个MapReduce作业每个作业只负责一个清晰的统计任务。第一个作业统计性别先验概率。Map端读取训练数据输出性别和计数1Reducer求和得到男、女各有多少样本再除以总样本数得到P(M)和P(F)。第二个作业统计条件概率。Map端读取一行训练数据性别作为key的一部分每个特征值作为key的一部分输出“性别_特征名特征值”和计数1。Reducer求和然后除以该性别总样本数得到P(特征|性别)。这个作业就是整个模型的核心。第三个作业用训练好的模型做预测。Map端读取待预测用户数据在setup阶段把模型文件加载到内存对每个用户计算两个性别的后验概率输出用户ID和预测结果。为什么拆成三个而不是把前两个合并从理论上看可以合并在同一个Reducer里既统计性别计数又统计条件概率但代码复杂度会增加不少。而且拆开后每份代码都只有一两个核心类调试起来非常方便。对课程设计或者初学者来说宁可多跑两个作业也要保证逻辑清晰可查。3.3 为什么不用一条SQL聚合直接搞定有人会说统计条件概率这个操作用Hive SQL一条GROUP BY就出来了为什么还要写Mapper和Reducer这里我要说句公道话如果目标是只得到条件概率表Hive确实更方便。但问题在于这个项目核心是理解MapReduce不是最快得到结果。写Hive SQL的过程中你接触不到Shuffle、Combiner、Partitioner这些概念。而真正在Hadoop上做模型训练时哪怕只是一个朴素贝叶斯也需要理解“为什么作业要分多个轮次”“为什么数据要跨节点传输”“什么时候该加Combiner”。这些理解只能靠手写Java代码来积累。另外预测阶段用SQL并不容易表达。给每个测试用户计算连乘概率SQL写起来会非常别扭尤其在不知道每个用户有多少特征的情况下很难用一个固定列的SQL完成。MapReduce里用Map遍历特征列表逻辑就自然很多。4. 核心源代码三个Job的前后端实现4.1 Job1性别先验概率统计先看第一个作业它的作用只有一个统计训练集中男性和女性用户各有多少人。Mapper的输出key是性别value是常量1。Reducer做求和。public class PriorCountMapper extends MapperLongWritable, Text, Text, IntWritable { private final Text genderKey new Text(); private final IntWritable one new IntWritable(1); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\t); if (fields.length 2) { return; } String gender fields[1].trim(); if (M.equals(gender) || F.equals(gender)) { genderKey.set(gender); context.write(genderKey, one); } } }Reducer如下public class PriorCountReducer extends ReducerText, IntWritable, Text, IntWritable { private final IntWritable total new IntWritable(); Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } total.set(sum); context.write(key, total); } }这个作业跑完后输出目录下会看到类似这样的结果F 3421 M 4652这就是全量样本的男女分布。在后续条件概率计算中会把这个结果作为配置参数传入避免再单独读一次文件。4.2 Job2条件概率表构建第二个作业是最关键的部分。Mapper读取的是已经扁平化过的样本数据每行格式为user_id gender feature_namefeature_value这里涉及一个细节一行用户可能有多条特征记录在预处理阶段已经把同一个用户的多个特征拆分成了多行。如果不想拆分也可以在Mapper里用自定义分隔符取出整个特征列表再逐个循环输出两种方式殊途同归但拆分后对后续写Combiner更友好。Mapper代码public class ConditionProbMapper extends MapperLongWritable, Text, Text, IntWritable { private final Text outKey new Text(); private final IntWritable one new IntWritable(1); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\t); // 期望格式: user_id gender feature_namefeature_value if (fields.length 3) { return; } String gender fields[1].trim(); if (!M.equals(gender) !F.equals(gender)) { return; } String feature fields[2].trim(); // 组合key形如: M_feature_view_genres动作 outKey.set(gender _ feature); context.write(outKey, one); } }Reducer统计每个“性别特征”组合的出现次数。这里要注意实际项目中这还没算完条件概率因为条件概率还需要除以对应性别的总样本数。我选择在Reducer里不除总数直接把计数输出然后由后续读取模型文件的代码来计算概率。这样做的好处是Runner可以用DistributedCache把先验概率文件传给作业在Reducer里结合总数一次性算出概率并输出模型文件。如果你在Reducer里拿不到先验概率也可以让Reducer只输出“gender_feature\tcount”然后在第三个Job加载模型时统一除以性别总数。我在项目里采用后一种方式模型文件读取逻辑更集中。Reducer代码public class ConditionProbReducer extends ReducerText, IntWritable, Text, DoubleWritable { private final DoubleWritable prob new DoubleWritable(); private double maleCount; private double femaleCount; Override protected void setup(Context context) throws IOException, InterruptedException { Configuration conf context.getConfiguration(); maleCount Double.parseDouble(conf.get(male.count, 1)); femaleCount Double.parseDouble(conf.get(female.count, 1)); } Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } String keyStr key.toString(); String gender keyStr.substring(0, 1); double total M.equals(gender) ? maleCount : femaleCount; // 拉普拉斯平滑的分母需要加上特征取值个数这里简化处理加1即可 double p (sum 1.0) / (total 1.0); prob.set(p); context.write(key, prob); } }为了在Reducer里拿到男女总数我在Driver里先读一次Job1的输出把两个总数放进Configuration传给Job2。具体做法是conf.set(male.count, maleCount)。Configuration传字符串没问题但要注意MapReduce作业真正调用时这些配置值会通过网络传到每个Task节点所以只适合传小体量参数不要传大文件。4.3 Job3用模型文件做预测第三个作业是预测。预测集的每行是一个待预测用户的数据格式和训练集类似只是性别字段为空20001 27 爱情,剧情 3.9 20,21,22预测的思路是在Mapper的setup阶段从HDFS读到模型文件Job2的输出加载进一个Map结构。然后对每行数据拆出特征列表对每个性别分别计算连乘概率。由于朴素贝叶斯本质上是求后验概率最大的类别不需要归一化直接比较P(M)和P(F)的大小即可。为了避免下溢出实际计算时会对连乘结果取对数变成累加。但为了代码易读且课程设计数据量不大我在下面示例中直接用了乘法。关键代码public class PredictMapper extends MapperLongWritable, Text, Text, Text { private final Text outKey new Text(); private final Text outValue new Text(); private final MapString, Double probMap new HashMap(); private double priorMale; private double priorFemale; Override protected void setup(Context context) throws IOException, InterruptedException { Configuration conf context.getConfiguration(); priorMale Double.parseDouble(conf.get(prior.male, 0.5)); priorFemale Double.parseDouble(conf.get(prior.female, 0.5)); URI[] cacheFiles context.getCacheFiles(); if (cacheFiles ! null) { for (URI path : cacheFiles) { Path filePath new Path(path); String fileName filePath.getName(); try (BufferedReader reader new BufferedReader( new InputStreamReader(new FileInputStream(fileName)))) { String line; while ((line reader.readLine()) ! null) { String[] parts line.split(\t); if (parts.length 2) { String featureKey parts[0]; // 例如 M_feature_view_genres动作 double prob Double.parseDouble(parts[1]); probMap.put(featureKey, prob); } } } } } } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\t); if (fields.length 5) { return; } String userId fields[0].trim(); String age fields[2].trim(); String genres fields[3].trim(); String rating fields[4].trim(); // 构造特征列表和训练阶段保持一致 ListString features new ArrayList(); features.add(feature_age_age discretizeAge(age)); for (String genre : genres.split(,)) { features.add(feature_view_genres genre.trim()); } features.add(feature_rating discretizeRating(rating)); double maleScore Math.log(priorMale); double femaleScore Math.log(priorFemale); for (String feature : features) { Double pMale probMap.get(M_ feature); Double pFemale probMap.get(F_ feature); if (pMale ! null) { maleScore Math.log(pMale); } if (pFemale ! null) { femaleScore Math.log(pFemale); } } String predict maleScore femaleScore ? M : F; outKey.set(userId); outValue.set(predict); context.write(outKey, outValue); } }这里用的就是上面提到的取对数累加避免概率连乘后下溢出。虽然课程设计数据量不大可能不取对数也没问题但养成好习惯总没错。4.4 Driver里的参数与HDFS路径管理Driver里需要依次提交三个Job第一个Job跑完以后读取输出目录下的part-r-00000文件解析出男女总数再把这些数字放进第二个Job的Configuration同时把Job2的输出目录加入第三个Job的DistributedCache。Driver的核心伪代码如下public class GenderPredictDriver extends Configured implements Tool { Override public int run(String[] args) throws Exception { String trainPath args[0]; String testPath args[1]; String outputBase args[2]; Path priorOutput new Path(outputBase /prior); Path conditionOutput new Path(outputBase /condition); Path predictOutput new Path(outputBase /predict); // Job1: 性别先验 Job job1 Job.getInstance(getConf(), gender prior count); job1.setJarByClass(getClass()); job1.setMapperClass(PriorCountMapper.class); job1.setReducerClass(PriorCountReducer.class); job1.setOutputKeyClass(Text.class); job1.setOutputValueClass(IntWritable.class); TextInputFormat.addInputPath(job1, new Path(trainPath)); TextOutputFormat.setOutputPath(job1, priorOutput); job1.waitForCompletion(true); // 读取先验输出 FileSystem fs FileSystem.get(getConf()); long maleCount 0, femaleCount 0; FileStatus[] statuses fs.listStatus(priorOutput); for (FileStatus status : statuses) { if (status.getPath().getName().startsWith(part)) { try (BufferedReader reader new BufferedReader( new InputStreamReader(fs.open(status.getPath())))) { String line; while ((line reader.readLine()) ! null) { String[] parts line.split(\t); if (parts.length 2) { if (M.equals(parts[0])) { maleCount Long.parseLong(parts[1]); } else if (F.equals(parts[0])) { femaleCount Long.parseLong(parts[1]); } } } } } } // Job2: 条件概率 Configuration conf2 getConf(); conf2.set(male.count, String.valueOf(maleCount)); conf2.set(female.count, String.valueOf(femaleCount)); Job job2 Job.getInstance(conf2, condition prob); job2.setJarByClass(getClass()); job2.setMapperClass(ConditionProbMapper.class); job2.setReducerClass(ConditionProbReducer.class); job2.setOutputKeyClass(Text.class); job2.setOutputValueClass(DoubleWritable.class); TextInputFormat.addInputPath(job2, new Path(trainPath)); TextOutputFormat.setOutputPath(job2, conditionOutput); job2.waitForCompletion(true); // Job3: 预测 Job job3 Job.getInstance(getConf(), gender predict); job3.setJarByClass(getClass()); job3.setMapperClass(PredictMapper.class); job3.setNumReduceTasks(0); job3.setOutputKeyClass(Text.class); job3.setOutputValueClass(Text.class); TextInputFormat.addInputPath(job3, new Path(testPath)); TextOutputFormat.setOutputPath(job3, predictOutput); job3.addCacheFile(new URI(conditionOutput /part-r-00000#model)); return job3.waitForCompletion(true) ? 0 : 1; } public static void main(String[] args) throws Exception { int exitCode ToolRunner.run(new GenderPredictDriver(), args); System.exit(exitCode); } }Driver里有几个容易踩的细节。一是setNumReduceTasks(0)让预测阶段只跑Map避免Reducer带来的额外排序开销也避免输出文件里出现part-r-。二是输出目录必须事先不存在否则作业会直接报错退出所以提交前最好检查并删除旧目录。三是DistributedCache的文件名要固定比如#model读取时直接用model就能从本地当前目录拿到文件。5. 从伪分布式到真集群部署过程中踩过的坑5.1 数据倾斜某类影视题材“刷榜”项目跑测试数据的时候一切正常但一放到完整训练集上第二个作业就慢得离谱。原因是数据倾斜一部分题材在样本里出现频率极高比如“动作”这种题材在男性用户的浏览记录里几乎人人都有导致对应特征key“M_feature_view_genres动作”的计数远大于其他key。所有这个key的数据都集中到一个Reducer上别的Reducer早跑完了就它还在慢慢汇总。这个问题的根因是MapReduce默认的HashPartitioner按key散列分区key里包含了特征值热门特征值自然会把大量数据送到同一个分区。解决办法有几种第一种是在组合key中加入一个随机前缀把同一个特征的计数分散到多个Reducer。比如把key变成“随机数_特征名”然后在Reducer端只做局部汇总最后再用一个单独的MapReduce作业按特征名聚合一次。这样做会增加一个作业但能显著缓解倾斜。第二种是使用Combiner。Combiner能在Map端先做一次局部求和减少shuffle的数据量。对于“M_feature_view_genres动作”这种高频key虽然最后仍然集中到一个Reducer但每个Map传给它的数据量已经从几万条变成了几条聚合结果。我在项目里直接复用了ConditionProbReducer作为Combiner代码一行不用改效果立竿见影。第三种是人工识别高频特征把它单独分区或走特殊计算路径。这个适合真实生产环境对课程设计来说有点过了。5.2 中文编码和分隔符问题训练数据里有很多中文题材比如“动作、爱情、喜剧”。Hadoop的TextInputFormat默认按UTF-8处理文本但如果你的原始文件是从Windows导出的CSV很可能编码是GBK。Map端读进来后中文全是乱码统计出来的特征名也乱七八糟。解决办法是在预处理阶段统一转码用iconv -f GBK -t UTF-8把原始文件转成UTF-8再上传到HDFS。千万不要指望在Mapper里给每一行重新解码因为TextInputFormat已经按字节流切分文件如果文件是GBK切分点可能正好落在某个汉字中间导致汉字被截断成乱码后面怎么解码都回不来。分隔符问题更隐蔽。原始数据用逗号分隔字段但题材列表内部也用逗号分隔比如“动作,科幻,犯罪”这样字段就会错位。我建议在预处理时统一把字段分隔符改成\t题材列表内部的逗号保留或者改成竖线|这样Mapper按\t拆分字段再按逗号拆题材列表逻辑清晰也不容易错。5.3 小文件过多导致的Map数暴涨这个项目在预处理阶段会产生很多小文件尤其是用Python脚本按用户拆分特征的时候一个用户一个文件最后整个目录下可能有几万个小文件。MR作业扫描输入路径时每个小文件都会占一个InputSplit最终生成一个MapTask几万个小文件意味着几万个MapTask光是启动Task的JVM开销就够集群喝一壶的。解决的思路是预处理后的文件不要存太多碎片尽量做到一个目录下只有几个大文件。可以用hdfs dfs -getmerge先把本地小文件合并成一个大文件再上传或者用hdfs dfs -put上传前在本地先cat *.part merged.tsv。如果数据已经在HDFS上也可以借助Hive的INSERT OVERWRITE TABLE ... SELECT把碎片表合并成压缩的大文件。5.4 内存溢出与Reducer个数设置训练作业跑久了会偶发OOM绝大多数情况发生在Reducer端。因为Reducer在内存里维护着所有key对应的value迭代器虽然MapReduce会定期把数据刷到磁盘但如果单个key的value太多或者Reducer数量太少内存压力就会变大。解决办法是先调mapreduce.reduce.memory.mb我一般从默认的1024M调到2048M或者3072M。前提是集群的每个NodeManager有足够的内存配额。另外Reducer数量要结合实际数据量设置。这个项目数据量不大我设置了10个Reducer就够用了。数据量大的时候也尽量不要让Reducer数量超过集群可用vcore数的一半否则很多Reducer排队等待整体耗时反而更差。6. 预测效果评估准确率高了不代表模型好6.1 混淆矩阵和AUC怎么算课程设计里最常见的误区是只看准确率。假设测试集里60%是男性用户那么一个“全部预测为男”的傻瓜模型准确率也有60%。所以评估时要把混淆矩阵列出来至少看一下四类数据真正男、假正男、真正女、假正女。在Hadoop项目里可以直接在预测阶段用一个统计作业输出混淆矩阵。Mapper读预测结果和真实标签输出“实际性别_预测性别”为key计数1Reducer求和。这样一个简单的作业就能得到完整的混淆矩阵F_F 1288 F_M 342 M_F 251 M_M 1897从混淆矩阵可以算出精确率和召回率。如果业务上更关心“把男性用户找出来”的覆盖率就要看男性召回率如果更关心“预测结果是男性的用户中有多少真是男性”就要看精确率。只报告一个数字容易被数据分布骗过去。AUC的计算在MapReduce里稍微麻烦一点需要按预测概率从大到小排序然后按序统计正负样本的秩。如果只是为了课程设计可以把预测结果导出到本地用Python的sklearn直接算AUC几万条结果很快。重点是要理解AUC代表的是“随机抽一个男用户和一个女用户模型把男用户排到前面的概率”而不是看懂一个函数调用。6.2 分层评估新用户和老用户分开看我在实际评估时发现如果把所有用户混在一起看准确率还不错但按用户行为丰富程度拆分后差异非常大。只用过一两个特征的新用户预测准确率可能只有55%和瞎猜差不多而行为特征丰富的用户准确率能到75%以上。原因很简单朴素贝叶斯依赖特征来计算后验概率特征太少就退化成了先验概率。所以在项目的最终报告里我建议把测试集按特征数量分层分别统计每个层级的准确率这样能清楚展示模型在什么条件下可用什么条件下不可用。这个分层评估的思路反过来也能指导产品方案对特征量极少且预测置信度不高的用户不要硬给性别标签宁可在前端标注为“未知”或者把预测结果用作候选池而不是最终判断。6.3 阈值调整与线上使用方式朴素贝叶斯输出的是两个后验概率的比较不是0到1之间的分数。默认取P(M)和P(F)哪个大就用哪个。但在实际使用中如果误判某个性别的代价更大可以调整阈值。比如业务方更怕把女生预测成男生那就只有当P(F)比P(M)高30%以上时才输出F否则输出“不确定”。MapReduce项目里给预测逻辑加阈值很简单在PredictMapper里比较两个分数时添加一个偏置项即可。比如double thresholdBias Double.parseDouble(conf.get(bias, 0)); String predict (maleScore thresholdBias) femaleScore ? M : F;这样不用改模型文件只需在提交作业时传一个参数就能反复测试不同的阈值效果。这个方法在生产环境中很实用因为模型上线后不一定只有“男/女”两个答案有时候“不判断”反而是最优策略。7. 让这套代码更有用的扩展思路7.1 Spark化从MR换到RDD和DataFrame如果以后数据量继续增长或者需要支持迭代式算法MapReduce的短板就很明显每个作业都要落盘中间结果重复读写HDFS。同样的逻辑用Spark写可能只有原来三分之一的代码量而且Spark的DataFrame API对特征工程更友好。但要注意从MapReduce换到Spark不等于简单把Mapper和Reducer替换成map和reduceByKey。Spark里更自然的做法是把样本封装成DataFrame用groupBy和agg完成统计再用VectorAssembler构造特征向量。虽然逻辑上仍然是朴素贝叶斯但工程组织方式完全不同。7.2 特征升级引入协同过滤输出和文本Embedding这个项目只用了用户自己的行为特征没有引入用户之间的相似性。想提高预测准确率下一步可以加入协同过滤的结果。比如先用ItemCF算出用户和某些典型男性向电影的匹配程度把一个连续值作为新特征。再比如对用户浏览过的电影标题做Embedding把平均向量作为特征输入。这些操作在MapReduce里写起来会比较吃力建议换到Spark或者用其他框架配合。7.3 增量更新定期重算还是在线学习电影网站的偏好是会变的今年流行的题材和去年可能差异很大。训练出来的模型不能一劳永逸。最简单的做法是定期离线重算比如每周跑一次完整流程重新生成模型文件。这正好是MapReduce擅长的场景全量扫描稳定可控。如果想更激进一点把每天的日志增量累加到统计计数里就需要维护模型文件的版本和状态复杂度明显上升不建议在课程设计里展开。按照我个人的实操体会这个项目的难度不在MapReduce代码本身而在数据准备和特征设计。很多同学拿着样例数据上来就写Mapper写到一半发现输入格式不对回头又去改预处理。建议第一步先用Python或Hive把数据结构彻底理清固定好TSV格式和特征枚举范围再动Hadoop相关代码。模型文件、预测输出、评估指标这些模块之间通过路径和配置参数解耦跑起来以后排查问题会轻松很多。本文还有配套的精品资源点击获取