Hadoop朴素贝叶斯文本分类器:MapReduce实现与源码解析
2026/9/23 21:39:32 网站建设 项目流程

简介:基于Hadoop的朴素贝叶斯文本分类器项目,采用MapReduce编程模型完成分类模型的训练与预测,适合作为Hadoop课程设计、毕业设计或算法入门实践。实验数据取自NBCorpus的CHINA和CANA两个类别,共五百余篇英文文本,按百分之七十与百分之三十划分为训练集与测试集。项目完整实现了序列文件初始化、类别文档计数、单类词频统计、总词频统计、模型生成与分类评估等多个MapReduce任务,并计算准确率、召回率和F1值。压缩包共552个文件,包括9个Java源码、518个txt格式的数据文本、Hadoop工程报告docx、说明文档md以及相关图表png,整体大小仅3.75MB,文件组织清晰,便于按处理流程对照学习。已有230人学习,代码经过验证可正常运行,附带的工程报告详细阐述了朴素贝叶斯原理、各阶段MapReduce设计思路、运行环境配置和实验结果分析,能帮助读者快速掌握分布式环境下文本分类的完整实现过程。

1. Hadoop 课程设计里的朴素贝叶斯:这份源代码到底帮你把活干到了哪一步

手头这份基于 Hadoop 的朴素贝叶斯文本分类器源代码,是一个完整的 Hadoop 课程设计交付物:6 个 Java 类,一个 IntelliJ 工程配置,外加一份 Word 报告。项目用 MapReduce 把朴素贝叶斯模型的训练过程拆成了三个统计作业和一个序列文件转换作业,再用一个独立的评估类计算 Precision、Recall 和 F1,训练和评估的链路是闭环的。

数据集用的是NBCorpus\Country目录下的 CHINA(255 篇)和 CANA(263 篇)两类文本,按 70%/30% 切分训练集和测试集。也就是说,你下载到的不是一堆零散算法片段,而是一套「数据预处理 → 模型训练 → 分类预测 → 指标评估」全流程代码。适合正在做 Hadoop 课程设计、毕业设计,或者想搞明白「朴素贝叶斯训练过程怎么拆成 MapReduce 作业」的读者直接复现。

2. 训练阶段的四个 Job 分工:从文本到贝叶斯模型的落盘过程

2.1 为什么项目先用 InitSequenceFileJob 把文本转成 SequenceFile

朴素贝叶斯训练的第一步,是把原始文本变成后续 Job 能稳定读取的输入格式。项目里的InitSequenceFileJob干的就是这件事:把NBCorpus\Country下的文本文件读进来,输出为一个 key-value 形式的 SequenceFile,key 是文档 ID,value 是文档内容。

直接用 TextInputFormat 读原始文本不行吗?可以,但有几个现实问题。原始文本的文件名在 MapReduce 里要通过FileSplit才能拿到,而且每个小文件会各自占用一个 InputSplit,后续 Job 如果还要反复按文档类别做分组,路径解析就变得非常啰嗦。转成 SequenceFile 之后,key 里可以直接带上类别和文档编号,后续三个统计 Job 只需要从 value 里取词,不用再关心文件来自哪个目录。

提示:SequenceFile 是 Hadoop 自带的二进制序列化格式,适合作为多个 MapReduce Job 之间的中间数据存储格式,比反复解析文本文件稳定得多。

2.2 三个统计 Job:文档数、单词频、类别总词数

朴素贝叶斯训练需要三类统计量,项目里对应三个 Job:

  • GetDocCountFromDocTypeJob:统计每个类别下有多少篇文档,用来计算先验概率 P(类别);
  • GetSingleWordCountFromDocTypeJob:统计每个类别中每个单词出现多少次,用来计算条件概率 P(词|类别);
  • GetTotalWordCountFromDocTypeJob:统计每个类别下所有单词的总数,作为条件概率的分母。

GetSingleWordCountFromDocTypeJob为例,它的核心逻辑就是一个典型的 word count 变体,只是 key 从「单词」换成了「类别 + 单词」。示意代码如下:

