夜雨聆风学习资料网

ARTICLE · 1090638

Hadoop 源码学习手册 06 · MapReduce:分布式计算框架|口袋通识馆

Hadoop 源码学习手册 06 · MapReduce:分布式计算框架|口袋通识馆

口袋通识馆 · 图解通识手册连载

《Hadoop 源码学习手册》第 06 章 · MapReduce:分布式计算框架

Map 分而治之 · Shuffle 拉取归并 · Reduce 聚合输出

  学习地图 

• 本模块讲什么:MapReduce 怎么把一个大计算拆成 Map + Shuffle + Reduce 三阶段

• 核心难点:Map 端的环形缓冲区(MapOutputBuffer)—— 用一个字节数组同时存数据和元数据

• 和上下游关系:运行在 YARN 上(MRAppMaster 是个 AM),输入输出走 HDFS

• 深入:Task/TaskAttempt 状态机 · RMContainerAllocator 资源申请 · Uber 小作业模式

关键类:MapTask · ReduceTask · MapOutputBuffer · MergeManagerImpl · TaskImpl · TaskAttemptImpl · RMContainerAllocator · MRAppMaster

比喻MapReduce 就像整理一堆乱糟糟的问卷调查。Map 阶段:每个人(Mapper)拿一摞问卷,把每张按问题分类记成小卡片(key=问题,value=答案)。Shuffle 阶段:把所有「同一个问题」的卡片收集到同一张桌子(Reducer)。Reduce 阶段:每张桌子的人把卡片按问题汇总统计(如数有多少人选 A)。精髓是「移动计算到数据」——不是把问卷搬到统计员那里,而是派统计员去问卷堆所在地。

三阶段范式

MapReduce 把分布式计算抽象成三个阶段,开发者只需写 Map 和 Reduce 两个函数。

 三阶段各干什么 

• Map:分而治之。每条输入记录独立处理,输出 (key, value) 对。Map 之间互不通信,天然并行。

• Shuffle:按 key 分区排序,把相同 key 的数据通过网络拉到同一个 Reducer。这是 MapReduce 的「网络瓶颈」所在。

• Reduce:归并。对每组相同 key 的数据做聚合(求和、计数、取最大值等),输出最终结果。

核心思想:「移动计算到数据」。与其把 TB 级数据拉到计算节点,不如把计算任务调度到数据所在的 DataNode 上。MapReduce 通过 YARN 的「数据本地性」调度实现这一点——优先把 Map task 分配到存有该数据块的节点。 

图:map → shuffle → reduce 数据流全景——中间结果走本地磁盘,最终结果才上 HDFS


MapOutputBuffer:环形缓冲区

这是 Map 端最精妙的设计——用一个字节数组同时存储数据和元数据,两区相向增长。

 内存布局 

单个字节数组 kvbuffer 被当成两区使用:左边是元数据区(向后增长),右边是数据区(向前增长),中间是空隙。两区相遇就触发溢出写。

// 内存布局示意(单个字节数组) // [kvmeta ← 元数据区(向后增长)] [......空隙......] [kvbuffer ← 数据区(向前增长)] //        ^                                    ^ //   metaStart                             dataStart(equator) // // 每条记录的元数据 (4个int = 16字节): //   kvmeta[offset+0] = VALSTART  (值的起始位置) //   kvmeta[offset+1] = KEYSTART  (键的起始位置) //   kvmeta[offset+2] = PARTITION (分区号,给哪个Reducer) //   kvmeta[offset+3] = VALLEN    (值的长度)
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/MapTask.java 
拆解逐字解读:kvmeta 是元数据区,每条记录占 4 个 int(16 字节),分别记录值的起始、键的起始、分区号、值长度。kvbuffer 是数据区,存真正的 key/value 字节。两区相向增长是为了共用一个数组避免分别分配——元数据从左往右长,数据从右往左长,中间空隙是可用空间。当空隙小于阈值(默认 80% 满)就触发 spill。为什么用元数据而不是直接存:排序时只需交换 16 字节的元数据指针(含分区号、起始位置),不用搬动整条数据,排序极快。
 Spill(溢出写)流程 
