1. 为什么逼自己啃一遍MapReduce源码而不是继续当黑盒用户很多写了好几年Java的大数据工程师对MapReduce的认知还停在“写个Mapper、写个Reducer、丢给集群跑”这一步。能用但一旦遇到数据倾斜、任务卡死、OOM、输出结果对不上就只能靠猜靠restart靠玄学调参数。你问我怎么知道的因为我就是这么过来的。后来我决定把MapReduce源码从头到尾过一遍前后花了大概三周每天两到三小时。说句实话这个过程比我想象的痛苦尤其是Hadoop那一堆类之间的继承关系绕得人头晕。但读完的收获也远超预期我不光搞清楚了Job从提交到落盘的完整链路还顺手弄明白了很多之前“调参靠蒙”的参数背后真正的作用机制。这篇文章我不打算给你做源码逐行的翻译机那个自己去看就行。我想跟你聊的是——源码该怎么读、重点看哪几个类、哪些代码值得反复琢磨、哪些地方看一眼就跳过。顺便把我踩过的坑和总结出来的学习方法一块儿分享出来。这篇内容适合谁适合已经能熟练编写MapReduce程序、但始终觉得“差点意思”的工程师也适合正在准备面试、需要把MapReduce机制讲透的求职者。如果你刚接触MapReduce不久建议先跑通几个实例再看这篇文章不然某些部分可能会觉得干。2. 开始之前先看清学MapReduce源码到底在学什么2.1 源码学习的核心主线一条作业的生命周期面对Hadoop这么庞大的代码库毫无章法地乱翻是大忌。我一开始就是没有规划直接在IDE里全局搜类名今天翻两行明天看三页一周下来脑子里全是浆糊什么都记不住。后来我做了一次彻底的思路调整。我意识到MapReduce源码学习的核心主线其实就一条一条作业从提交到结束全流程经历了什么。想清楚这一点整个阅读框架瞬间就清晰了代码不再是零散的类而是串成一条链路上的一个个节点。大致的主线是这样的客户端调用submit()把Job交给JobSubmitter。JobSubmitter向ResourceManager申请一个Application ID把作业的资源文件jar包、配置、分片信息写入HDFS。ResourceManager收到后调度到某个NodeManager上启动一个MRAppMaster。MRAppMaster根据输入分片数量计算Map任务数向ResourceManager申请容器。MapTask在容器里执行先读输入分片经过Mapper处理后写入环形缓冲区溢写、分区、排序。ReduceTask分阶段拉取属于自己的分区数据合并排序后交给Reducer处理。最终结果写到HDFS或直接丢弃整个作业在MRAppMaster的监督下结束。这里面的每一个节点都对应若干核心类。你只要把这条链路在脑子里串起来后面读所有源码都有定位感。哪怕某个类看不太懂你也能知道它处在哪个环节、是干嘛用的。2.2 版本和环境的选择2.x优先别跟老代码较劲第二个需要提前想清楚的问题是读哪个版本的源码市面上很多分析文章还停留在Hadoop 1.x聊的是JobTracker和TaskTracker。如果你的工作环境用的是CDH 6.x或者Apache Hadoop 3.x直接读1.x源码会把你看懵因为Job的调度逻辑发生了很大变化。我个人建议直接读Apache Hadoop 2.10.x或3.3.x的源码。原因很简单第一目前主流发行版几乎都是基于2.x或3.x的YARN架构读了用得上第二网上针对2.x新架构的分析文章更多遇到不懂的地方方便找参考。我自己的阅读用的是hadoop-3.3.6这个分支JDK 8即可编译IDE用IDEA就可以直接打开整个工程。源码在哪GitHub上搜apache/hadoop直接拉hadoop-mapreduce-project/hadoop-mapreduce-client这个子模块就够了。如果只想看核心逻辑不用导入整个Hadoop工程因为里面有几百个子项目光索引就要转半天。我最初把整个仓库都导入过IDEA疯狂建立索引风扇嗡嗡响体验极差。2.3 先懂编程实例再碰源码从写对到看对的阶梯在读源码之前我强烈建议你先有一个“能跑通”的MapReduce程序作为参照物。学校实训或者练习里经常做的WordCount、日志清洗、词频统计都可以。为什么因为源码是抽象的而实例是具体的。你要知道WordCount里那二十几行代码背后对应的源码调用链就是你理解MapReduce的最佳锚点。你写了一个Mapper的map()方法源码里是谁在调用它调用之前做了什么准备调用之后数据去了哪里带着这些问题去看源码和漫无目的地翻源码效率完全是两个档次。所以我画了一条学习路径供你参考写一个简单Job并成功提交到本地模式或集群。给Mapper和Reducer的代码打上断点以Debug模式跑一次看调用栈。顺着调用栈一层层点进去看源码。这个Step 2是关键。很多人不会用Debug方式跑MapReduce觉得Debug只能用于普通Java程序。其实本地模式下MapReduce程序是可以直接Debug的叫setLocalMode或者本地运行。我第一次在Mapper.run()方法上打了一个断点看着程序停在那一行一行往下走的时候那种感觉真的和纯阅读完全不同——所有抽象概念瞬间被盘活了。3. MapTask的源码拆解从run方法出发摸清数据流动的秘密3.1 第一个必须精读的类MapTask整个MapReduce中我建议你第一个精读的类是org.apache.hadoop.mapred.MapTask。为什么是它因为Map阶段的所有核心逻辑都汇聚在这个类里它像是一个总装车间输入端接的是分片数据输出端接的是shuffle的数据准备。打开这个类你会看到两个字段特别显眼mapOutputCollector和sortPhase。前者负责收集Mapper的输出后者标志着排序阶段是否结束。MapTask的run()方法里有一段逻辑根据任务类型走不同的分支。我们需要关注的是runNewMapper()方法。runNewMapper()里面有个很重要的判断是否定义了MapRunner。这个MapRunner是Java的Class对象如果配置了就用自定义的Runner否则使用默认的MapRunner。默认情况下MapRunner.run()会做这样几件事取到输入分片的RecordReader。通过RecordReader.nextKeyValue()持续读取键值对。每读取一对调用一次mapper.run(mapContext)。而mapper.run()的内部才真正调用你写的map()方法。我把这个调用关系简化成一句话input - (key, value) - Mapper.map()。别看这几个调用语义简单整条链路涉及的实现细节特别多。比如nextKeyValue()是怎么处理长记录跨split边界的LineRecordReader瞬间被分成几类TextInputFormat又是如何计算分片的这些你在调试一次之后都会恍然大悟比如一个文件只有4个字节为什么Map数可能是1而不是2——因为文件长度小于mapreduce.input.fileinputformat.split.maxsize时根本不会分片。3.2 环形缓冲区MapTask里最精妙的设计如果你读MapTask源码只想记住一个设计我会选环形缓冲区。它是MapOutputBuffer这个内部类实现的。请看这个类的字段和结构private byte[] kvbuffer; // 实际数据存储区 private byte[] kvmeta; // 元数据存储区 private int kvstart, kvend, kvindex; private int equator, bufindex, bufmark;注意Hadoop的环形缓冲区由数据区和元数据区组成两个区共用同一个大的字节数组各自从两头往中间增长。这是Hadoop比较厉害的设计节省内存、无缝支持溢写。数据区和元数据区的分界由equator这个指针控制。数据从数组左边的equator往右写元数据从数组右边的equator向左写两个方向同时扩展直到两者的index相遇触发溢写。你怎么判断该溢写了源码里的逻辑是bufindex超过kvbuffer.length或者将要碰撞kvmeta时。这个还在细讲环形缓冲区设计细节太占篇幅但有两个点值得你重点关注第一序列化。Mapper输出的key/value会被序列化成字节数组写进kvbuffer。你写的Writable对象最后都变成了字节。第二分区和排序。每一条记录写在kvbuffer里的同时还会在kvmeta里追加16字节的元数据包含value偏移量、key长度、value长度和分区号。这些元数据在后面做分区、排序、溢写时起到决定性作用。我把MapOutputBuffer的collect核心代码简化梳理一下synchronized void collect(K key, V value, int partition) throws IOException { // 检查是否需要溢写即空间不足以存下新记录时 if (bufindex intLimit || kvindex kvstart) { spill(); } // 序列化key和value keySerializer.serialize(key); valueSerializer.serialize(value); // 在kvmeta中记录这条记录的元数据 int kvmetaIndex ((kvindex - 1) * 4) kvmeta.length - 16; // 写入长度、偏移、分区号等信息 }这段代码的逻辑本身不复杂但也有个很妙的点溢写并非写到磁盘就完事而是把当前缓冲区的数据交给spill()方法。spill()内部会先做快排把数据按partition和key排序然后写溢写文件。读到这里我当年恍然大悟Map端这看似轻描淡写的排序其实是MapReduce性能的核心开销之一也是为什么我们经常看到Map阶段有spill records指标。3.3 Sort和Spill到底谁先谁后源码会给最准确的答案有一个很常见的误区很多人以为Map阶段是先全部处理完再一次性地排序溢写其实不是。spill()触发时机是环形缓冲区快满的时候而不是数据处理完之后。而且每一次spill都涉及一次全量的快排这是个大成本。spill()之后这几个步骤是确定的根据kvmeta里保存的partition信息把数据按分区切分。每个分区内部按key做一次排序快速排序。如果设置了Combiner则可能做一次局部合并。写入一个溢写文件spill文件。当MapTask处理完所有输入后还有一个mergeParts()的收尾过程把多个溢写文件合并成一个最终的输出文件。合并时依然要分配合适的buffer避免OOM。这块代码在三段式框架里相对有难度建议你重点看MergeQueue的实现。这里我特别想提醒你注意一个细节排序是Map阶段默认就有的不是Reduce阶段才发生的。而且这个排序的算法在不同阶段是不一样的。Map端溢写用的是快排Reduce端合并用的是堆排序两者的语义目标也不同。你如果能在面试里把这个差异讲清楚面试官会觉得你真读了源码不是光背概念。4. Shuffle是灵魂MapReduce源码中最值得反复咀嚼的一段4.1 理解Shuffle之前先搞懂Partitioner和Combiner的调用时机Shuffle是MapReduce里最容易被问倒、也最值得深挖的一环。我这里说的Shuffle指的是从Map端产生输出开始到Reduce端拿到输入为止的全过程。先看Partitioner。在MapTask的MapOutputBuffer.collect()被调用时有一行代码决定了你的键值对属于哪个分区int partition partitioner.getPartition(key, value, numPartitions);Partitioner默认实现是HashPartitioner源码也就几行public int getPartition(K key, V value, int numPartitions) { return (key.hashCode() Integer.MAX_VALUE) % numPartitions; }这里有个细节为什么要 Integer.MAX_VALUE因为Java的hashCode()可能返回负数对负数取模会得到负分区号导致分区无效。先过滤符号位保证分区号非负。这个细节你在OJ里写自定义Partitioner时一定要注意我自己就踩过坑自认为写了一个完美的自定义分法结果分区数对不上最后发现是负hashCode引起的。Combiner呢它本质上是一个运行在Map端的Reducer源码位置在MapTask的collect相关的配置文件读取处。Combiner只有在mapreduce.map.combine.minspills的最小溢写数达到了之后才可能触发默认值是3。也就是说溢写文件数量至少达到3个Combiner才参与合并。你把minspills调成1的时候每次spill都会执行Combiner本地聚合效果更明显但代价是CPU占用上涨。在生产排错时这个参数经常是被忽略的调节杠杆。4.2 从Map端到Reduce端的完整链路Map端的输出最终落成一个数据文件file.out和对应的索引文件file.out.index。ReduceTask启动时会调用类似shuffle的入口来抓取属于自己的那部分数据。在具体源码里这段逻辑由org.apache.hadoop.mapreduce.task.ReduceContextImpl和内部的Shuffle类完成但核心的抓取逻辑涉及Fetcher它是ReduceTask里一个实现Runnable的线程。Fetcher要做的事是从MRAppMaster获取已完成的MapTask列表。与对应的NodeManager建立HTTP连接请求map输出数据。一边抓取一边将数据写入Reduce端的内存缓冲区。当缓冲区满到阈值或者Map输出数据总量过大时溢写至磁盘。我记得源码里有个ShuffleScheduler负责决定“接下来该抓谁的数据”。它会做黑名单管理——如果某个节点连续多次抓取失败就会把它临时加入黑名单换一个节点重试。这个容错逻辑千万不要错过因为它在真实生产环境里对作业稳定性的贡献非常大。4.3 副本抓取与合并排序Reduce端的“叠叠乐”逻辑Reduce端抓来的数据并不是有序的。不同MapTask产生的输出分区虽然各自内部有序但你从多个Map端拿回来后混杂在一起就是乱序的。所以Reduce端第一件事就是做多路归并排序。源码里对应的是org.apache.hadoop.mapred.ReduceTask内部的MergeManager其中核心是createKVIterator方法。它会把内存中的数据段、磁盘中溢写的数据段统一交给一个PriorityQueue来做归并最终产生一条全局有序的键值流。有序之后凡是相同key的value会被连续送到Reducer的reduce()方法里。这就是为什么你在写Reducer时可以用IterableVALUES拿到同一个key的所有value的原因——它实际上是一个迭代器内部指向的是一条由归并排序产生的数据流不是一个真实的集合。这也是很多初学者容易误解的地方真正落到reduce()方法里时这个Iterable可能在迭代过程中跨文件读取底层是不断从磁盘拉数据的。如果你在写Reducer时对Iterable做了多遍遍历性能会成倍下降。所以最佳实践是在reduce方法里只遍历一次把需要的数据抽到自定义结构中尽量不要再回头取。这里我给你总结一张Shuffle核心类的对照表方便阅读时定位阶段核心类主要职责Map端收集MapOutputBuffer缓存、分区、桶化数据Map端排序ExternalSorter / QuickSort对缓冲区内数据排序Map端溢写BspWriter / IfileWriter写spill文件Map端合并MergeParts多个spill文件合并为一个Reduce端抓取Fetcher / ShuffleScheduler动态抓取map输出Reduce端合并MergeManager / PriorityQueue多路归并排序数据读取KVIterator / RawKeyValueIterator向Reducer提供有序键值流5. ReduceTask和作业调度把源码读成一张完整的流程图5.1 ReduceTask的run方法藏着多少你没想到的细节ReduceTask的总体流程可以概括为四个阶段shuffle - sort - reduce - write。源码上ReduceTask.run()方法与MapTask类似也会根据新旧API选择不同入口而runNewReducer()是新一代API的执行方法。runNewReducer()中几个值得关注的点ShuffleRunner会把RawKeyValueIterator最终包装成一个ReduceContext。ReduceContext.nextKey()负责移动当前的key同时检查key是否发生变化。Reducer.run()循环调用reduce()方法处理完同一个key的所有values之后再取下一个key。有几处细节特别有意思。比如ReduceContext.nextKeyValue()内部其实非常讲究——它维护了currentKey、currentValue等字段还会通过nextKeyIsSame这个布尔变量判断当前key是否发生变化从而决定是否跳出循环。你看到这段代码之后就能明白为什么reduce方法里的Iterable是一次性的了它底层的cursor在推进完之后不会再自动返回。严格来说ReduceTask还有一个容易被忽略的阶段判断逻辑它需要从MRAppMaster获取一段“shuffle已结束”的通知。这个通知在分布式环境下通过ShuffleConsumerPlugin的close()方法触发。这意味着如果某个ReduceTask的shuffle阶段一直没能完成后面的reduce阶段根本不会启动。这个机制平时很难遇到问题但如果你在超大规模集群上作业hang住了这一环绝对是排查的重点区域。5.2 MRAppMasterMapReduce作业的“总导演”很多人读MapReduce源码时容易把目光全部聚焦在MapTask和ReduceTask上从而忽略了调度中枢MRAppMaster。但要真正理解整个作业的运行机制这个类一定要看。MRAppMaster启动后要做的事初始化Dispatcher事件分发器。注册各类事件处理器比如JobEvent、TaskEvent、JobHistoryEvent。根据输入分片元信息计算任务数。动态为Map和Reduce任务申请资源并下发任务启动命令。其中最有意思的设计是事件驱动模型。MRAppMaster内部维护了一套异步的事件循环不同的模块通过发送事件进行交互彼此之间没有强依赖调用。举个例子某个MapTask运行完成它给Dispatcher发送一个TaskAttemptEventMRAppMaster收到后更新任务状态再根据并发度决定是否调度下一个任务。这种模型的好处是扩展性好但也带来一个问题日志中经常只能看到事件流转的蛛丝马迹出现故障时直接看代码调用栈反而不好使。我建议你在读MRAppMaster时把TaskAttemptEvent、JobFinishedEvent这几个事件相关的类也顺带过一遍否则很多行为会看得一头雾水。5.3 Job提交的入口JobSubmitter干了哪些脏活累活最后再往前绕一步回到客户端。当我们调用Job.waitForCompletion(true)时实际执行流程先是submit()方法进入JobSubmitter.submitJobIntercepted()。这个类在提交作业时做了若干工作校验作业输出目录是否存在避免覆盖。将作业的jar包上传到HDFS。计算输入分片把分片信息写入job.split文件。把作业配置conf写入HDFS的job.xml。应用setupJob钩子允许开发者提交前做一些特殊处理。这里有一个经典面试题“MapReduce作业的split切片规则是什么”在JobSubmitter的writeSplits()方法里会调用InputFormat.getSplits()。FileInputFormat的默认分片逻辑是目标分片大小 max(minSize, min(maxSize, blockSize))。默认情况下blockSize往往是64MB或128MB所以分片基本等于一个块大小。源码里的computeSplitSize()方法仅有几行但决定了很多集群调优的方向。比如你有一个很大的压缩文件但由于压缩格式不支持切分如某些不支持splittable的压缩格式整个文件会变成一个不能split的输入分片导致单Map处理全部数据、负载极度不均。这种问题不看源码光看监控面板很难定位。6. 常见问题与排查技巧实录6.1 本地Debug模式跑不出结果多半是这些坑学习源码阶段很多人第一件事就是本地Debug但会遇到一些典型的坑。第一个坑没设置mapreduce.framework.namelocal导致程序还按YARN模式找集群直接报Connection refused。解决方式是在代码里加Configuration conf new Configuration(); conf.set(mapreduce.framework.name, local); conf.set(fs.defaultFS, file:///);第二个坑Debug模式下输入路径如果设在HDFS会因为本机没有HDFS环境而报FileSystem错误。建议直接把数据放在本地文件系统用本地路径跑通就开始打断点。第三个坑断点打的位置不对。很多人喜欢在map()方法第一行打断点当然有效但要真正观察源码链路建议断点下在这三个位置MapTask.runNewMapper()的入口。Mapper.run()的while (context.nextKeyValue())那一行。MapOutputBuffer.collect()方法。这里能看到整个环形缓冲区操作的初期状态。6.2 任务执行成功但结果不对从源码检查这3个环节作业能跑完但结果不符合预期的情况我遇到很多次。这种问题源码知识往往比业务排查更管用。按我的经验优先检查这三个环节第一Partitioner是否自定义。如果你写了自定义Partitioner但分区数和Reduce数不匹配很可能某些key被分到了不存在的分区导致数据丢失。源码里可以印证分区号大于numPartitions时getPartition直接抛出异常但前提是你的numPartitions设置合理。第二Combiner是否被错误使用。Combiner必须满足交换律和结合律否则Map端局部合并的结果会和全局合并不一致。源码层面的combiner.run()其实调用的就是Reducer方法它没有做任何正确性校验。也就是说你把一个不满足交换律的Combiner传进去框架不会报错只会静默地产生错误结果。这类问题真到线上就是事故级别的。第三output format的输出路径。如果Reduce输出目录已经存在作业会在提交流程的checkOutputSpecs()阶段直接失败。源码里这一段的校验逻辑缜密到连“目录是文件还是目录”都会检查。如果你遇到FileAlreadyExistsException别再折腾别的先清空输出目录。6.3 从“看热闹”到“看门道”高效阅读源码的三个方法最后分享几个让我受益最大的源码阅读方法。第一个是**“调用栈逆推法”**。遇到不懂的方法不要从入口找出口而是直接在Debug状态下查看这个方法被谁调用了。通过IDEA的Call Stack面板往上找一层调用者往往能快速定位到当前方法的设计意图。比如我当时看不懂MapOutputBuffer的adjustSpillIndex就是通过逆推发现它只是为了防止数据区与元数据区交叉的一个“边界修正”。第二个是**“问题驱动法”**。不要为了读而读从实际问题出发找源码。比如你遇到“Map阶段堆内存溢出”那就去查MapOutputBuffer的初始化参数io.sort.mb是怎样影响环形缓冲区大小的。源码里有一行int size conf.getInt(io.sort.mb, 100) * 1024 * 1024;一看到你就明白了这个参数决定的不只是“缓存大一点”而是整个溢写边界的起点。问题解决完这个类你也读得八九不离十了。第三个是**“画图辅助记忆法”**。MapReduce的调用关系极其庞大纯看代码记忆效果很差。我建议你在读每个类时一边读一边画调用关系图。画图不用什么专业工具在纸上或白板上画箭头就够了。一张好的调用链图胜过十篇笔记。比如我把MapTask的调用链画成“run - runNewMapper - MapRunner.run - Mapper.run - map()”这张图直到现在我面试时还能直接默写出来。7. 关于源码学习节奏的一点真实体会源码学习的最大门槛不是代码本身而是心态。我在前面说了自己最初三周几乎是“无效阅读”直到确定主线之后才走上正轨。如果你现在正被各种类名绕得头大我建议你先放下深挖的执念把整条Job生命周期跑通一遍哪怕很多内部细节暂时不懂先把大框架立住。我个人觉得比较合理的学习节奏是第一周只看JobSubmitter和MRAppMaster搞清楚作业怎么提交、任务怎么调度第二周攻MapTask和环形缓冲区第三周攻Shuffle和ReduceTask。每天不用贪多稳扎稳打研究一两个核心类就够了。读完之后你再去写MapReduce程序完全是对代码有掌控感的状态你写的map()方法不再是黑盒里的函数而是你亲手从源码里看过的调用链上的一环。还有一个加分项值得提读源码过程中你不经意积累的这些细节在面试时是非常自然的谈资。当面试官问你“Map端为什么要排序”的时候别人只能背概念你可以直接说“因为MapOutputBuffer溢写前会调QuickSort这是为了和Reduce端的多路归并衔接”这种答案的区分度是立竿见影的。