public class WordCountByClassMapper extends Mapper<Text, Text, Text, IntWritable> { private final Text outKey = new Text(); private final IntWritable one = new IntWritable(1); @Override protected void map(Text docId, Text content, Context context) throws IOException, InterruptedException { // 假设 docId 格式为 CHINA_001,第一段是类别 String[] parts = docId.toString().split("_"); String category = parts[0]; // 按空白字符分词,实际项目中可以换成分词器 String[] tokens = content.toString().split("\\s+"); for (String token : tokens) { if (token.trim().isEmpty()) { continue; } // 输出的 key 是 "类别 + 单词",在 shuffle 阶段自动按类别分组 outKey.set(category + "_" + token); context.write(outKey, one); } } }

这个 Mapper 的输入就是 2.1 里InitSequenceFileJob产出的 SequenceFile,所以 Mapper 的泛型参数是<Text, Text>,而不是常见的<LongWritable, Text>。key 按类别_单词拼接,这样 shuffle 之后同一个类别同一个单词的所有计数会自动汇总到同一个 Reducer。

Reducer 侧直接累加即可。理论上可以用 Combiner 做本地聚合,减少网络传输,但这个数据量(总共 500 多篇文档)不加 Combiner 也完全跑得动。如果你的数据集换成了几十万篇,建议在main方法里加一句job.setCombinerClass(IntSumReducer.class),能肉眼可见地缩短训练时间。

2.3 模型如何落盘:计数结果就是模型文件

朴素贝叶斯的「模型」不是传统意义上的权重矩阵,而是一组计数表。类别先验概率、每个类别下每个单词的出现次数、每个类别的总词数,把这三张表存下来,模型就固化了。

项目里这三个 Job 的输出目录是分开的,训练完成后你在 HDFS 上能看到类似下面的目录结构:

目录内容后续用途
/model/doccount每个类别的文档数计算 P(类别)
/model/wordcount每个类别下每个单词的频次计算 P(词
/model/totalcount每个类别的总词数计算 P(词

这里有个容易被忽略的设计点:wordcounttotalcount必须使用同一套分词规则。如果训练时GetSingleWordCountFromDocTypeJob用空格分词,而GetTotalWordCountFromDocTypeJob用了另一个分词器,两个 Job 统计出来的词表对不上,条件概率的分母和分子就不在同一个量纲下,分类结果会莫名其妙地偏向某个类别,而且你很难从结果里看出来。

我在复现时的做法是:把分词逻辑抽成一个静态工具方法,两个 Job 的 Mapper 都调同一个方法,保证分词规则绝对一致。这个项目因为是课程设计,分词直接用split("\\s+")处理英文语料问题不大,但如果你要处理中文文本,务必先统一分词器,再谈训练。

3. 分类与评估:贝叶斯判定逻辑和 P/R/F1 的计算细节

3.1 贝叶斯判定的核心公式和下溢处理

训练完成后,GetNaiveBayesResultJob读取模型文件,对每个测试文档计算它属于每个类别的后验概率,取概率最大的类别作为判定结果。

朴素贝叶斯的判定公式是:

P(类别 | 文档) ∝ P(类别) × Π P(词 | 类别)

其中 P(词|类别) 在项目里直接用频率估计:

P(词|类别) = 该词在类别中的出现次数 / 该类别总词数

这个公式在实现时有一个经典坑:一篇文档包含几百个词,几百个小于 1 的小数连乘,结果会无限趋近于 0,直接 double 乘法会下溢成 0,导致所有类别的后验概率都是 0,分类结果变成随机猜。常见的解决办法是取对数,把连乘变成连加:

log P(类别 | 文档) = log P(类别) + Σ log P(词 | 类别)

项目里GetNaiveBayesResultJob的 Reducer 侧做计算时,应该就是按这个方式实现的。我建议你检查一下代码里有没有Math.log,如果看到的是直接乘法,那这个代码在小数据集上可能碰巧能跑出结果,但换个大点的测试集大概率翻车。

3.2 未登录词与拉普拉斯平滑的取舍

训练集里没出现过的词,测试文档里是可能出现的。如果某个词在 CHINA 类里出现 0 次,在 CANA 类里出现 20 次,那么这个词的加入会直接让 CHINA 类的条件概率变成 0,进而让整个文档被判到 CANA 类,哪怕其他所有词都强烈指向 CHINA。

解决方法是拉普拉斯平滑,给分子分母同时加一个平滑项:

P(词|类别) = (该词出现次数 + α) / (类别总词数 + α × 词典大小)

α 一般取 1,词典大小是两个类别去重后的总词数。这个项目的数据集只有两个类,我建议你动手把 α 换成 0.5、1、2 各跑一遍,对比 F1 值的变化。这是一个很讨巧的「进阶实验」,写在课程设计报告里能直接加分——因为它说明你理解了平滑项对未登录词的处理逻辑,而不是只会调用现成库。

注意:如果项目代码里没做平滑,先别急着改。跑一遍基线结果,记录 F1,再加上平滑跑一遍,用两组数字的对比来说明平滑的必要性,这比直接上一份改好的代码更有说服力。

3.3 Evaluation.java 里的 P/R/F1 计算逻辑

Evaluation.java负责把分类结果和真实类别做比对,计算三个指标。首先要明确混淆矩阵的四个值:TP(预测为 CHINA 且实际是 CHINA)、FP(预测为 CHINA 但实际是 CANA)、FN(预测为 CANA 但实际是 CHINA)。Precision、Recall、F1 的定义如下:

public class Evaluation { public static void main(String[] args) throws Exception { // args[0]:分类结果文件,每行格式为 "文档ID 预测类别" // args[1]:真实类别文件,每行格式为 "文档ID 真实类别" Map<String, String> predictions = loadResult(args[0]); Map<String, String> truths = loadResult(args[1]); int tp = 0, fp = 0, fn = 0; for (Map.Entry<String, String> entry : predictions.entrySet()) { String docId = entry.getKey(); String predicted = entry.getValue(); String actual = truths.get(docId); if ("CHINA".equals(predicted)) { if ("CHINA".equals(actual)) { tp++; } else { fp++; } } else { if ("CHINA".equals(actual)) { fn++; } } } double precision = (tp + fp) == 0 ? 0.0 : (double) tp / (tp + fp); double recall = (tp + fn) == 0 ? 0.0 : (double) tp / (tp + fn); double f1 = (precision + recall) == 0.0 ? 0.0 : 2 * precision * recall / (precision + recall); System.out.printf("Precision=%.4f Recall=%.4f F1=%.4f%n", precision, recall, f1); } }

注意两个细节。第一,这里只统计了正类(CHINA)的 P/R/F1,CANA 类的指标需要把混淆矩阵反过来再算一遍。如果报告里只写一个类的三个数字,答辩时大概率会被问「另一个类的指标呢」,提前把两类的都算出来列成表格最稳妥。第二,loadResult解析文件时如果文档 ID 在两组数据里对不上,算出来的指标会偏小,我建议在这个方法里加一个Set记录原始 ID 数量,最后输出比对成功了多少条,方便排查数据错位。

4. 从零复现:环境准备、打包、运行顺序与参数说明

4.1 环境准备与工程导入

这份工程里有.iml文件,说明它是标准 IntelliJ IDEA 工程。复现的第一步是在 IDEA 里直接Open这个目录,让 IDEA 按.iml恢复模块配置。Hadoop 相关依赖需要本机已经装好,建议直接在Project Structure → Libraries里把 Hadoop 的share/hadoop/commonshare/hadoop/mapreduceshare/hadoop/hdfs下的 jar 包全部加进来。

环境方面,伪分布式是最省事的方案。Hadoop 2.x 和 3.x 都能跑,关键在core-site.xmlhdfs-site.xmlyarn-site.xml三个配置里把fs.defaultFS指向hdfs://localhost:9000,别让它默认走本地文件系统——否则hadoop jar提交作业时读写路径的行为会和预期不一样。Windows 上跑伪分布式还要额外配HADOOP_HOME环境变量和winutils.exe,这个坑比较多,我在第 5 章专门写。

4.2 数据准备与目录规划

原始数据在NBCorpus\Country下,先把 CHINA 和 CANA 两个文件夹上传到 HDFS:

hdfs dfs -mkdir -p /input/CHINA /input/CANA hdfs dfs -put NBCorpus/Country/CHINA/* /input/CHINA/ hdfs dfs -put NBCorpus/Country/CANA/* /input/CANA/

注意:训练集和测试集的 70%/30% 切分,要在上传之前就完成。合理做法是在本机把两个文件夹各自随机抽 70% 放进train子目录,剩下 30% 放进test子目录,再分别上传。如果直接把全部数据传上去,然后在 HDFS 上用mv来回挪文件切分,后续写文档的时候很难说清楚切分逻辑,答辩时也容易留下漏洞。

4.3 五个 Job 的运行顺序与参数传递

整个流程按顺序执行六个步骤,前五步是 MapReduce 作业,最后一步是本地 Java 类:

# 1. 把原始文本转成 SequenceFile hadoop jar nb-classifier.jar InitSequenceFileJob \ /input/train /seq/train # 2. 统计每个类别的文档数 hadoop jar nb-classifier.jar GetDocCountFromDocTypeJob \ /seq/train /model/doccount # 3. 统计每个类别下每个单词的出现次数 hadoop jar nb-classifier.jar GetSingleWordCountFromDocTypeJob \ /seq/train /model/wordcount # 4. 统计每个类别的总词数 hadoop jar nb-classifier.jar GetTotalWordCountFromDocTypeJob \ /seq/train /model/totalcount # 5. 读模型文件,对测试集分类 hadoop jar nb-classifier.jar GetNaiveBayesResultJob \ /seq/test /model/result # 6. 本地运行评估类,输出 P/R/F1 hadoop jar nb-classifier.jar Evaluation \ /model/result/part-r-00000 /input/test/truth.txt

运行顺序上有先后依赖:前四步是训练链,必须按 1→2→3→4 来,因为后三个 Job 的输入都是第 1 步产出的/seq/train;第 5 步需要前四步的产物全部就绪,少一个目录都会在读取模型时报 FileNotFoundException。最后一步纯属本地计算,读分类结果和真实类别文件,跑完在控制台打印三个指标。

提示:如果某个 Job 参数写错了,不需要全部重跑。比如第 4 步跑挂了,模型目录里/model/doccount已经存在,只要把第 4 步重新执行即可,前两个 Job 的产物可以直接复用。前提是代码逻辑没改,改了就必须从头重跑。

参数传递这块有个值得注意的细节:InitSequenceFileJob在用FileInputFormat.addInputPath()添加输入路径时,可以直接把/input/train这个父目录传进去,Hadoop 会自动递归扫描下面的CHINACANA子目录。但这样做有个隐患——GetDocCountFromDocTypeJob统计文档数时拿到的文档 ID 如果只包含文件名不包含上级目录名,CHINA_001CANA_001这种同名文件会互相覆盖。稳妥的做法是上传数据前就把文档重命名带上类别前缀,比如CHINA_001.txtCANA_001.txt,一劳永逸。

5. 常见问题排查:我在跑这个项目时踩过的五个坑

5.1 大量任务失败,日志提示容器内存不足

现象:跑GetSingleWordCountFromDocTypeJob时,Map 任务大面积失败,YARN 界面里看到 Container 被 kill,日志里出现「Container is running beyond physical memory limits」。

原因:伪分布式环境默认给每个容器分配的内存很小(通常是 1GB),而GetSingleWordCountFromDocTypeJob这个 Job 的 Mapper 把整个文档内容放进 value 里做分词,遇到长文档时堆内存会顶到上限。YARN 检测到物理内存超标,直接杀掉容器。

解决:在yarn-site.xml里调大容器内存上限,改完必须重启 YARN 才生效:

<property> <name>yarn.nodemanager.vmem-pmem-ratio</name> <value>3</value> </property> <property> <name>yarn.nodemanager.pmem-check-enabled</name> <value>false</value> </property>

我当时的做法是临时关掉物理内存检查,等作业跑完再恢复。线上环境不建议这么干,但课程设计里这是最快见效的手段。

5.2 分类结果里所有文档都被判到同一个类别

现象:GetNaiveBayesResultJob跑完,打开结果文件发现 CHINA 和 CANA 的测试文档全部被分成了一类,分布极其不均衡。

原因:最常见的是训练和测试的数据切分有问题,比如两个类别的训练集文档数差距悬殊,导致先验概率 P(类别) 严重失衡。也有可能是代码里没做对数运算,条件概率连乘下溢成 0,最后比较的是两个都是 0 的概率,Reduer 里的>比较符自然会让代码走向第一个分支。

解决:先看/model/doccount里的统计结果,确认两个类别的训练文档数符合 70% 切分预期;再把GetNaiveBayesResultJob里算概率的部分加上Math.log,并输出每个文档在各类别下的 log 后验值,肉眼对比到底差在哪里。曾经真的遇到过一次代码里把>写成了>=,某个类在 log 值相等时始终被选中,改回来就正常了。

5.3 分词不一致导致的条件概率错乱

现象:训练完成后用测试集评估,F1 值只有 0.5 左右,比随机略好。看/model/wordcount里的单词表,发现 CHINA 类下冒出了大量不相关的单词。

原因:前面说过,GetSingleWordCountFromDocTypeJobGetTotalWordCountFromDocTypeJob必须用同一套分词规则。我当时改过其中一处的分隔符正则,从\\s+改成了[^a-zA-Z]+,结果两个 Job 的词表量级差了一倍,条件概率的分母和分子根本不对应。

解决:把分词逻辑抽成公共工具方法,两个 Job 的 Mapper 都调同一个方法。这个项目规模小,直接在 Mapper 里定义static方法就行。从那以后我每次跑完训练都会先对比 wordcount 和 totalcount 两个输出里单词表的前 20 个词,词表对不上就直接重跑,不做无用功。

5.4 在 Windows 上提交作业报错:Failed to locate the winutils binary

现象:在本地 IDE 里直接跑main方法报错,提示找不到winutils.exe,或者 Hadoop 文件系统操作全部返回null

原因:Hadoop 的本地库是 Linux 原生的,Windows 下跑伪分布式必须有一份本地适配。很多人在这卡住,因为HADOOP_HOME配了,但winutils.exe没放进对应目录。

解决:下载对应 Hadoop 版本的winutils.exe放到$HADOOP_HOME/bin下,同时在 IDEA 的 Run Configuration 里加 VM 参数:

-Dhadoop.home.dir=D:/hadoop-3.3.4

如果还报权限相关的问题,多半是缺hadoop.dll,把可执行文件同目录下的动态库一起放进去再试。这个问题和 HDFS 路径无关,优先级最高,先解决它再谈跑 Job。

5.5 测试集与训练集样本重叠,F1 虚高导致的误判

现象:评估结果 F1 高达 0.98,怎么看都不真实。检查数据发现训练集里出现了和测试集一模一样的文档。

原因:切分训练集和测试集时没有做样本去重。NBCorpus 这个数据集本身是按主题归档的,有的新闻稿是同一篇文章的转载,标题相同、内容只有个别字差异,切分时很容易让同一篇文档同时出现在训练集和测试集里。

解决:切分前先对文档内容做一次去重,用文档内容的前 100 个字符做哈希,哈希值相同的只保留一条。我当时为了省事用文件名去重,结果漏掉了转载的情况,被答辩老师一眼看穿。去重逻辑不要放进 MapReduce 里,直接在本地 Python 或 Java 里处理完再上传 HDFS 最稳妥。

6. 进阶验证:手动算一个样例,把分类器的黑匣子打开

6.1 手工复算一遍贝叶斯判定

跑通只是第一步,真正理解这份代码的 方式是拿一个测试文档,手动算一遍分类结果,和GetNaiveBayesResultJob的输出对比,确认一致。这个验证过程也能帮你发现模型文件里可能存在的隐藏问题。

假设训练模型统计出如下数据:

| 类别 | 文档数 | 总词数 | P("trade")|类别) | P("oil")|类别) | |---|---|---|---|---| | CHINA | 178 | 52000 | 0.0021 | 0.0018 | | CANA | 184 | 48000 | 0.0008 | 0.0032 |

现在有一个测试文档,分词后包含trade出现 3 次、oil出现 2 次,先验概率 P(CHINA)=178/362≈0.4917,P(CANA)=184/362≈0.5083。

取对数计算:

log P(CHINA|doc) = log(0.4917) + 3×log(0.0021) + 2×log(0.0018)

≈ -0.71 + 3×(-6.17) + 2×(-6.32) ≈ -31.76

log P(CANA|doc) = log(0.5083) + 3×log(0.0008) + 2×log(0.0032)

≈ -0.68 + 3×(-7.13) + 2×(-5.75) ≈ -33.57

CHINA 类的 log 后验值更大,分类器判定为 CHINA。打开/model/result对应行,如果输出也是 CHINA,说明模型的概率计算链路是通的;如果输出不一致,优先检查分母用的总词数是不是和单词频次来自同一个训练集。

6.2 从二分类到多分类的扩展思路

这个项目固定了 CHINA 和 CANA 两个类,但代码设计上统计 Job 的输出 key 都是「类别 + 词」的形式,天然支持多类别。想扩展成三分类,只需要把新的类别文件夹放进来重新跑训练链,GetDocCountFromDocTypeJob会多输出一行文档数,GetNaiveBayesResultJob里的类别列表如果是从模型文件动态读出来的,那 Reducer 侧不用改任何代码。

如果类别是从代码里硬编码的,建议改成遍历模型文件里的 key 前缀,这样加类别时不用重新编译。调整拉普拉斯平滑系数 α 也是一个值得做的实验:α 从 0.1 到 2.0 按步长 0.1 各跑一次,画一条 F1 随 α 变化的曲线,放进报告里比任何文字都直观。

验证时有一件事要留意:GetNaiveBayesResultJobEvaluation.java是拆开的两个程序,中间靠 HDFS 文件传递数据。中间结果一旦被覆盖,评估阶段输入的可能是上一轮的分类结果。所以每次重跑训练链之前,先把/model/目录删干净,再从头执行,宁可多花两分钟也不要拿脏数据去评估。

从那以后,我每次复现这类分类项目,都会强制走一遍这个流程:跑通全链路 → 手工算一个文档核对分类结果 → 对比两个类别的 P/R/F1 → 删干净输出目录重跑一次。四个步骤做完,代码里的问题基本能暴露八成。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询