// 1. 缓冲区使用达阈值 (默认 80%, mapreduce.map.sort.spill.percent) // 2. 加锁,启动 spill 线程 // 3. spill 线程: 分区排序 → 合并键值对 → (可选 Combiner) → 写入本地磁盘 // 4. Map 线程继续写另一半缓冲区 // 5. 所有 map 完成后,合并所有 spill 文件为一个有序文件 // // 重要: 数据写的是本地磁盘,不是 HDFS!只有最终结果才写 HDFS。 // Combiner 最少 spill 次数: minSpillsForCombine = 3 //   (少于 3 个 spill 文件不触发 Combiner)
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-core/src/main/java/org/apache/hadoop/mapred/MapTask.java 
拆解逐字解读:spill 触发后,spill 线程先把缓冲区按分区号排序(同一分区内再按 key 排序),然后可选地跑 Combiner(预聚合,减少输出量),最后写到本地磁盘(不是 HDFS)。为什么写本地磁盘:spill 是临时中间结果,Map 结束后会被合并;写 HDFS 太慢且没必要。为什么 minSpillsForCombine=3:合并时如果 spill 文件太少(<3),跑 Combiner 的开销可能大于收益,所以跳过。
易错易错点:很多人以为 Map 输出直接发给 Reduce。实际上 Map 输出先写本地磁盘(spill),Reduce 通过 HTTP 主动拉取(Shuffle)。Map 和 Reduce 之间隔着「Shuffle + Sort」阶段。

图:MapOutputBuffer 环形缓冲区——元数据区与数据区相向增长,空隙不足 80% 即触发溢出写


Shuffle:Reduce 端拉取与归并

ReduceTask 分三阶段:拉数据(SHUFFLE)→ 排序合并(SORT)→ 调用 Reducer(REDUCE)。

 关键类 
类
角色
所在位置
ShuffleHandler
运行在 NM 上的 HTTP 服务器,提供 Map 输出数据给 Reduce 拉取
NM 端
ShuffleConsumerPlugin
Reduce 端拉取插件,抽象出拉取逻辑
Reduce 端
MergeManagerImpl
合并管理器:小数据内存合并,大数据溢出磁盘 + 外部合并
Reduce 端
Fetcher
拉取线程,从 Map 任务的 HTTP 端口拉属于自己的分区
Reduce 端
  调用链 · Reduce 端 Shuffle 三阶段 

• ReduceTask.run() 启动,进入 SHUFFLE 阶段           

为什么先 Shuffle:Reduce 的输入是所有 Map 的输出,必须先拉齐才能处理。

• ShuffleSchedulerImpl + Fetcher 线程池从 Map 任务的 HTTP 端口拉取属于自己的分区数据           

为什么用 HTTP 拉取而非推送:Map 众多且完成时间不一,Reduce 主动拉取能按需控制节奏,且 HTTP 走标准端口易穿越防火墙。

• MergeManagerImpl 管理拉来的数据:小数据放内存合并,大数据溢出到磁盘做外部合并           

为什么分内存和磁盘:数据量可能远超内存,外部归并排序能在有限内存下处理任意大数据。

• 进入 SORT 阶段:所有数据归并成一个全局有序的流           

为什么要排序:Reduce 要求相同 key 的数据连续到来,排序后才能按 key 分组迭代。

• 进入 REDUCE 阶段:按 key 分组迭代,调用 Reducer.reduce()

为什么按 key 分组:reduce 函数签名是 reduce(key, Iterable<Value>),需要同 key 的所有 value 聚在一起。

最终:Reduce 输出通过 OutputFormat 写到 HDFS(最终结果才上 HDFS)。

图:Reduce 端 Shuffle 拉取与归并——Fetcher 按分区拉取,内存/磁盘分级合并,最终归并成全局有序流


作业状态机

