我花了几天时间把 MapReduce 完整过了一遍,一边看源码一边手动搭建环境跑实例,算是把 Hadoop 分布式计算框架这条线彻底理清了。网上讲 MapReduce 的资料很多,但大多偏理论,真正从入门到能独立完成一个任务闭环的笔记不多。这篇分享就从我的角度出发,把你需要知道的原理、关键细节、实操步骤还有那些文档里不写但你必须知道的坑,一次讲清楚。
这篇笔记适合谁看?正在学大数据、准备相关岗位面试、或者工作中需要上手 Hadoop 生态但还没完全吃透 MapReduce 的开发者。我的目标是让你读完既能理解 MapReduce 为什么这样设计,又能照着步骤亲手跑通一个 WordCount 实例,甚至能自己排查作业失败、数据倾斜这类典型问题。
1. 项目概述:MapReduce 到底是什么
1.1 从一大段需求说起
如果你手头有 100 GB 的日志文件,想统计里面每个关键词出现的次数,你会怎么做?单机读一遍可能要几个小时,而且内存根本扛不住。这时候就需要把数据拆成很多份,分给很多台机器一起算,算完再把结果收回来合并。这个“分下去、并行算、收回来”的过程,就是分布式计算的本质。
MapReduce 就是这个思想的标准化实现。它不是一门编程语言,而是一套编程模型,跑在 Hadoop 集群之上,专门用来处理海量数据的批量计算。你只需要写两个函数,Map 负责“分”,Reduce 负责“合”,剩下的任务调度、节点通信、故障恢复全部由框架完成。
我第一次学的时候有个误区,以为 MapReduce 是一个需要手动开线程、管网络通信的工具,后来才发现它更像一套“约定”:你把数据想怎么处理、怎么聚合告诉框架,框架帮你把脏活累活全干了。
1.2 为什么到现在还要学 MapReduce
很多人会说,现在有 Spark、Flink 这些更快的计算引擎,MapReduce 是不是过时了?这个问题我在面试里也经常被问到。
我的看法是,MapReduce 仍然值得学,而且值得认真学。原因有三点。
第一个原因是底层逻辑的通用性。MapReduce 的“分而治之”“移动计算而非移动数据”这些思想,Spark 和 Flink 都继承了。你把 MapReduce 吃透,再学 Spark 会轻松很多,因为那些高频概念,比如 Shuffle、分区、合并,本质上是同一套东西,只是实现方式不同。
第二个原因是生态地位。Hive 早期默认的底层执行引擎就是 MapReduce,虽然现在 Hive 可以用 Tez、Spark 作为执行引擎,但很多老旧的集群和生产环境仍然直接跑着 MapReduce 作业。你不会,出了问题就束手无策。
第三个原因是面试。MapReduce 的相关问题几乎是所有大数据岗位面试的必考项。从 WordCount 原理到数据倾斜优化,我见过太多人挂在 Shuffle 细节上。这关过不去,后面聊 Spark、Flink 也会底气不足。
当年我自己学的时候,就是靠死磕 MapReduce 的性能优化,才逐步建立起对分布式计算的整体认知框架。很多 Spark 调优的思路,回头去看,其实都能从 MapReduce 里找到原型。
2. 核心设计思路:框架背后的关键决策
2.1 “分而治之”与“移动计算”
MapReduce 最核心的思想是“分而治之”。一个大任务被拆成无数个小任务,分发到不同节点上并行执行。但这里有一个关键细节值得多说几句:数据是分布在不同机器上的,处理逻辑分发到数据所在的位置去执行,而不是把海量数据搬到一台机器上处理。这就是“移动计算而非移动数据”。
想象一下你在一个大仓库里理货,货物分散在不同货架上。最蠢的办法是把所有货物搬到中间一张大桌子上慢慢分类,聪明的办法是让每个货架旁边的人都地完成初步整理,再把整理好的小批量结果拿到中间合并。云计算框架的核心逻辑与此完全一致。
这样做的好处很明显:省下了海量的网络传输时间。因为移动程序代码比移动数据要快得多,尤其当数据量达到 TB 级别时,这个优势会被无限放大。在实际集群中,MapReduce 会优先把计算任务调度到数据所在的节点上,这就是“数据本地性”优化,能极大减少网络开销。
在 Hadoop 的源码实现里,这个优化体现在任务调度时会尽量选择数据块所在的节点或机架。如果你的数据在节点 A,你的 Map 任务被分配到节点 A 上执行,就能直接读本地磁盘上的数据,速度快一个数量级。
2.2 为“不可靠”设计的容错机制
分布式环境下有个必须面对的现实:节点随时可能宕机,网络随时可能抖动。MapReduce 的设计思路,从一开始就假设硬件是不可靠的,然后通过软件机制来保证整个作业能完成。
具体怎么做?答案就是“任务重试”和“推测执行”。
MapReduce 会把任务运行状态实时汇报给 AppMaster,一旦某个任务长时间没响应或执行失败,AppMaster 会把这个任务重新调度到另一个节点上执行。如果是同一个任务反复失败,框架不会无限重试,而是会标记失败并结束整个作业,同时输出错误日志供排查。
推测执行更有意思。如果集群中有个节点因为硬件老化等问题跑得特别慢,拖慢了整个作业的进度,MapReduce 会在其他节点上启动一个同样的任务作为备份,谁先跑完就采用谁的结果,另一个直接杀掉。这个机制在集群繁忙时非常有用,但也带来一个副作用:集群资源会额外占用,所以在新版 Hadoop 中已经改为默认开启但策略更加灵活。
我觉得这套容错设计最大的价值在于:你写 MapReduce 程序时,不需要关心节点故障,框架自动帮你处理。这种“把复杂留给自己,把简单留给用户”的理念,也是后续所有大数据框架的设计基调。
2.3 数据本地性优化:让计算去找数据
在上面提到“移动计算而非移动数据”时,我简单提了一句数据本地性。这里展开说一下。
HDFS 会把大文件切分成 128 MB 的块,每个块默认有三个副本,分布在不同的节点上。MapReduce 在调度 Map 任务时,会尝试把任务分配到数据块副本所在的节点上,这样 Map 任务就可以直接读本地磁盘上的数据,不需要通过网络拉取远程数据块。
数据本地性有多个层级,从高到低分别是:
- 节点本地:数据块就在任务所在节点上,速度最快。
- 机架本地:数据块在同一机架的另一台节点上,需要走一次机架内网络。
- 跨机架:数据块在远端机架的节点上,需要走核心网络,速度最慢。
在写代码时,DataNode 会通过心跳机制上报自己持有的数据块列表给 NameNode,当 AppMaster 要调度任务时,会结合这个列表尽量把任务放到数据所在节点上。不过这里有个值得注意的细节:如果某个节点数据量特别大,为了负载均衡,框架也不一定把所有任务都放在同一个数据节点上,而是会权衡数据本地性和负载均衡之间的关系。
所以在实际生产环境中,如果你发现一个作业的 Map 任务全部是本地读取,磁盘 IO 和网络 IO 都比较正常,那说明调度得不错;如果发现大量跨机架读取,就得考虑是否数据分布不均衡,或者集群本身物理拓扑比较特殊。
另外,我们经常说“任务在数据所在的节点上运行”,指的是 Map 任务的输入数据。Reduce 任务的输入数据本来就来自所有 Map 任务的输出,所以 Reduce 任务不存在数据本地性这个概念,它必须通过网络拉取属于自己分区的数据。
3. 核心细节解析与实操要点
3.1 Map 阶段:怎么做初步处理
Map 阶段是整个 MapReduce 作业的起点。InputFormat 会读取数据源,调用我们实现的 Mapper 类的 map 方法,对每一条输入记录(默认以行为单位)做处理,输出零到多条键值对。
这里面有一个容易踩的坑:Map 阶段的输入并不是直接把原始文件丢给 map 方法,而是需要 TextInputFormat 先对文件做一次切分,然后一行一行地迭代,每行的偏移量作为 key、该行内容作为 value。如果你之前没接触过 Hadoop,可能不理解为什么 key 是偏移量,其实这是 Hadoop 源码里 TextInputFormat 的默认行为,wordcount 计算不需要偏移量,所以通常只用 value。
Map 函数的输出会被写入本地磁盘,而不是直接通过网络发送给 Reduce。这是 MapReduce 的一个关键设计决策:Map 输出是中间结果,只有经过 Shuffle 之后,属于自己的分区数据才会被 Reduce 拉走。
在写 Mapper 时,有几点经验可以分享。
第一点是合理使用 context.write 方法而不是自己搞一个 List 收集再一次性输出。Hadoop 的 OutputCollector 会在一定量时自动溢写,自己收集反而容易内存溢出。
第二点是尽量在 Map 阶段做“提前过滤”和“提前合并”。比如你要统计一定时间段的数据,在 Map 阶段就过滤掉不需要的行,而不是全量输出给 Reduce,让 Reduce 去过滤。这样能大幅减少 Shuffle 的数据量,效果立竿见影。
第三点是在 Context 对象的 setup 和 cleanup 方法中做资源的初始化和释放,而不是在 map 方法中反复操作,这样能减少开销。setup 方法在整个 task 开始前执行一次,cleanup 在整个 task 结束后执行一次。
3.2 Shuffle 与排序:整个框架最复杂的部分
Shuffle 是 MapReduce 中最重要的概念,也是面试最爱考的内容。很多人对 Shuffle 的理解停留在“把数据从 Map 发到 Reduce”的层面,但细节远比想象中复杂。
在 Map 端,Map 函数的输出先被写入一个环形缓冲区(默认 100 MB,通过 mapreduce.task.io.sort.mb 配置),当缓冲区使用量达到阈值(默认 80%)时,后台线程开始将数据溢写到磁盘。在溢写之前,数据会根据 Partition 分区并按键排序。如果一个 Map 任务输出的数据太多,可能会发生多次溢写,最终在 Map 任务完全结束前,把这些溢写文件合并成一个大文件,同时生成一个索引文件供后续 Reduce 拉取数据时查找。
Reduce 端会启动若干个 copy 线程,主动去各个 Map 节点拉取属于自己的分区数据。拉取过来的数据同样会先放到内存中(shuffle 阶段的内存缓冲),不够时溢写磁盘,最后全部拉取完毕后,对这些数据做一次归并排序,把相同 key 的键值对聚集在一起,然后交给 Reduce 方法处理。
我在学 Shuffle 的时候,拿快递柜做了个类比:Map 端相当于每个快递驿站把包裹按片区(分区)分好、按地址排好序,Reduce 端相当于每个片区派件员从所有驿站把自己片区的包裹拉回来,再按地址整理好一次性派件。这样理解就特别直观。
关于 Shuffle 有几个常见面试题,我在这里顺手整理一下:
- 为什么要排序?因为排序后可以让相同 key 的数据连续排列,方便后续的 Reduce 直接顺序处理。对于 BinaryComparator 比较规则,Hadoop 默认使用自然排序。
- 为什么 Map 端要做一次本地归并?因为一个 Map 任务可能产生多次溢写文件,Reduce 每次拉取一个文件不现实,本地归并可以减少 Reduce 端拉取文件的数量,大幅提升效率。
- Combiner 是什么?它是在 Map 端做的“预 Reduce”,把 Map 输出中相同 key 的 value 先合并一次,减少传输到 Reduce 的数据量。但 Combininer 的输入和输出必须与 Reduce 的键值类型兼容,否则会产生类型错误。
- 如果不想用 Combiner,可以不配置,框架也能跑,只是效率会下降。
3.3 Reduce 阶段:最终结果如何产生
Reduce 方法接收的输入是排好序后的数据,框架会保证相同 key 的所有 value 会放在同一个列表中连续传给 reduce 方法。这个保证是整个编程模型的核心,也是分布式计算能够正确聚合结果的基础。
在实现 reduce 方法时,要特别注意 key 相同但 value 数量很大的情况。比如某个单词在 1 亿行中都出现,它的 value 迭代器可能非常长,如果一次性把所有 value 放进 List,极易内存溢出。正确做法是直接在迭代器循环中累加或做处理,不要保存到内存里。
Reduce 阶段也可以有多个 Reduce Task。默认情况下只有一个,如果你想提高聚合效率,可以手动指定多个,但要注意这会导致最终输出被分成多个文件。如何保证相同 key 全部进入同一个 Reduce,是由 Partitioner 决定的。默认的 HashPartitioner 对 key 做哈希取余,保证相同 key 进入同一个分区。当你输出结果文件数量超过预期时,很可能就是你改了 Reduce Task 数量和分区逻辑。
另外,Reduce 输出默认写到 HDFS 上,文件名为 part-r-00000 这样的形式。如果你需要自定义输出格式,可以通过 OutputFormat 扩展,比如要输出到数据库或自定义文件格式。
3.4 WordCount 实例详解
学习 MapReduce 最好的方式就是跑通 WordCount。虽然网上示例很多,但很多人只是照抄代码,没有真正理解每一步在做什么。我在这里完整写一遍,并附上关键注释。
首先是 Mapper 类:
public class WordCountMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private final Text word = new Text(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // key 是行偏移量,value 是这一行的内容 String[] words = value.toString().split("\\s+"); for (String w : words) { if (w.isEmpty()) { continue; } word.set(w); context.write(word, one); } } }然后是 Reducer 类:
public class WordCountReducer extends Reducer<Text, IntWritable, Text, IntWritable> { private final IntWritable result = new IntWritable(); @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; // values 迭代器中包含了 Map 阶段输出的所有相同 key 的 value for (IntWritable val : values) { sum += val.get(); } result.set(sum); context.write(key, result); } }最后是 main 方法,配置 Job 并提交:
public class WordCount { public static void main(String[] args) throws Exception { if (args.length != 2) { System.err.println("Usage: WordCount <input path> <output path>"); System.exit(-1); } Job job = Job.getInstance(new Configuration(), "word count"); job.setJarByClass(WordCount.class); job.setMapperClass(WordCountMapper.class); job.setCombinerClass(WordCountReducer.class); job.setReducerClass(WordCountReducer.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.setCombinerClass(WordCountReducer.class),因为这里的 Reducer 同时也是合适的 Combiner,可以提前对 Map 端的输出做一次求和,显著减少网络传输量。
但也正好说明我刚才提到的坑:并不是所有 Reduce 函数都适合做 Combiner。举一个反例:如果要计算一组数字的平均值,Map 端输出(key, 值),Reduce 端求平均值,那 Combine 时如果同样直接求平均值,多个子平均值再加权平均会得出错误结果,因为每个分组的大小可能不同。所以不要盲目复用 Reduce 作为 Combiner。
3.5 自定义数据类型与序列化
Hadoop 并没有使用 Java 自带的序列化机制,因为它太重量级,序列化后的字节数大,性能差。Hadoop 自己定义了一套序列化接口 Writable,实现了更高效的二进制协议。
如果你需要在 Map 和 Reduce 之间传递自定义对象,比如订单信息、用户行为日志,就需要定义自己的 Writable 类型。具体需要继承 Writable 接口,实现 write 和 readFields 两个方法,字段顺序必须完全一致。还有一个细节,如果你需要把这个对象作为 key 使用,想让它在 Shuffle 阶段能排序,还得继承 WritableComparable 接口并实现 compareTo 方法。
我强烈建议在学习阶段就手写一两个自定义 Writable 类型的例子,因为面试经常问,而且这在真实项目中几乎一定会用到。当你处理多维数据时,默认的 Text 类型太局限,自定义对象能帮你更好地组织数据。
4. 实操过程与核心环节实现
4.1 环境准备:本地模式与集群模式
在跑 MapReduce 作业之前,你需要先有一个可以运行的 Hadoop 环境。环境有两种选择:伪分布式模式和完全分布式模式。
如果你是学习阶段,资源和机器有限,我推荐用伪分布式模式。所谓伪分布式,就是在单台机器上同时启动 NameNode、DataNode、ResourceManager、NodeManager 等所有角色,模拟出一个迷你集群。优点是启动简单、调试方便,缺点是性能弱,无法体验真实集群的任务调度和数据本地性优化。
完全分布式则至少需要三台机器以上,一般一台做 NameNode 和 ResourceManager,两台做 DataNode 和 NodeManager。生产环境会更大,还会有主备 NameNode 的高可用架构,但那是进阶话题,这里不展开。
无论哪种模式,都需要先安装 JDK 并设置 JAVA_HOME 环境变量。Hadoop 3.x 版本要求 JDK 8 或 JDK 11,你在安装前先确认版本匹配。
Hadoop 官方下载页有时访问较慢,可以使用国内镜像源下载。下载解压后需要配置 Hadoop 目录下的 etc/hadoop/hadoop-env.sh,指定 JAVA_HOME 路径。注意这里的 JAVA_HOME 要写成绝对路径,不能直接用环境变量代指,否则启动可能出现问题。
接下来修改四个核心配置文件:
- core-site.xml:配置 fs.defaultFS 为 hdfs://localhost:9000 等。
- hdfs-site.xml:配置副本数为 1(伪分布式)或默认 3(生产集群)。
- mapred-site.xml:配置 mapreduce.framework.name 为 yarn。
- yarn-site.xml:配置 yarn.nodemanager.aux-services 为 mapreduce_shuffle。
配置完成后,需要格式化 NameNode:
hdfs namenode -format注意,这个命令只能执行一次。如果集群出了问题需要重新格式化,要先删除临时目录下的数据,否则会出现集群 ID 不一致的问题。当时我第一次搭伪分布式,反复格式化导致 NameNode 启动不了,排查了半天,就是因为在 /tmp 目录下残留了旧数据。
然后启动 HDFS 和 YARN:
start-dfs.sh start-yarn.sh用 jps 命令检查进程,如果看到 NameNode、DataNode、ResourceManager、NodeManager 这几个进程都在,说明启动成功。
4.2 编写与打包代码
环境准备好以后,就可以写代码了。建议直接用 Maven 管理项目依赖,在 pom.xml 中引入 Hadoop Client 依赖(具体版本号和你安装的 Hadoop 版本保持一致)。
依赖配置好以后,把上一节给的三个类写好,然后用 Maven 打包成 jar 包。执行 mvn clean package 即可。如果你开发环境没有配置 Maven,也可以直接使用 javac 编译,但依赖管理会麻烦一些,不建议。
打包好之后,把 jar 包上传到集群节点。上传方式可以使用 scp 命令。然后准备一份待统计的文本文件,上传到 HDFS。
hdfs dfs -mkdir -p /wordcount/input hdfs dfs -put test.txt /wordcount/input然后提交 MapReduce 作业:
hadoop jar wordcount-1.0.jar com.example.WordCount /wordcount/input /wordcount/output这里要注意:输出目录在提交前必须不存在。MapReduce 会自己创建输出目录,如果已经存在,作业直接报错。这是 Hadoop 对数据安全的一种保护,防止你把之前的输出覆盖掉却不自知。
执行后,命令行会输出作业进度,包括 Map 和 Reduce 的完成百分比。你还可以打开 YARN 的 Web 界面的 ResourceManager 查看作业详情,可以看到每个任务的启动时间、运行日志、数据本地性等信息,非常有用。
4.3 查看运行结果与日志
作业执行完毕后,可以用以下命令查看 HDFS 上的输出结果:
hdfs dfs -cat /wordcount/output/part-r-00000对于 WordCount 这种小任务,输出应该是每行一个单词加频率。初次跑通非常有成就感,我记得当时看到一个一个单词和频率列出来的时候,对 MapReduce 是怎么工作的才算真正有了感觉。
日志排查是更重要的技能。如果作业失败,可以找到对应任务日志的存放位置。YARN 会把任务日志收集到日志聚合目录,使用命令:
yarn logs -applicationId <application_id>可以查看应用程序的所有日志。这是排查问题的第一入口。如果日志显示某个任务反复失败,通常说明代码有 bug,或者数据有脏数据格式不符合解析逻辑。
我碰到最多的错误是 Input path does not exist,原因通常有两种:一是 HDFS 上的路径写错了,二是 hdfs dfs -put 之后没有保存到正确路径却误以为在本地。检查方法很简单,执行 hdfs dfs -ls /wordcount/input 看文件是否存在即可。
4.4 作业参数调优:让程序跑得更快
虽然学习阶段不一定需要调优,但了解一些常见参数能让你后续在生产环境少踩坑。
首先是内存参数。Map 和 Reduce 任务默认分配的内存可能偏小,如果处理的数据量更大,会频繁触发 GC,甚至直接 OOM。常见配置有 mapreduce.map.memory.mb 和 mapreduce.reduce.memory.mb,以及对应的 mapreduce.map.java.opts 和 mapreduce.reduce.java.opts。注意 java.opts 配置的是 JVM 堆内存,需要比 container 内存小一些,为 JVM 自身留出空间。
其次是并行度。Map Task 的数量通常由输入数据的分片数决定:128 MB 一个分片,1 GB 的文件会有 8 个 Map 任务。你可以通过调整 mapreduce.input.fileinputformat.split.maxsize 和 minsize 来改变分片大小,进而改变 Map 数量。Reduce Task 数量通过 mapreduce.job.reduces 配置,默认 1,如果你想加快处理但后续不需要单文件输出,可以适当调大。不过 Reduce 数量不是越大越好,因为每个 Reduce 都要拉取所有 Map 的输出数据,太多了反而增加 Shuffle 开销。
最后是 Shuffle 相关参数。mapreduce.task.io.sort.mb 控制 Map 端环形缓冲区大小,增大可以降低溢写频率。mapreduce.reduce.shuffle.parallelcopies 控制 Reduce 端拉取数据的并发数,增大能加快拉取,但要注意不能让所有 Reduce 同时把网络打满。
调优的思路永远是这样的:先定位瓶颈,再针对性调节。不要一上来就盲目调大所有参数,那样容易造成资源浪费甚至引入新问题。
5. 常见问题与排查技巧实录
5.1 数据倾斜:任务慢的根源
数据倾斜是分布式计算中最经典的问题。现象是:大部分 Reduce 任务早就跑完了,但有一个或几个 Reduce 任务卡了很久,整个作业迟迟不能结束。本质原因是数据分布不均匀,某个 key 的数据量远大于其他 key,导致处理它的 Reduce 任务负载过重。
在 MapReduce 中最常见的数据倾斜类型有三种:某个 key 的数据量巨大;某个 key 的 value 计算非常复杂;以及 Map 端的输出数据明显集中到某个分区中。
解决办法要先分析原因,再对症下药。
如果是某个 key 的数据量太大,可以考虑加盐:给这个 key 加上一些随机前缀,让它被分发到多个 Reduce 中处理。比如在 Map 阶段把 key 变成 key + 随机数,Reduce 阶段再摘掉前缀做聚合。这个方法对 wordcount 这种求和类问题很有效,但要注意,如果最终输出需要保持 key 的唯一性,还需要再经过一轮聚合去盐。
如果是某个 key 的计算本身很复杂,比如用户维度的计算依赖一个巨大的正则匹配,那可以考虑采用 Combiner 提前合并,减少该 key 进入 Reduce 的 value 个数。如果还不行,只能考虑拆逻辑,把复杂的计算拆成多阶段作业。
这里要提醒一点:加盐方案只适用于 value 之间相互独立、最终结果可以通过二次聚合得到的场景。如果你的业务需要看到某一个特定 key 的完整明细,那就不能简单加盐了,需要换思路,或用数据和业务特点做针对性调整。
5.2 小文件过多:任务调度的大敌
小文件过多是 Hadoop 文件系统非常忌讳的场景。因为 HDFS 的每个文件、目录、数据块在 NameNode 内存中都会对应一条元数据记录。如果文件数量特别多,NameNode 的内存就会不够用,进而影响整个集群的稳定性。
MapReduce 处理小文件时也有问题。默认情况下,每个文件或文件块对应一个 Map 任务,如果有一万个 1 KB 的小文件,就会启动一万个 Map 任务,任务调度和启动开销远大于实际计算开销,性能极差。
解决办法有几种:如果小文件已经存在,可以在作业前先用 SequenceFile 或 CombineFileInputFormat 合并小文件。CombineFileInputFormat 可以把多个小文件合并成一个分片交给一个 Map 任务处理,这样就从源头上减少了 Map 任务的数量。如果是流式写入导致的小文件,比如 Spark Streaming 写 HDFS,那么应该在写入端就做好合并控制,比如按时间窗口滚动成较大的输出文件。
我在实际项目中处理过最极端的一次:某个业务目录下有几十万个小文件,导致 NameNode 告警。最后我们写了个简单的 MapReduce 作业,把大量小文件合并成少量大文件,整个过程花了大约十几分钟,NameNode 的压力才降下来。
这个问题的根源往往是上游写入策略不合理,所以在设计阶段就要考虑好文件大小管理,而不是等出问题再救火。
5.3 内存溢出与 GC 频繁
Map 或 Reduce 任务报 java.lang.OutOfMemoryError 是常见的失败原因,可能发生在堆内,也可能发生在堆外。
堆内存溢出通常是 JVM 堆大小不够,可以通过调大 mapreduce.map.java.opts 或 mapreduce.reduce.java.opts 来解决,注意这个值应该比容器内存 mapreduce.map.memory.mb 略小。举个例子,如果你把 mapreduce.map.memory.mb 设为 2048 MB,java.opts 可以设为 -Xmx1536m,留出几百 MB 给非堆内存和 JVM 自身使用。
堆外内存溢出,最典型的是 Map 端溢写时缓冲区不够。你可以通过调大 mapreduce.task.io.sort.mb,同时适当调低溢写阈值比率 mapreduce.map.sort.spill.percent,让溢写发生得更频繁但每次量更小。
GC 频繁则通常是对象创建过于密集,比如在 Map 方法里直接 new Text 对象,而不是重用对象。我在 WordCount 例子里用了 private final Text word,就是这个原因。一个 task 要处理几百万行数据,每行都 new 两个大对象,GC 压力会非常大。这种代码层面优化,效果往往比调 JVM 参数更明显。
还有一点容易被忽略:Reduce 阶段的归并排序需要占用内存,如果同时有数量非常大的 Map 输出文件要合并,Reduce 堆内存设太小会频繁触发多次溢写,导致性能断崖式下降。
5.4 任务反复失败:从日志定位根因
作业失败是最让人头疼的问题,但也是最能锻炼排查能力的问题。遇到失败,第一步永远是看日志,不要瞎猜。YARN 的任务日志提供了从用户代码到系统级的所有信息,往往能在 Caused by 部分看到真正的异常。
常见的有这样几类:
第一类是代码异常,比如空指针、解析异常。这类最容易定位,日志里会直接打印异常堆栈,往往指向你代码的某一行。
第二类是类型不匹配,比如 Map 输出的 key 类型和 Reduce 输入的 key 类型不一致。这种情况作业在提交时会报错,提示类型不匹配。记得 Map 输出的类型必须和 Reduce 输入的 key/value 类型严格一致。
第三类是文件锁冲突,比如在一个集群里同时提交了多个写同一个 HDFS 路径的作业,报了 FileAlreadyExistsException。这种是路径冲突问题,换输出路径即可。
第四类是任务被 kill 但没有明显的异常。这种情况通常是内存超限,NodeManager 主动杀掉了容器,日志中会包含物理内存超限的提示。此时需要调大内存参数。
排查时的思路也分享下:先把日志从 ApplicationId 到任务级别一层层点进去,定位是哪个 task 失败,再看是失败在哪一个阶段,是 map 还是 reduce,再切换日志类型看 stderr 和 syslog。不要只盯着 stderr,很多框架自身的警告和调试信息在 syslog 里。
5.5 常见问题速查表
我把实践中最常遇到的问题整理成一张速查表,方便快速查阅。
| 现象 | 可能原因 | 解决方案 |
|---|---|---|
| 作业提交时报输出目录已存在 | 输出路径被占用 | 换输出路径,或清理旧路径 |
| Input path does not exist | HDFS 路径错误 | 用 hdfs dfs -ls 确认路径 |
| Map 任务大量失败 | 输入数据有脏格式 | 查看日志定位反序列化异常 |
| 个别 Reduce 任务特别慢 | 数据倾斜 | 加盐分散或使用 Combiner |
| 任务被 kill,日志提示内存超限 | 容器内存不足 | 调大 memory.mb 和 java.opts |
| GC 频繁,任务执行缓慢 | 代码创建过多对象 | 重用对象,避免频繁 new |
| 启动集群时 NameNode 起不来 | 格式化残留旧数据 | 清理临时目录后重新格式化 |
| 输出文件数量过多 | Reduce 数量过多或分区逻辑异常 | 调整 reduce 数量或自定义 Partitioner |
这张表本质上是排查思路的浓缩版,建议你先照着思路走,再对照表查找。
另一个实际经验是:调优前一定要有基线数据。跑作业前先记录默认参数下任务的耗时、资源占用、数据量级,然后每次只调整一个参数,再对比效果。我见过太多人一次调十几个参数,最后出了问题根本不知道是哪一项导致的。一次只调一个参数,是分布式系统调优的黄金法则。
还是那句话,这些问题的根源都脱离不了 MapReduce 本身的机制。你把分片、Shuffle、排序、分区这些核心机制搞明白了,遇到问题时自然能快速定位到具体环节,而不是病急乱投医。
我在学习 MapReduce 的过程中最大的体会是:其实它并不复杂,复杂的是你还没搞清楚它的执行流程就急着去写代码、跑任务。先把一张图装进脑子——数据怎么切分、Map 怎么处理、Shuffle 怎么搬运、Reduce 怎么聚合,每一步会发生什么,再去动手做,你会发现一切都顺理成章。
最后再分享一个小技巧。如果你身边有多个学习伙伴,可以尝试互相出题,比如让对方描述一个场景,你来设计 Mapper、Reducer、Combiner 和 Partitioner 该怎么搭配。这种“脑内推演”比闷头写代码更能加深理解,面试的时候也更容易在对话中展现出真实的掌握程度。MapReduce 只是大数据计算的第一块敲门砖,把它吃透了,后面的路会好走很多。