熟悉 Hadoop 中 MapReduce 原理的人都知道,整个执行流程里至少会涉及三次排序,分别是溢写快速排序、溢写归并排序以及 Reduce 拉取后的归并排序,而且这些排序默认开启,也就是系统自带的天然排序。那么,MapReduce 为什么要这样设计?先给出结论:核心目的是让整体执行更稳定,同时让输出结果满足大多数业务场景。前者体现在采用的是 sortShuffle 而不是 hashShuffle,后者则体现在中间结果的预处理能力上。要知道,经过排序后的数据在后续处理时会方便很多,比如 Reduce 端拉取数据时,排序结果就像索引一样,能够明显提升归并与计算效率。

在分析这种设计原因之前,先完整理解一下 MapReduce 的执行过程。在map阶段,系统会根据预先定义好的partition规则对数据进行分区。map任务首先把输出结果写入缓存,当缓存数据达到阈值后,就会将结果spill到磁盘。每发生一次spill,磁盘上就会生成一个对应的spill文件,因此一个 map task 可能会产生多个spill文件,并且在每次spill时,都会先对key进行排序。接着进入shuffle阶段,当map输出全部写完后,map 端还需要执行一次merge操作,按照partition以及每个partition内部的key进行归并排序(即合并+排序),此时每个partition内的数据已经按key整体有序。随后开始第二次merge,这次发生在reduce端。在这个过程中,数据可能同时存在于内存和磁盘中。严格来说,这个阶段的merge并不只是单纯排序,而是将多个局部有序的大文件继续进行合并+排序,最终生成一个更大的有序文件并完成整体排序流程。理解完这个过程后,你可能会觉得:如果自己实现一个MapReduce框架,似乎直接用HashMap输出 map 内容也可以。但在大数据处理场景下,真实系统设计往往需要优先考虑稳定性、扩展性和后续归并效率。
2.1 MapTask运行机制详解
整个流程图如下:

详细步骤:
首先,读取数据的组件
InputFormat(默认是TextInputFormat)会通过getSplits方法对输入目录中的文件进行逻辑切片规划,得到多个splits。有多少个split,通常就会对应启动多少个MapTask。默认情况下,split与block的关系是一对一。输入文件被切分为
splits之后,由RecordReader对象(默认是LineRecordReader)负责读取数据,以n作为分隔符,每次读取一行内容,并返回。其中,Key表示每一行首字符的偏移量,value表示这一行的文本内容。读取
split并返回之后,数据会进入用户自定义并继承的Mapper类中,执行用户重写的map函数。RecordReader每读取一行,这里就会调用一次。map逻辑执行完成后,会通过context.write把每条结果进行collect收集。在collect阶段,首先会进行分区处理,默认使用HashPartitioner。MapReduce提供了Partitioner接口,它的作用是根据key或value以及reduce任务数量,决定当前这对输出数据最终交给哪个reduce task处理。默认实现是先对key hash,再对reduce task数量取模。这样的默认分配方式主要是为了尽量均衡各个reduce的处理压力。如果用户对Partitioner有更明确的业务需求,也可以自定义实现并配置到job中。接下来,数据会被写入内存,这块内存区域被称为环形缓冲区。缓冲区的作用是批量收集
map输出结果,减少频繁磁盘IO带来的性能开销。我们的key/value对以及Partition结果都会被写入这个缓冲区。当然,在写入之前,key和value都会先被序列化成字节数组环形缓冲区本质上可以理解为一个数组,里面保存着
key、value的序列化数据,以及key、value对应的元数据信息,例如partition、key起始位置、value起始位置和value长度等。所谓环形结构,本质上是一种抽象设计。缓冲区的大小是有限制的,默认值为
100MB。当map task输出结果很多时,就有可能把内存撑满,因此需要在一定条件下把缓冲区中的数据临时写入磁盘,然后再次利用这块内存。这个从内存写入磁盘的过程被称为Spill,中文通常翻译为溢写。溢写由独立线程完成,不会影响向缓冲区继续写入 map 结果的线程。为了避免溢写线程启动时阻塞map输出,整个缓冲区设计了一个溢写比例spill.percent。该比例默认是0.8,也就是说当缓冲区数据达到阈值(buffer size * spillpercent = 100MB * 0.8 = 80MB)时,溢写线程就会启动,锁定这80MB内存执行溢写过程,而Maptask的输出结果仍然可以继续写入剩余的20MB内存,两者互不影响。
当溢写线程启动之后,需要对这
80MB空间中的key进行排序(Sort)。排序是MapReduce模型中的默认行为,也是 Shuffle 机制的重要基础。如果在
job中配置过Combiner,那么此时就会发挥作用。它会把相同key的key/value对中的value做局部聚合,例如求和,从而减少溢写到磁盘的数据量。Combiner能够优化MapReduce的中间结果,因此在整个模型中可能会被多次调用。那么,哪些场景适合使用
Combiner呢?从这里可以看出,Combiner的输出会作为Reducer的输入,因此它绝不能改变最终计算结果。Combiner只适用于Reduce输入key/value类型与输出key/value类型完全一致,并且不会影响最终结果的场景,比如累加、求最大值等。Combiner的使用一定要谨慎,用得好可以显著提升job执行效率,用得不当则可能影响reduce最终结果。
合并溢写文件:每次溢写都会在磁盘上生成一个临时文件(写入前会判断是否启用
combiner)。如果map输出结果很大,多次溢写就会导致磁盘上存在多个临时文件。当整个数据处理结束后,系统会对这些磁盘临时文件进行merge合并,因为最终输出文件只有一个。合并完成后再写入磁盘,并为该文件生成一个索引文件,用来记录每个reduce对应数据的偏移量。
2.2 ReduceTask运行机制详解