MapReduce 用状态机管理作业、任务、任务尝试三层生命周期。

 JobImpl 状态机 

作业状态:NEW → INITED → SETUP → RUNNING → COMMITTING → SUCCEEDED/FAILED/KILLED

// JobImpl 关键状态转换 (MRAppMaster 驱动): //   InitTransition:      NEW → INITED      创建 Map/Reduce 任务 //   StartTransition:     INITED → SETUP    开始执行 //   CommittingTransition: RUNNING → COMMITTING  输出提交 // // TaskAttemptImpl 状态机: //   NEW → ASSIGNED → RUNNING → COMMIT_PENDING → SUCCEEDED/FAILED // // DefaultSpeculator: 备份执行,6 个哨兵条件检测 straggler(拖后腿)任务 // OutputCommitter: AM端 vs Task端分工 //   setupJob/commitJob 在 AM  (作业级) //   setupTask/commitTask 在 Task (任务级)
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/MRAppMaster.java 
拆解逐字解读:InitTransition 在作业初始化时根据 InputFormat 的 split 数量创建对应数量的 Map task,再按配置创建 Reduce task。DefaultSpeculator 会监控 task 进度,如果某个 task 明显慢于平均水平(6 个哨兵条件判断),会启动一个备份 task 同时跑,谁先完成用谁的结果。OutputCommitter 把提交逻辑分两层:作业级(在 AM)负责整体提交/清理,任务级(在每个 Task)负责单个 task 的临时文件清理。
  调用链 · MapReduce 作业从提交到完成 

• Job.waitForCompletion(true) 用户入口           

为什么从这开始:这是用户代码调用点,true 表示打印进度。

• JobSubmitter.submitJobInternal() 提交前准备           

为什么先准备:要算 split 数、上传 jar、写配置,让所有节点都能拿到。

• InputFormat.getSplits() 计算输入分片,决定 Map 任务数           

为什么算 split:每个 split 对应一个 Map,split 大小影响并行度。

• 上传 jar/配置到 HDFS,YARNRunner.submitApplication() 提交到 YARN           

为什么上 HDFS:YARN 可能在任意节点启动 AM,资源必须放共享存储。

• RM 启动 MRAppMaster 容器,MRAppMaster.serviceInit() 初始化 AsyncDispatcher + JobImpl           

为什么用事件驱动:AM 内部组件解耦,AsyncDispatcher 统一分发事件。

• 发 JOB_INIT 事件 → InitTransition 创建 Map/Reduce task           

为什么事件驱动初始化:状态机模式,便于在状态转换时插入 hook。

• RMContainerAllocator 向 RM 申请容器,ContainerLauncher 在 NM 上启动 task 容器           

为什么分两步:先申请资源(RM 决定在哪启动),再启动(NM 执行)。

• MapTask.run() 跑 MapOutputBuffer → spill → 本地磁盘           

为什么先 Map:Map 的输出是 Reduce 的输入,必须先产出。

• Shuffle:Reduce 端 Fetcher 拉取 Map 输出 → MergeManagerImpl 归并           

为什么 Shuffle:把分散在各 Map 的同 key 数据聚到同一 Reducer。

• ReduceTask.run() 调 Reducer.reduce() → OutputFormat 写出最终结果到 HDFS           

为什么最后写 HDFS:最终结果要持久化共享,中间结果是临时的不上 HDFS。

最终:作业状态变 SUCCEEDED,结果文件落在 HDFS 指定输出目录。


TaskImpl:任务级状态机

Job 是作业级,Task 是单个 Map/Reduce 任务级。TaskImpl 管一个任务的多次尝试(TaskAttempt)——失败了重试,慢了起备份。

 任务状态与尝试调度 

• NEW → SCHEDULED:任务创建后进调度,等容器分配。

• SCHEDULED → RUNNING:第一次尝试在容器里跑起来。

• RUNNING → SUCCEEDED:尝试成功,任务成功。

