ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

用MapReduce实现朴素贝叶斯文本分类:完整Hadoop源码解析

用MapReduce实现朴素贝叶斯文本分类:完整Hadoop源码解析 简介基于Hadoop的朴素贝叶斯文本分类器项目完整实现MapReduce训练、分类与评估流程适合Hadoop课程设计、毕业设计或机器学习入门者参考学习。数据集选取NBCorpus中CHINA与CANA两类共518篇英文文本按70%/30%比例划分训练集与测试集项目覆盖贝叶斯模型训练、测试文档分类以及Precision、Recall、F1值计算可直接用于理解朴素贝叶斯在分布式环境下的落地方式。压缩包共552个文件以txt文本数据、java源文件、png图片、docx/pdf说明文档为主总大小约3.75MB另含README与项目配置信息目录结构清晰便于查阅。已有230人学习下载。资源内代码经运行验证适合在此基础上扩展或二次开发是完成课设、毕设及初期项目演示的实用参考。1. 从课程设计到可跑通的朴素贝叶斯这份 Hadoop 源码到底怎么用如果你正在做 Hadoop 课程设计或毕业设计多半会被“用 MapReduce 实现一个分类器”这类要求卡住。理论课上朴素贝叶斯半小时就能讲完但真要放到 HDFS 上跑训练、出模型、做评估细节远比想象多。这份源码的特别之处在于它不是一个 Demo而是一个完整可运行的 Hadoop 朴素贝叶斯文本分类器训练阶段拆成多个 MapReduce Job预测和评估也都有对应实现数据用的是 NBCorpus 的 CHINA 与 CANA 两类新闻文本合计 518 篇文档按 70/30 划分训练集与测试集。适合两类人一是计算机相关专业做课设、毕设的学生可以直接跑通作为项目基线二是想搞清楚“朴素贝叶斯在分布式环境下到底怎么拆成 MapReduce 步骤”的工程师源码里的 Job 划分方式本身就值得读一遍。本文会先讲清楚这个分类器的数据流和六个核心类各自负责什么再逐步拆解每个 MapReduce Job 的输入输出、代码关键点最后给出从打包到评估的完整命令和常见坑。2. 贝叶斯训练拆成 MapReduce统计什么、怎么统计2.1 朴素贝叶斯在文本分类上的数学形式朴素贝叶斯分类器做的事情很简单给定一篇文档 (d)计算它属于类别 (c) 的后验概率 (P(c|d))取概率最大的类别作为预测结果。根据贝叶斯定理[ P(c|d) \propto P(c) \times \prod_{i1}^{n} P(w_i|c) ]其中 (P(c)) 是类别先验概率(P(w_i|c)) 是单词 (w_i) 在类别 (c) 下的条件概率。文本场景下文档被表示为词袋模型也就是忽略词序只统计词频。为了处理训练集中未出现的单词条件概率一般用拉普拉斯平滑加一平滑计算[ P(w_i|c) \frac{count(w_i, c) 1}{count(c) |V|} ]这里 (count(c)) 是类别 (c) 下所有单词的总数(|V|) 是整个训练集的词汇表大小。对于二分类来说只需要比较两个类别的后验概率所以训练阶段的任务本质上是统计四类数值文档总数、每个类别的文档数、每个类别下每个单词的出现次数、每个类别下所有单词的总次数。2.2 训练阶段为什么需要拆成多个 MapReduce Job一个朴素贝叶斯训练过程看起来只需要一次遍历统计词频但在 Hadoop 上实现时受限于 reduce 阶段按 key 分组的机制直接在同一个 Job 里统计不同粒度的聚合值会让代码很别扭。比如“每个类别文档总数”的 key 是类别名“每个单词在每个类别下次数”的 key 是“类别单词”两套统计维度混在一起就需要在 reduce 里做类型判断。这份源码的做法是把不同维度的统计拆成独立的 Job每个 Job 只负责一类聚合逻辑清晰且易于验证中间结果。具体来说训练阶段涉及四个 JobGetDocCountFromDocTypeJob统计每个类别下的文档数量输出类别名到文档数的映射GetSingleWordCountFromDocTypeJob统计每个类别下每个单词的出现次数输出“类别单词”到词频的映射GetTotalWordCountFromDocTypeJob统计每个类别下所有单词的总次数输出类别名到总词数的映射InitSequenceFileJob把训练文本转换为 Hadoop SequenceFile 格式方便后续 Job 读取。前三个 Job 是典型的“一个维度一个 Job”的思路最后一个 Job 是数据预处理环节。理解了这个划分方式再看每个类的实现就不会迷路。2.3 InitSequenceFileJob文本数据怎么变成 Hadoop 可读格式在 MapReduce 中直接读取原始文本文件当然可以但默认的TextInputFormat会把每行作为一个 value而这里的训练单元是“一篇文档”——文档可能跨多行行数还不固定。直接按行读取会破坏文档完整性所以需要先把文档转换成 key-value 形式key 是文档名或文档IDvalue 是文档全文内容。SequenceFile 是 Hadoop 原生的二进制 key-value 存储格式适合作为中间结果在多轮 MapReduce 之间传递。InitSequenceFileJob的核心逻辑是输入路径指向原始数据集目录输出为一个 SequenceFile。Mapper 用FileSplit获取当前文件路径作为 key逐行读取并拼接为完整文档内容作为 value。关键代码如下public class InitSequenceFileMapper extends MapperLongWritable, Text, Text, Text { private Text docId new Text(); private Text docContent new Text(); private StringBuilder content new StringBuilder(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { FileSplit fileSplit (FileSplit) context.getInputSplit(); String fileName fileSplit.getPath().getName(); docId.set(fileName); content.append(value.toString()).append(\n); docContent.set(content.toString()); context.write(docId, docContent); } }这段代码的逻辑说明context.getInputSplit()获取当前输入分片getPath().getName()拿到正在处理的文件名作为文档IDmap方法被每一行调用通过 StringBuilder 累积内容。需要注意的一点是按行累积到Text对象时value 会被反复覆盖当文本文件很大时所有行会在同一个map方法调用链中累积吗实际上不是——每个输入分片对应一个 Mapper 实例但map方法对分片内的每一行调用一次StringBuilder是 Mapper 实例的成员变量所以一个文件的多行会持续累积到同一个StringBuilder中。但这也意味着一个 Mapper 实例处理完一个文件后状态不会自动清空如果同一个 Mapper 被分配了多个文件内容会混在一起。解决办法是在setup阶段或写入后重置StringBuilder源码中通常直接在 write 之后调用content.setLength(0)清空。参数说明输入路径是训练数据的根目录Hadoop 会递归读取其下所有文件输出路径是 SequenceFile 的存放目录。TextOutputFormat默认的行分隔输出在这里不适用必须设置为SequenceFileOutputFormat否则后续 Job 读取时格式不匹配。3. 训练模型生成三类统计 Job 的输入输出与参数设定3.1 GetDocCountFromDocTypeJob类别文档数统计这个 Job 是最简单的一个。输入是InitSequenceFileJob产出的 SequenceFilekey 是文档名value 是文档内容。Mapper 的输入 key-value 已经是一篇完整文档所以无需在 Mapper 里解析内容只需要提取文档所属类别。类别信息从哪来数据集的目录结构是NBCorpus/Country/CHINA/xxx.txt和NBCorpus/Country/CANA/xxx.txt即上级目录名就是类别名。在 Mapper 中通过FileSplit的路径拿到父目录名作为输出的 keyvalue 写一个固定值 1public class DocCountMapper extends MapperText, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text category new Text(); Override protected void map(Text key, Text value, Context context) throws IOException, InterruptedException { String filePath ((FileSplit) context.getInputSplit()).getPath().getParent().getName(); category.set(filePath); context.write(category, one); } }Reducer 端直接把相同 key 的 value 累加输出类别名与文档总数。这里的逻辑说明因为输入已经是 SequenceFilegetInputSplit()拿到的路径还是原始文件路径吗在 SequenceFile 作为输入时FileSplit.getPath()返回的是 SequenceFile 本身而非原始数据路径。这会导致getParent().getName()拿到的是 SequenceFile 目录的父目录而不是CHINA或CANA。这个细节是实操中最容易出错的地方。解决办法有两种第一种是在InitSequenceFileJob输出 key 时就把类别编码进去比如 key 设计为CHINA/xxx.txtMaper 解析 key 中分隔符前的部分第二种是DocCountJob的输入不走 SequenceFile而是直接读原始文本目录用TextInputFormat默认按行读取通过FileSplit拿类别——此时不需要文档完整性只需要数量。源码中实际用的是第二种。所以这个 Job 的输入路径应指向原始数据集目录而不是 SequenceFile 目录。参数建议mapreduce.job.reduces设为 2有两个类别或保持默认让框架自动决定。输出格式为TextOutputFormat每行是类别名 \t 文档数。3.2 GetSingleWordCountFromDocTypeJob类别下单词词频统计这是训练阶段最核心的 Job。输入是完整文档Mapper 要做分词和类别提取两件事。分词部分需要引入一个分词器——这里没有用 IK Analyzer 或 HanLP而是简单的空格与标点切分因为 NBCorpus 的英文文本大小写混杂直接 split 会导致CHINA和china被当成两个词。所以在分词前需要统一转小写再用正则切除非字母字符public class WordCountMapper extends MapperText, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text wordWithCategory new Text(); Override protected void map(Text key, Text value, Context context) throws IOException, InterruptedException { String category ((FileSplit) context.getInputSplit()).getPath().getParent().getName(); String[] words value.toString().toLowerCase() .replaceAll([^a-z\\s], ).split(\\s); for (String word : words) { if (word.isEmpty()) continue; wordWithCategory.set(category \t word); context.write(wordWithCategory, one); } } }输出 key 是类别 制表符 单词value 是 1。Reducer 端按这个组合 key 聚合词频。但这里有一个问题Reducer 的输入 key 是复合的输出时需要拆开以适配模型格式——训练模型需要的是“类别、单词、词频”三列。可以在 Reducer 里用key.toString().split(\t)拆分再写回三个字段的文本行。参数上需要注意的是这里的mapreduce.map.output.compress可以设为 true默认 false因为词频统计的中间数据量较大压缩能减少 shuffle 阶段的网络传输。压缩格式推荐org.apache.hadoop.io.compress.SnappyCodecmapreduce.map.output.compresstrue mapreduce.map.output.compress.codecorg.apache.hadoop.io.compress.SnappyCodec3.3 GetTotalWordCountFromDocTypeJob类别总词数的两种实现路线这个 Job 相对取巧。如果已经跑完了GetSingleWordCountFromDocTypeJob那么每个类别下的总词数就等于该类别所有单词词频之和——可以直接对词频 Job 的输出做二次聚合而不是重新扫描原始数据。这也是为什么源码里单独有一个GetTotalWordCountFromDocTypeJob它的输入是GetSingleWordCountFromDocTypeJob的输出Mapper 只需取 key 的第一个字段类别名value 取词频Reducer 做累加。这样做的好处是节省一次全量数据扫描代价是多一个 MR Job 的调度开销。在数据量只有几百篇文档的场景下后者的开销完全可以接受所以这个设计是合理的。public class TotalWordCountMapper extends MapperLongWritable, Text, Text, IntWritable { Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\t); context.write(new Text(fields[0]), new IntWritable(Integer.parseInt(fields[2]))); } }逻辑说明fields[0]是类别名fields[1]是单词fields[2]是词频。因为这类 Job 的输入是纯文本输出格式直接跳过 SequenceFile 的转换。Reducer 端只需 sum 即可。这个 Job 的输出只有两行两个类别如下所示CHINA 86453 CANA 87911数值是示意。有了这三个 Job 的输出训练模型的全部参数就齐了类别文档数、每个类别的词汇分布、每个类别的总词数。3.4 模型文件怎么组织训练完成后模型不是一个单独的文件而是三个目录下的多个 part 文件。实际应用时需要写一个小工具把这些 part 文件合并加载到内存 HashMap 中结构大致如下MapString, Double categoryDocCount; // 类别 - 文档数 MapString, MapString, Long wordCount; // 类别 - (单词 - 词频) MapString, Long categoryTotalWordCount; // 类别 - 总词数在测试分类时对每个测试文档做同样的分词处理遍历其单词集合逐词计算两个类别的条件概率对数之和加上类别先验的对数取较大者。4. 测试与评估单机 Java 走通流程再上 MapReduce4.1 测试阶段的两种预设源码允许测试阶段是单机 Java 程序也可以是 MapReduce 程序。单机方式的优势是调试方便可以直接断点看概率计算过程MapReduce 方式的优势是测试集很大时能并行处理。这里先说单机方式因为它是理解预测逻辑最快的手段。加载模型后预测一个文档类别的核心代码如下public String predict(String docText, SetString vocabulary) { String[] words docText.toLowerCase() .replaceAll([^a-z\\s], ).split(\\s); double bestLogProb Double.NEGATIVE_INFINITY; String bestCategory null; for (String category : categoryDocCount.keySet()) { double logProb Math.log(categoryDocCount.get(category) / totalDocCount); long totalWords categoryTotalWordCount.get(category); MapString, Long cateWordCount wordCount.get(category); for (String word : words) { long count cateWordCount.getOrDefault(word, 0L); logProb Math.log((count 1.0) / (totalWords vocabulary.size())); } if (logProb bestLogProb) { bestLogProb logProb; bestCategory category; } } return bestCategory; }这段代码的逻辑说明外层循环遍历两个类别内层循环遍历测试文档的每个单词。vocabulary是整个训练集的词汇表大小由所有类别的单词集合取并集得到。在计算条件概率时分子加一拉普拉斯平滑分母加词汇表大小这两处参数是关键如果分母只加类别自己的词汇量而不是全局词汇量会导致两个类别的概率不具备可比性——因为各自的分母不同相当于两套不同的尺度。这是朴素贝叶斯实现中最常见的错误。数值稳定性方面直接用Math.log连乘会造成下溢因为条件概率通常远小于 1几百个词乘下来会变成 0.0。所以必须使用对数加法log(a*b) log(a) log(b)。上面代码里用的是对数概率累加这正是它不直接乘概率的原因。4.2 Evaluation 类Precision、Recall、F1 怎么算评估阶段Evaluation类的职责是把测试文档的真实类别与预测类别做对比产出混淆矩阵。基于混淆矩阵计算Precision TP / (TP FP)Recall TP / (TP FN)F1 2 * Precision * Recall / (Precision Recall)对于二分类需要选定一个正类比如 CHINA。实现上可以维护一个MapString, MapString, Integer confusionMatrix外层 key 是真实类别内层 key 是预测类别直接统计两两组合的数量。计算指标时TP confusionMatrix.get(CHINA).get(CHINA)FP 其他类被预测为 CHINA 的数量之和FN CHINA 被预测为其他类的数量之和。输出格式建议做成三行表格样式方便写进实验报告。4.3 GetNaiveBayesResultJobMapReduce 版预测实现要点MapReduce 版预测同样遵循“一个测试文档一条记录”的原则。需要提前把测试文档也转成 key-value 格式或者直接用一个自定义 InputFormat 按文件读取。Mapper 端的核心逻辑与单机版的predict方法一样但模型参数的传递方式不同——不能把整个 HashMap 塞进每个 Mapper常见做法是放到 DistributedCache或 Hadoop 3 的共享缓存中。hadoop jar naive-bayes.jar GetNaiveBayesResultJob \ -files hdfs:///model/categoryDocCount.txt,hdfs:///model/wordCount.txt,hdfs:///model/totalWordCount.txt \ /input/test_seq /output/test_result参数说明-files后的每个文件会被分发到每个 Mapper 节点的本地工作目录Maper 在setup方法中读取并加载到内存。这里的categoryDocCount.txt是前面三个训练输出合并后的模型文件。测试 Job 的输出为每行一个文档的分类结果包含文档名、真实类别从目录路径提取、预测类别三个字段。5. 模型应用与排错从合并模型文件到处理低准确率5.1 训练输出怎么快速合并成单文件模型实训时三个训练 Job 的输出是分散的 part 文件直接用 Hadoop 命令一键合并hdfs dfs -getmerge /output/doc_count /local/model/doc_count.txt hdfs dfs -getmerge /output/word_count /local/model/word_count.txt hdfs dfs -getmerge /output/total_word_count /local/model/total_word_count.txtgetmerge的参数说明第一个路径是 HDFS 上的目录第二个路径是本地目标文件。它会按 part 文件名的字典序拼接内容不是按行内容排序。如果后续加载依赖有序数据需要在合并前用sort处理或者干脆在 MapReduce 输出时把分区数设为 1setNumReduceTasks(1)确保只有一个 part 文件从根上避免乱序问题。这三个文件需要进一步合并成一个模型文件格式建议采用 JSON 或自定义文本协议这里给出一个简单的文本模型文件格式CATEGORY_DOC_COUNT CHINA 255 CATEGORY_DOC_COUNT CANA 263 WORD_COUNT CHINA economy 125 WORD_COUNT CHINA trade 98 TOTAL_WORD_COUNT CHINA 86453 TOTAL_WORD_COUNT CANA 87911加载时逐行解析用前缀区分记录类型。这个文本格式有一个好处直接hadoop fs -cat就能查看内容排查数据是否正确不需要写额外工具。5.2 分类结果准确率低时的排查路径跑通流程之后最容易遇到的问题是测试集准确率远低于预期。按照下面的顺序逐个排查基本能定位到根因。第一检查训练阶段三个统计 Job 的输出是否合理。例如GetTotalWordCountFromDocTypeJob输出两个类别的总词数相差几倍说明数据集本身不均衡需要看是否需要做类别权重调整。第二检查分词逻辑是否丢了关键词。replaceAll([^a-z\\s], )会去掉所有非字母字符数字、连字符、缩写都会被切断像U.S.会变成u s两个词。如果数据集里领域术语多建议只去标点不动字母数字。第三检查测试集是否和训练集重叠。用getmerge合并训练集测试集文档名列表用comm -12对比重复项如果重叠超过 1%评估出来的 F1 会很虚高但实际没有参考价值。第四检查拉普拉斯平滑参数。默认加一是常规做法但如果词汇表特别大几万词加一平滑会严重压低所有条件概率这时可以把这个参数调到 0.5 或更小观察 F1 的变化。5.3 一个容易被忽略的参数vocabulary 的并集计算时机在 4.1 节的单机预测代码里vocabulary是作为参数传入的。它的计算方式应该是把 CHINA 类单词表和 CANA 类单词表做并集得到全局词表大小。如果在训练时分别记录了两个类别的词汇量想当然地用wordCount.get(CHINA).size() wordCount.get(CANA).size()会重复计算两个类别共有的词导致分母虚大。正确做法是构造一个HashSet把两个类别的 keySet 全部 add 进去再size()。这个 Bug 在代码里很难肉眼看出来因为vocabulary本身只影响条件概率的分母在几百个词的文档下虚大的分母会把所有类别的分数同时压低最终的类别判断反而可能不被改变——所以它只影响概率绝对值不影响排名。但如果你后续要做阈值截断或置信度输出这个误差会直接影响结果。5.4 处理数据倾斜当某个类别的文档特别多如果数据集不是均衡的比如 CHINA 有 5000 篇而 CANA 只有 200 篇那么先验概率会偏向 CHINA导致小类别的召回率奇低。在 MapReduce 训练不变的前提下可以在模型加载阶段修改先验概率的计算方式用对数先验加一个权重因子double prior Math.log(categoryDocCount.get(category) / totalDocCount); prior * 0.8; // 削弱先验影响这个 0.8 是经验值调参时需要观察 Precision/Recall 曲线变化。注意这里调整的是先验而不是条件概率因为类别不平衡主要影响先验项条件概率是在类别内统计的天然与类别文档数无关。5.5 伪分布式与集群环境下的内存参数设置最后提一个运维层面常见的问题GetSingleWordCountFromDocTypeJob在跑较大数据集时容易 OOM。原因是 Mapper 端的分词结果没有做 combine几万篇文档的中间结果全部涌向 shuffle。解决方法是加一个 Combiner复用 Reducer 类job.setCombinerClass(WordCountReducer.class);如果已经加了 Combiner 还 OOM检查 Mapper 端是否设置了mapreduce.task.io.sort.mb这个参数控制 map 端排序缓冲区的最大值默认 100MB。数据量大时把它调到 200MB可以显著减少 spill 次数。注意这个参数在yarn-site.xml中容易被覆盖启动 Job 时通过-D传入最稳妥hadoop jar naive-bayes.jar GetSingleWordCountFromDocTypeJob \ -D mapreduce.task.io.sort.mb200 \ /input/train_seq /output/word_count-D参数说明它只对当前 Job 生效优先级高于配置文件中的同名设置。伪分布式环境下这个参数效果不明显真正有价值的是在 10 个节点以上的集群上shuffle 的数据量会决定 Job 的执行时间。6. 验证模型效果重构分类主程序并输出评估报告整个项目跑通后我一般会额外写一个小工具把训练和评估串联起来方便反复实验。这个工具不要求是 MapReduce用伪分布式集群上的hadoop jar加载模型后在本地直接做预测和指标计算反而更快。核心是一个双重循环外层遍历测试集目录下的待分类文档内层调用 4.1 节的predict方法。当某篇文章的分类结果显示在confusionMatrix中把结果写入评估输出文件。一个值得试验的边界是把测试集文档去重后多次输入predict输出结果应该完全一致——概率计算是确定性的没有随机因素。这可以用来验证模型加载是否正确。如果发现两次结果不同检查代码中是否有使用HashMap的迭代顺序影响输出——朴素贝叶斯的类别比较与迭代顺序无关但如果你在输出格式化时依赖了entrySet()的顺序可能出现不稳定显示。评估报告的最终输出建议包含以下字段测试文档总数、每类的 TP/FP/FN、宏平均与微平均的 Precision/Recall/F1。这里有一个实用技巧把混淆矩阵输出为 CSV 格式导入 Excel 后直接做透视表比在控制台看日志直观很多。输出代码用String.format拼接即可注意浮点数保留 4 位小数避免报告里出现一长串小数位。至此从InitSequenceFileJob到Evaluation整个朴素贝叶斯分类器的训练、预测、评估链路就完全打通了。如果你需要把这份源码作为课程设计提交建议把三个训练 Job 的中间输出数量如CHINA类单词种类数、总词数贴到实验报告里作为训练正确性的佐证。本文还有配套的精品资源点击获取
返回列表