
简介这是一份基于Hadoop MapReduce实现的朴素贝叶斯文本分类器完整项目面向正在学习大数据与文本挖掘的在校学生以及需要完成课程设计或毕业设计的开发者。项目覆盖贝叶斯分类器训练、测试集分类和Precision、Recall、F1评估的完整流程选用NBCorpus中CHINA与CANA两类文本作实验样本按70%和30%划分训练集与测试集并配套说明文档解释各个MapReduce作业的用途。压缩包内共552个文件、约3.75MB其中518个txt为实验语料与分类结果9个java为核心源码另有png截图、md与docx/PDF文档便于查看运行流程与项目报告。已有230人学习浏览适合Hadoop课程作业、期末项目或入门实战源码结构清晰稍加修改即可迁移到其他文本分类任务。1. Hadoop 上的朴素贝叶斯文本分类器当单机内存撑不住时这是最务实的一条路训练语料从几千篇涨到上百万篇时单机版的朴素贝叶斯文本分类器会先让内存告急再让训练时间以小时计。很多人第一反应是换复杂模型但落地时最先见效的往往是先把朴素贝叶斯在 Hadoop 上跑通——它对 MapReduce 的契合度几乎是天然的。这篇笔记围绕「基于 Hadoop 开发实现的朴素贝叶斯文本分类器 源代码 文档说明」这个交付物展开从原理和选型讲到中文文本预处理与词频统计再到训练与预测的 MapReduce 代码走读最后落在参数调优和踩坑记录。适合正在做课程设计、毕业设计或刚把数据处理迁移到 Hadoop 上的工程师照着复现。2. 朴素贝叶斯为什么适合 Hadoop先弄清要算的两批概率2.1 训练阶段算先验与条件概率公式拆成 MapReduce 视角先看一个反直觉的现象很多人以为朴素贝叶斯太「朴素」、不可能是生产级方案但它在垃圾邮件过滤、新闻分类、短文本情感判断这些任务上只要特征工程做得扎实效果能逼近甚至超过调参不当的 BERT 基线。原因在于文本分类的决策边界往往非常宽朴素贝叶斯对特征独立性的「错误假设」反而在小样本下起到了正则化作用。理解这一点你才不会在实现过程中动不动怀疑算法本身。朴素贝叶斯的分类决策依据是后验概率给定一篇文档 d它属于类别 c 的概率 P(c|d) 正比于 P(c) 乘以 P(d|c)。在文本分类里P(d|c) 会被拆成文档中每个词的条件概率的乘积也就是 P(w1|c) × P(w2|c) × … × P(wn|c)。这种拆分依赖特征独立性假设它假设文档里的词在给定类别下互不影响。这个假设在严格意义上不成立但实践中足够好用计算量也友好。于是整个训练过程要算的就是两批数。第一批是先验概率 P(c)等于类别 c 的文档数除以总文档数。第二批是条件概率 P(wi|c)等于词 wi 在类别 c 的文档里出现的总次数除以类别 c 下所有词出现次数之和。这里用的是多项式朴素贝叶斯模型它把词频信息完整保留下来。也有人用伯努利模型只看词是否出现不考虑频次在长文档分类时多项式模型通常表现更好。具体计算时文档 d 属于类别 c 的得分是 log P(c) sum(log P(wi|c))。加 log 是为了把连乘变连加防止几十个概率相乘导致浮点下溢。这一点在 Hadoop 实现里尤为重要因为文档越长、词表越稀疏原始概率乘积越容易直接变 0.0。我在后面第 4 章的代码走读里模型文件里存的就直接是 log 概率预测时只需做加法。套到前面的课程设计和面试题里先验概率是「类别文档数 / 总文档数」条件概率是「词频 / 类词频总和」再加上平滑项。这三个数都带有「除以总数」的归一化而这正是 reduce 阶段收拢所有统计之后最自然的动作。2.2 MapReduce 的天然契合点统计就是最标准的并行场景再往下拆一层。朴素贝叶斯的训练数据无非是大量打好类别标签的文本。单机跑的时候你要么把所有文本读进内存维护一个 HashMapString, Long 做词频累加要么用数据库临时表一行一行 insert。前者在几百万篇文档、几十万词表的规模下很容易把堆内存打满后者则慢到让人怀疑人生。Hadoop 的做法完全不同。HDFS 把训练文本切成若干分片每个分片交给一个 map 任务做局部统计map 输出「类别词 → 次数」交给 reduce 任务归并。归并完成之后每个类别下每个词的全局词频就拿到了。拿到的结果本身就是键值对形态可以直接序列化成模型文件。整个过程中单个节点只处理自己那部分数据内存压力被分拆掉吞吐量随节点数近似线性扩展。有人会问既然统计这么简单为什么不用 Spark 或者 Flink这个问题的答案取决于你的环境。如果集群里已经跑着 Spark用 DataFrame 做 groupBy 确实更省事但如果你的生产环境就是 Hadoop Hive 的存量设施新增一个 Spark 组件要过审批、要维护资源队列MapReduce 反而零额外依赖。而且朴素贝叶斯训练是单次扫描、单次归并MapReduce 的两次串行在这里没有太多效率损失。对比之下迭代式算法KMeans、LR 这种在 MapReduce 上每轮都要落盘一次那才是真正的痛点。朴素贝叶斯恰恰避开了这个弱点。另一个隐藏的好处是模型天然小。朴素贝叶斯训练完模型只包含「每个类别的先验」和「每个类别下每个词的 log 概率」几十万词表也就是几十 MB 的文本文件。这个文件可以直接从 HDFS 拉回本地、放进 Redis 或者本地内存不依赖任何 Hadoop 组件就能做预测。相比动辄几个 GB 的神经网络 embedding这个特性让它特别适合课程设计、轻量离线服务和资源受限的团队。不过这里要澄清一个常见误解朴素贝叶斯在 Hadoop 上的优势在训练不在预测。预测阶段对单条文本算后验概率计算量极小本地脚本几毫秒就完成。如果你把预测也做成一个 MapReduce 作业反倒会因为作业启动开销JVM 启动、资源申请拖慢响应。所以典型做法是用 Hadoop 做批量训练和离线预测比如每天对当天新增文档打标签线上实时预测继续用本地加载的模型文件。这个分工后面会展开。初学阶段在伪分布式上跑通整套流程是最快的入门路径。伪分布式的坑主要在 core-site.xml 和 hdfs-site.xml 的路径冲突、datanode 起不来、NodeManager 误杀容器这几类这些不属于本项目的核心范围但如果你搭建环境时反复翻车可以优先检查fs.defaultFS是否写成了hdfs://localhost:9000、dfs.replication是否设为 1、YARN 的虚拟内存检查是否关掉。环境稳住之后再把同样的代码放到多节点集群上你会发现真正要调的参数从「能不能跑」变成了「内存给多少、reducer 设几个、要不要加 Combiner」。注意MapReduce 的并发度同时受mapreduce.job.reduces和 YARN 队列容量限制。伪分布式上所有任务在同一 JVM 里串行问题不明显上了集群shuffle 和数据倾斜问题会全部冒出来。训练输入的文件格式建议用 HDFS 上的普通文本文件一行一条样本格式固定为「类别\t分词后的正文」。这里类别在前、正文在后中间用 Tab 分隔是为了让 Mapper 能用最简单的split(\t, 2)切开避免正文里混入分隔符导致解析错位。为什么不用空格分隔因为正文分词后就可能包含空格再拿空格当分隔符会拆错。为什么不用逗号因为中文逗号和西文逗号在新闻正文里随处可见。Tab 是纯文本里最不容易出现的字符这一点看着小实际能帮你省掉大量解析 bug。3. 数据预处理与词频统计把文本变成 Hadoop 能吃的格式3.1 中文分词与停用词过滤从原始语料到「类别文本」的规范输入任何一个文本分类项目数据预处理花的时间都超过模型训练。中文文本没有天然空格必须先分词。常见的做法是接一个分词工具Java 项目用 IKAnalyzer 或 HanLPPython 项目直接上 jieba。我在 Hadoop 训练链路里通常的做法是先用 Python 脚本离线把原始语料分好词输出「类别\t空格分隔的词序列」到本地文件然后hdfs dfs -put上传到 HDFS。训练代码里不再做分词只按空格切词。这样做的好处是分词逻辑和训练逻辑解耦换分词器只需要重新跑一遍预处理脚本不碰 MapReduce 代码。下面这个 Python 脚本是预处理阶段的最小可运行版本兼容 Python 3#!/usr/bin/env python3 # -*- coding: utf-8 -*- 预处理原始语料 - 类别\t分词正文 import re import sys # 实际项目里替换成 jieba.cut 可以提高中文分词质量 def tokenize(text: str) - list[str]: # 简易分词按 2 个连续汉字切词同时过滤数字和英文 tokens re.findall(r[\u4e00-\u9fa5]{2,}, text) return tokens STOPWORDS set(的 了 和 是 在 我 你 他 这 那 有 就 不 也 都 而 说 中 为 与.split()) def main(): for line in sys.stdin: line line.strip() if not line: continue # 原始输入格式类别\t正文正文未分词 label, _, content line.partition(\t) if not content: continue words [w for w in tokenize(content) if w not in STOPWORDS] if words: print(f{label}\t{ .join(words)}) if __name__ __main__: main()这段代码的核心是partition(\t)而不是split(\t, 1)partition 只切第一段剩余内容全部保留在其中不会因为正文里混入 Tab 而出错。tokenize 函数里用了正则[\u4e00-\u9fa5]{2,}匹配连续两个以上的汉字把单个字和纯数字、英文过滤掉。这个策略对新闻文本足够用但如果你做的是商品评论或短文本建议换 jieba 并把用户词典加载进去否则像「不划算」「贼拉好」这类口语化词会被切碎。跑完脚本后输出长这样体育 姚明 退役 之后 首次 现身 上海 参加 公益 活动 财经 央行 发布 报告 称 物价 总体 平稳每一行对应一篇文档。把这些行合并成一个文本文件一行一条样本上传到 HDFS 上。文件太小会让 Hadoop 切分不充分一般建议凑到至少 128 MB对应默认 block size不足的话在 HDFS 上用hdfs dfs -put上传时不需要特殊处理但 NameNode 上小文件数量会膨胀。关于小文件的问题第 5 章会专门讲。注意HDFS 默认不压缩存储训练文件传上去占用双倍空间本地一份、HDFS 三副本。如果语料很大可以先把文件用 gzip 压缩再上传Hadoop 会自动解压 gzip 后缀的文本文件不需要改代码。3.2 词频统计的 MapReduce 骨架第一个可运行的 Job预处理之后先跑一个词频统计 Job 验证环境和数据链路。这一步不是浪费它能帮你确认输入文件能读、HDFS 路径配置正确、Reducer 能正常写出结果。如果这个最小 Job 跑不通后面训练代码跑不通时你根本分不清是环境问题还是算法问题。下面是标准的 WordCount 变体输出「词 → 全局词频」。我用 Java 写对应 Hadoop 2.x/3.x 的 APIimport org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.Reducer; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; public class WordCount { public static class TokenizerMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text word new Text(); Override public void map(Object key, Text value, Context context) throws IOException, InterruptedException { // 输入文件每行是类别\t词1 词2 词3 ... String[] parts value.toString().split(\t, 2); if (parts.length 2) return; String[] words parts[1].split(\\s); for (String w : words) { word.set(w); context.write(word, one); } } } public static class SumReducer extends ReducerText, IntWritable, Text, IntWritable { private IntWritable result new IntWritable(); Override public void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) { sum val.get(); } result.set(sum); context.write(key, result); } } public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, word count); job.setJarByClass(WordCount.class); job.setMapperClass(TokenizerMapper.class); job.setCombinerClass(SumReducer.class); job.setReducerClass(SumReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }这里有两处值得说明。第一Mapper 里我按\t把行切成「类别」和「正文」两段只统计正文部分避免类别标签混进词表。如果你用的是公共分类语料类别名也是合法中文词不切的话「体育」「财经」这些词会被当成特征词分类时形成奇怪的误导。第二我设置了setCombinerClass(SumReducer.class)combiner 会在 map 端做一次局部求和把「词 → 大量 1」合并成「词 → 局部计数」再进入 shuffle。这一步能把网络传输量减少一到两个数量级是整个训练场景里性价比最高的优化。第 6 章会继续展开 Combiner 的细节。编译和运行命令打包成 jar 后# 打包假设项目根目录是 nbproject mvn clean package -DskipTests # 上传训练数据到 HDFS hdfs dfs -mkdir -p /nb/data hdfs dfs -put train.txt /nb/data/ # 提交 job输出目录必须不存在 hadoop jar target/nbproject-1.0.jar WordCount /nb/data/train.txt /nb/data/wc_out # 查看前几行输出 hdfs dfs -cat /nb/data/wc_out/part-r-00000 | head -20有一个坑在这里提前说输出目录/nb/data/wc_out在提交前必须不存在否则 Hadoop 会直接报FileAlreadyExistsException。很多新手在迭代调试时习惯固定一个输出路径第一次跑通之后第二次就翻车。我一般会在脚本里先执行hdfs dfs -rm -r -f清掉旧目录再提交作业。这个习惯到训练阶段会帮你省下大量返工。4. 训练与预测的完整实现两个 Job 跑出朴素贝叶斯模型4.1 训练 Job 的 Mapper/Partitioner/Reducer 设计与代码正式的朴素贝叶斯训练我拆成两个 MapReduce Job各干一件事第一个统计「类别 → 文档数」输出先验概率所需的分母第二个统计「类别词 → 词频」输出条件概率所需的分子。最后在本地合并两边结果成模型文件。为什么拆成两个而不是在一个 Job 里做完因为一个 Job 里要拿全局文档总数算先验概率要么依赖 Counter 的最终聚合值要么依赖 Reducer 的执行顺序这些在 YARN 上都有不确定性。分开两个 Job每个都只做纯分组统计逻辑最简单也最容易验证。Job1 的 Mapper 和 Reducer 如下。Mapper 每读一行解析出类别输出类别 → 1Reducer 累加后输出类别 → 文档数public class DocCountJob { public static class DocCountMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text category new Text(); Override public void map(Object key, Text value, Context context) throws IOException, InterruptedException { String[] parts value.toString().split(\t, 2); if (parts.length 2) return; category.set(parts[0].trim()); context.write(category, one); } } public static class DocCountReducer extends ReducerText, IntWritable, Text, IntWritable { Override public void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { int sum 0; for (IntWritable val : values) sum val.get(); context.write(key, new IntWritable(sum)); } } public static void main(String[] args) throws Exception { Configuration conf new Configuration(); Job job Job.getInstance(conf, doc count); job.setJarByClass(DocCountJob.class); job.setMapperClass(DocCountMapper.class); job.setCombinerClass(DocCountReducer.class); job.setReducerClass(DocCountReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }这个 Job 几乎就是 WordCount 换了个 key。真正有意思的是 Job2需要按类别分区让同一个类别的所有词进同一个 Reducer。这样每个 Reducer 只在本地就能拿到「类别 c 下所有词的词频总和」直接算出条件概率不用二次扫描。Job2 的 Mapper 解析出类别和词输出类别词 → 1public class TermCountJob { public static class TermCountMapper extends MapperObject, Text, Text, IntWritable { private final static IntWritable one new IntWritable(1); private Text outKey new Text(); Override public void map(Object key, Text value, Context context) throws IOException, InterruptedException { String[] parts value.toString().split(\t, 2); if (parts.length 2) return; String category parts[0].trim(); String[] words parts[1].split(\\s); for (String w : words) { if (w.isEmpty()) continue; outKey.set(category w); context.write(outKey, one); } } } // ... Partitioner / Reducer 见下文 }然后是按类别哈希的 Partitioner。这里的关键是类别必须被完整哈希不能拿类别词整体哈希否则同一个类别会散到多个 Reducer 里public static class CategoryPartitioner extends PartitionerText, IntWritable { Override public int getPartition(Text key, IntWritable value, int numPartitions) { String keyStr key.toString(); int atIndex keyStr.indexOf(); String category keyStr.substring(0, atIndex); return (category.hashCode() Integer.MAX_VALUE) % numPartitions; } }Reducer 端是本 Job 的核心。reduce 方法把同一个类别词的计数累加存进本地 HashMap同时维护该类别的词频总和。cleanup 方法在这个 Reducer 处理完所有 key 之后调用一次此时本地 HashMap 里已经有了这个类别完整的词统计public static class TermCountReducer extends ReducerText, IntWritable, Text, Text { private MapString, Long wordCount new HashMap(); private long totalWordCount 0; private String currentCategory null; Override public void reduce(Text key, IterableIntWritable values, Context context) { String keyStr key.toString(); int atIndex keyStr.indexOf(); String category keyStr.substring(0, atIndex); String word keyStr.substring(atIndex 1); if (currentCategory null) currentCategory category; long sum 0; for (IntWritable val : values) sum val.get(); wordCount.merge(word, sum, Long::sum); totalWordCount sum; } Override protected void cleanup(Context context) throws IOException, InterruptedException { if (currentCategory null) return; long vocabSize wordCount.size(); Text outKey new Text(); Text outValue new Text(); for (Map.EntryString, Long entry : wordCount.entrySet()) { // 拉普拉斯平滑分子加 1分母加词表大小避免零概率 double prob (entry.getValue() 1.0) / (totalWordCount vocabSize); outKey.set(currentCategory \t entry.getKey()); outValue.set(String.valueOf(prob)); context.write(outKey, outValue); } } }这里有个细节Reducer 的currentCategory只取第一个遇到的类别名。按分区器的逻辑一个 Reducer 只会收到一个类别的数据所以这种做法成立。但如果在集群上别人改了分区器或者 reducer 数量比类别数少两个类别就会被合并到一个 ReducercurrentCategory只会是第一个第二个类别的词统计会被错误地算进第一个类别里。保险做法是让currentCategory赋值发生在每次进入一个新的类别时然后把多个类别分别处理。我在实际代码里用MapString, MapString, Long直接存多类别但为了这里讲清楚先用单类别版本。4.2 模型文件格式落盘之后用脚本直接加载Job2 输出的每行是「类别、词、条件概率」三列用 Tab 分隔。类别和词中间不再是而是独立的两列这样后续脚本读取时按\tsplit 成三份即可不会混淆。先验概率还没算进去需要结合 Job1 的「类别 → 文档数」输出。合并模型的过程我习惯用一个 Python 脚本在本地做把两个 Job 的输出从 HDFS 拉下来按类别对齐把先验概率和条件概率都取 log 后写进最终模型。模型文件最终格式长这样财经 __PRIOR__ -2.302585 财经 央行 -6.214608 财经 物价 -7.130899 体育 __PRIOR__ -1.609438 体育 姚明 -5.521461__PRIOR__是一个保留词条专门存先验概率的 log 值。预测时初始化每个类别的分数为先验 log 值然后逐个词累加条件概率的 log 值。这样做的好处是模型文件本身就是纯文本、按行读取、不需要反序列化。一个几百 MB 的模型文件用cat model.txt | head就能抽查避免了黑匣子式的加载流程。合并脚本很短核心逻辑是读两个文件、按类别分组、输出 log 概率# 从 HDFS 拉取两个 Job 的输出 hdfs dfs -getmerge /nb/doc_count_out local_doc_count.txt hdfs dfs -getmerge /nb/term_count_out local_term_count.txt # 合并脚本 python3 merge_model.py local_doc_count.txt local_term_count.txt model.txt#!/usr/bin/env python3 import math import sys doc_count_file, term_count_file, out_file sys.argv[1], sys.argv[2], sys.argv[3] # 读取每个类别的文档数 doc_counts {} with open(doc_count_file, encodingutf-8) as f: for line in f: cat, cnt line.rstrip(\n).split(\t) doc_counts[cat] int(cnt) total_docs sum(doc_counts.values()) with open(out_file, w, encodingutf-8) as out: # 写先验概率 for cat, cnt in doc_counts.items(): prior math.log(cnt / total_docs) out.write(f{cat}\t__PRIOR__\t{prior}\n) # 写条件概率 with open(term_count_file, encodingutf-8) as f: for line in f: cat, word, prob line.rstrip(\n).split(\t) log_prob math.log(float(prob)) out.write(f{cat}\t{word}\t{log_prob}\n)交付物里的文档说明也应该围绕模型格式展开。README 里至少写清三件事输入文件格式、模型文件每一列的含义、训练与预测两个 Job 的启动命令。我在交付课程设计时还会附一张参数表标明每个可调参数改在哪个文件里的哪一行。有时候老师或面试官不会一行行读代码但一定会看文档能不能让人照着跑起来。4.3 预测阶段的两种做法Hadoop 批量预测与本地脚本预测训练完成后预测分两条路。离线批量预测把模型文件放进 DistributedCacheMapper 在 setup 中加载模型对每条待预测文本输出「文档ID、预测类别、得分」。实时预测则不经过 Hadoop直接把 model.txt 拉回本地用 Python 加载进字典单条文本几毫秒出结果。本地预测脚本是最容易验证结果的实现#!/usr/bin/env python3 def load_model(path): model {} with open(path, encodingutf-8) as f: for line in f: cat, word, score line.rstrip(\n).split(\t) model.setdefault(cat, {})[word] float(score) return model def predict(model, words): best_cat, best_score None, float(-inf) for cat, probs in model.items(): score probs.get(__PRIOR__, 0.0) for w in words: # 未登录词直接跳过在取 log 的条件下等价于加 log(1)0 score probs.get(w, 0.0) if score best_score: best_cat, best_score cat, score return best_cat, best_score if __name__ __main__: model load_model(model.txt) # 假设已经按第 3 章的方式分好词 words 央行 发布 报告 称 物价 平稳.split() print(predict(model, words))这段代码唯一要留意的点是未登录词直接在循环里跳过。因为在 log 空间里某个词在某个类别下不存在等价于概率为 0取 log 是负无穷做加法会把整个分数拉成负无穷。而实际朴素贝叶斯里对训练集之外的词我们用平滑后的概率处理这里直接跳过等于假设它的 log 概率是 0虽然不严格但对排序结果影响很小。如果你要更严谨可以在 load_model 时为每个类别额外记一个__UNK__概率取所有未登录词共享一个极小值。5. 避坑指南Hadoop 跑朴素贝叶斯常见的 5 个问题5.1 零概率把整篇文档判死拉普拉斯平滑是底线现象模型训练完之后随便拿一条训练集里的样本去预测居然分错类细看某个类别的分数直接是负无穷。原因测试文本里出现了训练时某个类别下没见过的词。如果不做平滑P(wi|c) 直接取到 0整篇文档的条件概率乘积变成 0log 一下就是负无穷。出现一个这样的词整个类别的后验分数就崩了。解决在条件概率计算时做拉普拉斯平滑。第 4 章代码里的公式(count 1) / (total vocabSize)已经处理了这个问题。注意分子加 1 是超参数可以用 add-1拉普拉斯或 add-k 进行微调词表很大的时候add-1 对低频词的惩罚仍然偏重可以把 k 调到 0.1 试试。这个参数在模型合并脚本里改一行就行。5.2 海量小文本文件压垮 NameNode用 CombineTextInputFormat现象训练语料是几万个小 JSON 或小 txt每个只有几百字节。上传到 HDFS 后集群变慢NameNode 内存涨得很快job 运行前 split 列表长长一串。原因HDFS 每个文件、每个目录都是一条元数据记录。NameNode 内存堆积而且FileInputFormat默认按 block 切分小文件每个都单独起一个 map 任务任务调度开销远大于计算开销。解决把多个小文件合并成大文件再上传是最直接的办法如果数据源持续产生小文件可以在代码里设置CombineTextInputFormat作为输入格式。它会把多个小文件合并到一个 split减少 map 数量。代码上只需要在 Driver 里加一行job.setInputFormatClass(CombineTextInputFormat.class);如果是万条级别的新闻语料我建议预处理时直接cat *.txt all.txt合并成单文件训练前再也不要碰散文件。5.3 中文编码乱码与回车符UTF-8 之外的隐形坑现象本地 Windows 上标好的训练数据传上去MapReduce 跑了不报错但词频统计里全是\ufeff或者\r中文词一个也没统计到。原因Windows 记事本保存 UTF-8 文件会带 BOM\ufeffCRLF 行尾让split(\t)后正文末尾多一个\r这两个字符在 Hadoop 里都不会被自动清洗。解决预处理脚本读写时显式指定encodingutf-8-sig来去除 BOM跑数据清洗时把\r替换掉。检查方法很简单hdfs dfs -cat /nb/data/train.txt | head -n 1 | xxd | head -n 2如果第一行开头是ef bb bf就是 BOM 没处理干净。血泪经验宁可预处理脚本多写 5 行清洗代码也绝不在 Hadoop 端依赖魔数处理。5.4 类别不均衡导致的数据倾斜分区策略与采样现象语料里「体育」类有 80 万篇「文化」类只有 2 万篇。训练 Job 跑的时候别的 Reducer 都完事了处理体育类的 Reducer 还要跑很久。原因前面第 4 章的 Partitioner 按类别哈希使得一个类别进一个 Reducer。类别文档数差异大每个类别下词频总和的差异也大自然出现长尾。解决有三个方向。可以在预处理阶段对多数类欠采样让类别比例不要超过 10:1也可以把分区粒度改到「类别内哈希取模」让一个类别的词分散到多个 Reducer代价是条件概率分母需要在合并阶段二次汇总最粗暴但有效的办法是按比例调mapreduce.job.reduces数量太少导致一个类别独占数量太多导致大量空 Reducer 空跑。具体设多少我在第 6 章的调参部分会给出一个经验区间。5.5 Reduce 端 OOM内存参数与 Combiner 的配合现象Job2 跑到 reduce 阶段某个 Reducer 报Java heap spaceJob 反复重试后失败。原因Reducer 的 cleanup 之前要把该类别所有词频缓存进 HashMap。如果一个类别的词表特别大几十万个词默认 1 GB 的 reducer 堆内存可能不够另外如果没有设置 Combiner海量的「词→1」shuffle 到 reduce 端也会放大内存压力。解决先检查是否设置 Combiner再加内存参数。常用配置如下property namemapreduce.reduce.memory.mb/name value2048/value /property property namemapreduce.reduce.java.opts/name value-Xmx1800m/value /property注意reduce.memory.mb是容器内存上限reduce.java.opts是 JVM 堆大小两者要配套堆大小必须小于容器上限一般留 10% 给 JVM 自身使用。如果你同时开了 10 个 reducer 且每个要 2 GBYARN 队列得能分配出这么多内存否则作业一直 pending这是另一个容易被误判为「卡死」的情况。6. 进阶与验证让分类器可解释、可调优的收尾技巧6.1 用 Combiner 做局部合并让 Shuffle 流量降一个量级训练 Job2 的 map 输出是类别词 → 1如果文章很长、词频很高shuffle 数据量会非常吓人。Combiner 在 map 端先把局部词频加起来输出就变成类别词 → 局部计数shuffle 数据量可以降到原来的 1/10 甚至更低。实现上直接复用求和 Reducer 即可只要 reducer 满足「输入输出键值类型一致」就能当 Combiner。需要注意的是 Combiner 不是二次 Reduce它不保证对同一个词只合并一次所以 Combiner 里的逻辑必须满足交换律和结合律求和没问题求均值就不行。验证 Combiner 效果最简单的方式是看 job 的 CounterCombine output records和Reduce input records。理想情况下后者比前者小一个量级。我一般会先在 1 万条样本上对比开与不开的运行时确定收益后再全量跑。关于调参这里给一个经验区间reducer 数量按类别数的 1 到 2 倍设置避免一个类别独占一个 reducermap 端内存mapreduce.map.memory.mb保持默认 1 GB 就够真正吃内存的是 reducer。如果类别数超过 100优先考虑分段统计而不是一个 reducer 全扛。6.2 交叉验证与混淆矩阵别只盯着准确率一个数字跑通之后最容易被忽视的是验证。拿一部分数据当训练集、一部分当测试集是最基础的更好的做法是 K 折交叉验证把数据切 K 份每次拿 K-1 份训练、1 份验证轮流做 K 次。在 Hadoop 上做交叉验证不用改代码只要预处理时按标签生成 K 组 train/test 文件循环提交 K 次作业即可。对课程设计来讲跑 5 折就足够说明效果了。混淆矩阵更能暴露问题。如果「体育」被大量误判成「娱乐」大概率是这两个类别的特征词重叠度高球员花边、节目综艺都出现人名。反过来如果某个类别精确率很高但召回率很低通常是训练样本太少或平滑参数过重。把这些数字打印出来挂在文档说明里比贴一句「准确率 91%」有说服力得多。我自己做这类项目踩过最大的坑是「拿训练数据当测试数据评估」——模型在训练集上准确率 98%换一批真实数据直接掉到 74%。后来养成的习惯是训练脚本和评估脚本分离评估时必须用预处理脚本的另一份输出绝不在同一个目录里混着。希望这个习惯和上面的调参经验能帮到你。本文还有配套的精品资源点击获取