• RUNNING → FAILED:尝试失败且重试次数耗尽(mapreduce.map.maxattempts / reduce.maxattempts,默认 4),整个任务失败。

• RUNNING → KILL_WAIT → KILLED:收到杀任务命令,等运行中的尝试退出。

 addAndScheduleAttempt:尝试的两种 Avataar 

每次创建新尝试时打上 Avataar 标签:VIRGIN(首次尝试)或 SPECULATIVE(推测备份尝试)。T_ADD_SPEC_ATTEMPT 事件由 Speculator 触发——当它判定某任务拖后腿时,发这个事件让 TaskImpl 起一个备份尝试。备份尝试和原尝试谁先成功用谁,另一个被杀。

易错易错:maxAttempts 默认 4 不是「失败 4 次才放弃」,而是「最多尝试 4 次」。第 1 次失败后还有 3 次机会。推测执行的备份尝试也算进 maxAttempts——所以频繁起备份可能反而耗尽重试次数。
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java 

TaskAttemptImpl:15 个状态的精细生命周期

TaskAttemptImpl 是最细粒度的状态机——一次尝试从申请容器到产出结果,要经历 15 个状态。其中两个「收尾中」状态(SUCCESS_FINISHING_CONTAINER / FAIL_FINISHING_CONTAINER)和容器回收逻辑是关键。

  调用链 · 一次 TaskAttempt 的关键转换 

• NEW → ASSIGNED:RequestContainerTransition 向 RMContainerAllocator 提交资源请求,带数据本地性提示

提示通过 RackResolver 算出输入 split 所在的节点/机架,RM 调度时优先分配到那——数据本地性省网络

• ASSIGNED → ContainerAssignedTransition:registerPendingTask 记录容器已分配但还没启动           

中间态:容器分配和进程启动之间有窗口,登记便于异常时清理

• → RUNNING:LaunchedContainerTransition 调 registerLaunchedTask,任务进程在容器里真正跑起来           

这一步才发 TASK_ATTEMPT_STARTED 事件给 JobImpl,作业据此更新进度

• RUNNING → SUCCESS_FINISHING_CONTAINER:任务逻辑成功完成,进入容器收尾           

为什么有「收尾中」状态:任务代码跑完不代表容器能马上回收,要等输出提交、日志聚合等收尾

• → SUCCEEDED:TaskAttemptFinishingMonitor 看门狗确认收尾完成           

看门狗防卡死:如果收尾阶段异常卡住,watchdog 超时后强制转 FAILED,避免永远悬在 FINISHING

• 失败路径:RUNNING → FAIL_FINISHING_CONTAINER → FAILED,同样有收尾过渡           

失败也要收尾:保留日志、释放容器资源,不能一失败就丢

15 个状态看似复杂,本质是「申请→分配→启动→运行→收尾(成功/失败)→终态」的主线,外加 COMMIT_PENDING(提交待确认)、KILLED 等分支。FINISHING 系列状态保证了容器回收的有序性。

易错易错:数据本地性是提示不是保证。RequestContainerTransition 告诉 RM「我希望跑在节点 A」,但如果 A 满了,RM 会降级到同机架(RACK_LOCAL),再不行跨机架(OFF_SWITCH)。MapReduce 容忍降级——只是慢一点,不会错。
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java 

RMContainerAllocator:AM 怎么向 RM 要资源

MRAppMaster 用 RMContainerAllocator 这个组件和 RM 交互。它负责把「作业需要多少 Map/Reduce 容器」翻译成对 RM 的资源请求,并处理分配结果。

 优先级与请求策略 

• 四级优先级:优先级数值越小越优先。PRIORITY_FAST_FAIL_MAP = 5(失败快速重试 Map,最高优先级)< PRIORITY_REDUCE = 10(Reduce)< PRIORITY_MAP = 20(常规 Map)< PRIORITY_OPPORTUNISTIC_MAP = 19(机会 Map)。