Reduce大致可以分为copy、sort、reduce三个阶段,重点通常集中在前两个阶段。copy阶段包含一个eventFetcher,用于获取已经完成的map任务列表,再由 Fetcher 线程去copy数据。在这个过程中会启动两个merge线程,分别是inMemoryMerger和onDiskMerger,前者负责将内存中的数据merge到磁盘,后者负责对磁盘中的数据继续进行merge。当数据copy完成后,copy阶段也就结束了,随后进入sort阶段。sort阶段的核心工作是执行finalMerge操作,本质上是完成最终归并排序。该阶段结束之后,就进入reduce阶段,调用用户自定义的reduce函数进行处理。详细步骤
2.2.1 Copy阶段
这一阶段可以理解为简单的数据拉取过程。Reduce进程会启动若干数据copy线程(Fetcher),通过HTTP方式向各个maptask请求属于自己的输出文件。
2.2.2 Merge阶段
Merge阶段。这里的merge和map端的merge动作类似,只不过数组中存放的是从不同map端copy过来的数据。Copy过来的数据会先放入内存缓冲区中,这里的缓冲区大小比map端更加灵活。merge通常有三种形式:内存到内存、内存到磁盘、磁盘到磁盘。默认情况下,第一种形式并不会启用。当内存中的数据量达到一定阈值时,就会启动内存到磁盘的merge。这与map端类似,本质上也是一次溢写过程。在这个过程中,如果设置了Combiner,同样也会被调用,随后在磁盘上生成多个溢写文件。第二种 merge 方式会一直持续运行,直到不再有新的map端数据到来,然后才会启动第三种磁盘到磁盘的merge方式,生成最终文件。
2.2.3 合并排序
当分散的数据被合并成一个大的数据集之后,还会继续对合并后的结果进行排序。随后,对排序后的键值对调用reduce方法。对于键相等的键值对,只会调用一次reduce方法,每次调用可能产生零个或多个新的键值对,最后再把这些输出结果写入到HDFS文件中。
从MapReduce的完整执行过程回头来看,就能更清楚地理解为什么一定要排序,以及为什么在Shuffle阶段采用SortShuffle。从架构设计角度来说,MapTask 和 ReduceTask 本质上是两个完全不同、运行在 Yarn 上的进程,它们之间主要通过内存或磁盘进行数据交互。为了降低程序之间的耦合度,并更好地支持失败重试、任务恢复等机制,就不能像 Kafka 那样让生产者和消费者存在强依赖或阻塞关系,否则一旦链路受阻,就可能影响整个集群的吞吐与稳定性。MapReduce 面向的是海量数据处理,因此系统需要尽量让 Map 端已经完成的数据及时落盘,同时又要兼顾整体执行效率。所以在 map 结束时,交给 reduce 的不仅是排好序的数据,还有一份索引文件。这样虽然会额外消耗一定 CPU 资源,但对于已经落盘的数据来说,Reduce 端在拉取、归并和处理时会更快、更稳定。换句话说,只要 Map 端顺利完成,理论上即使中途停机,也可以继续执行 ReduceTask 来完成整个任务。至于为什么不采用 HashShuffle,原因也很直接:在大数据场景下,HashShuffle 对内存的占用通常更高,容易导致内存压力过大甚至溢出,从而影响集群计算稳定性。因此,无论是从 Hadoop 优化、MapReduce 原理,还是从 Shuffle 排序机制来看,SortShuffle 都是更适合大规模分布式计算的设计选择。大数据开发,更多关注查看个人资料