• Reduce 优先级高于常规 Map:按数值看 Reduce(10) 比 Map(20) 更优先。但实际 Reduce 通过 slowstart 延迟申请(mapreduce.job.reduce.slowstart.completedmaps,默认 Map 完成 5% 后才开始申请 Reduce),所以运行时观感是 Map 先跑——这是申请时机控制,不是优先级控制。

• 失败 Map 最优先:PRIORITY_FAST_FAIL_MAP(5) 数值最小,抢在 Reduce 和常规 Map 前面调度,尽快补上失败任务。

• 请求合并:同一节点/机架的多个请求会合并成批量请求,减少 RPC 次数。

易错易错:别以为「Map 优先级高于 Reduce」。按优先级数值,Reduce(10) 其实比常规 Map(20) 更优先。运行时 Map 先跑是因为 Reduce 的申请被 slowstart 延后了——这是时机差异,不是优先级差异。优先级只在同时有请求时才决定谁先分到容器。
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/RMContainerAllocator.java 

Uber 模式:小作业就地跑

如果作业特别小(几个 Map、数据量小),MRAppMaster 会走 Uber 模式——不向 RM 申请独立容器,直接在 AM 自己的进程里串行跑所有 Map 和 Reduce。省去了申请容器、启动进程、Shuffle 网络传输的开销。

  调用链 · makeUberDecision 判断条件 

• uberEnabled:配置 mapreduce.job.ubertask.enable=true(默认 false)           

必须显式开启——Uber 牺牲了并行性和容错,默认关

• smallNumMapTasks:Map 数 ≤ mapreduce.job.ubertask.maxmaps(默认 9)           

Map 太多串行跑会很慢,9 是经验上限

• smallNumReduceTasks:Reduce 数 ≤ mapreduce.job.ubertask.maxreduces(默认 1)           

多 Reduce 要多路 Shuffle,Uber 串行跑不划算

• smallInput:输入总大小 ≤ mapreduce.job.ubertask.maxbytes(默认 1 个 HDFS 块,约 128MB)           

数据太大 AM 内存装不下,且串行处理慢

• smallMemory + smallCpu:任务所需资源 ≤ AM 容器资源           

任务在 AM 进程里跑,不能超 AM 的资源配额

• notChainJob:不是链式作业(多个 MR 串联)           

链式作业有多个 Mapper 串联,Uber 不支持这种结构

全部满足才进 Uber。Uber 模式下用 LocalContainerAllocator(不向 RM 申请,本地伪造容器)和 LocalContainerLauncher(在 AM 进程内直接调 MapTask.run/ReduceTask.run),并关闭推测执行——串行单份,没有备份可言。

易错易错:Uber 不是「小作业加速器」,它是「省去 YARN 调度开销」。对极小作业(毫秒级计算),申请容器+启动 JVM 的开销可能比计算还大,Uber 让这些开销归零。但作业稍大就退回正常模式,串行反而更慢。
 hadoop-mapreduce-project/hadoop-mapreduce-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/MRAppMaster.java 

设计精髓:MapReduce 把「分布式计算」简化成两个函数(Map/Reduce)。开发者不用管数据怎么分片、任务怎么调度、节点怎么通信、失败怎么重试——框架全包了。代价是灵活性低(只有 Map 和 Reduce 两个算子),所以后来 Spark 用 RDD + 多算子取代了它。但 MapReduce 的 Shuffle 思想被所有大数据框架继承。 
易错修改指南:要调 Map 缓冲区,改 mapreduce.map.sort.spill.percent(spill 阈值)和 mapreduce.task.io.sort.mb(缓冲区大小)。spill 阈值调高能减少 spill 次数但增加单次排序内存压力;要加 Combiner,确保操作满足结合律(sum/max 可以,average 不行)。要调 Reduce 拉取并发,改 mapreduce.reduce.shuffle.parallelcopies。

AI 著书,人间讲义

口袋通识馆 · Hadoop 源码学习手册 · 第 06 章 · 油墨未干

全书 14 章已印毕 · 关注看连载 · 下一册写什么,留言区点单

相关学习资料