深入 Hadoop 15 - MRv2 on YARN ApplicationMaster 与 Task Attempt
上一篇讲了 Shuffle。本篇讲 MapReduce 在 YARN 上的具体实现——MRv2。这是 MapReduce 计算模型篇的最后一篇。
MRv2 常被介绍成"MapReduce 在 YARN 上的版本"。这个描述对应了部署形态但漏了内部架构。准确的说法是:MRv2 把 Hadoop 1.x 的 JobTracker 拆成 ResourceManager(资源管理)和 MRAppMaster(作业编排)两部分,MRAppMaster 是一个普通的 YARN ApplicationMaster,运行在 container 里。每个 Task 也是一个 container(YarnChild),通过 umbilical RPC 与 MRAppMaster 通信。MRv2 的 Speculative Execution、Counter、Task Recovery 等机制都是 MRAppMaster 实现的,与 YARN 基础设施解耦。
本篇只抓一个问题:MRAppMaster 内部是怎么编排 Map / Reduce Task 的、Task Attempt 怎么向 AM 报告进度、Speculative Execution 的触发条件和代价、AM 重启时怎么恢复 Task 状态。
MRAppMaster 的内部结构
MRAppMaster 是 MapReduce 框架的 ApplicationMaster 实现。它本身是一个普通 Java 进程,跑在 YARN 分配的 container 里(container-0)。内部结构可以拆成几个核心组件:
1 | |
Job Init 阶段做几件事:
读取作业配置文件(job.xml,从 HDFS 加载)。配置包含所有 MapReduce 参数——Mapper 类、Reducer 类、InputFormat、OutputFormat、Reduce 数、Combiner 等。
作业提交阶段使用 InputFormat.getSplits(JobContext) 计算 InputSplit 列表,并把切分元数据写入 staging 目录。InputFormat 是抽象类,定义 getSplits(...) 和 createRecordReader(...) 两个核心方法;FileInputFormat 是文件输入的抽象基类,提供通用的 getSplits(...) 实现,TextInputFormat、SequenceFileInputFormat 这类具体类再决定记录读取和可切分边界。CombineFileInputFormat 也是抽象类,它返回 CombineFileSplit,实际作业通常使用 CombineTextInputFormat、CombineSequenceFileInputFormat 或自定义子类。每个 InputSplit 包含数据位置和长度,MRAppMaster 读取这些 split 元数据后做数据本地化调度,优先在持有数据的节点上启动 Map Task。
初始化作业状态机——作业进入 RUNNING 状态,准备启动 task。
Job Tracker(internal)是 MRAppMaster 内部的 task 调度器。它决定哪个 task 在哪个节点上跑、什么时候启动 Reduce Task、什么时候触发 Speculative Execution。注意这个 Job Tracker 与 Hadoop 1.x 的 JobTracker 完全不同——前者是作业内调度(一个 MRAppMaster 一个),后者是集群级调度(一个集群一个)。命名巧合,不是同一个东西。
Job Tracker(internal)的调度逻辑:
1 | |
这种"Map 优先 + Reduce 提前启动"的策略是 MapReduce 性能调优的关键维度之一。slowstart 设得低(例如 0.05)让 Reduce 提前启动,可以隐藏 Reduce Task 的启动开销;设得高(例如 0.5)让 Reduce 等更久但避免浪费资源。
Task Attempt 的生命周期
每个 Task(例如某个具体的 Map Task “task_…_m_000000”)在执行过程中可能有多次 Attempt。第一次 Attempt 失败时 AM 会启动新 Attempt,直到达到最大重试次数(默认 4,mapreduce.map.max.attempts / mapreduce.reduce.max.attempts)。
Task Attempt 的状态机:
1 | |
NEW:刚创建,未分配 container。
UNASSIGNED:已分配 container,AM 准备启动。
ASSIGNED:container 已分配到某个 NM,AM 已通过 RPC 让 NM 启动。
RUNNING:container 进程已启动,开始执行 MapTask 或 ReduceTask。
COMMIT_PENDING:Task 完成,请求 AM 提交输出(防止多个 Attempt 同时写 OutputFormat)。
SUCCEEDED:Task 输出已提交,Task 成功完成。
FAILED:Task 失败,AM 决定是否启动新 Attempt。
KILLED:Task 被强制终止(Speculative、用户 kill、抢占)。
每次 Attempt 启动都创建一个全局唯一的 AttemptId(例如 attempt_1721900000000_0001_m_000000_0)。这个 ID 后缀 _0、_1、_2 表示第几次 Attempt。
Attempt 失败的常见原因:
container 进程崩溃(OOM、JVM crash、被 cgroups 杀)。
container 与 AM 心跳超时(默认 mapreduce.task.timeout = 600000ms 即 10 分钟)。这种"挂起的 Task"AM 没收到心跳,标记为 FAILED。Task 实际可能还在跑,但 AM 已经启动新 Attempt。
Task 主动抛出未捕获异常(用户 Mapper/Reducer 代码 bug)。
NM 故障(节点宕机),AM 通过 RM 知道该 NM dead,把该 NM 上所有 Task Attempt 标记 FAILED。
Umbilical RPC:Task 与 AM 的通信
Task 进程(YarnChild)启动后通过 umbilical RPC 与 MRAppMaster 通信。这是一个简单的 RPC 协议,承担三类信息:
进度报告:Task 周期性向 AM 报告当前进度(0.0 - 1.0 浮点数)。AM 用这个进度判断 Task 健康度(长时间没进度更新可能是 Task 挂起)。
Counter 上报:Task 把 Counter(Map 输入记录数、Spilled records 等)实时上报给 AM。AM 聚合所有 Task 的 Counter,在 Web UI 展示。
完成通知:Task 完成后通过 umbilical RPC 调用 AM 的 done() 方法,告诉 AM 自己的输出可以提交。
umbilical RPC 的默认心跳间隔 1 秒(mapreduce.task.progress.report.interval),超时 10 分钟(mapreduce.task.timeout)。
注意 umbilical 是 Task 主动调用 AM 的 RPC,不是 AM 拉取 Task——这与 NM 与 RM 的关系(NM 主动 RPC RM)一致。Task 是 RPC client,AM 是 RPC server。
如果 Task 不报告进度超过 timeout,AM 标记该 Task 为挂起并启动新 Attempt。这个机制防止 Task 死循环或 GC 停顿导致的无限等待。代价是误判——长时间 GC 停顿的 Task 可能被错误标记 FAILED。处理这种场景的方法是调大 timeout(mapreduce.task.timeout),或者让用户代码在长时间操作前调用 context.progress() 报告进度。
Speculative Execution:对抗慢节点
集群里某些节点比平均慢——磁盘老化、网络抖动、CPU 负载高。如果一个 MapReduce 作业有 1000 个 Map Task,999 个在 10 分钟内完成,1 个慢节点上的 Map Task 跑了 1 小时,整个作业延迟被这个慢 Task 拖到 1 小时。
Speculative Execution(推测执行)是 MapReduce 对抗慢节点的机制。AM 检测哪些 Task 比平均慢,启动备份 Attempt。哪个 Attempt 先完成就用哪个,另一个被 kill。
具体检测算法(DefaultSpeculator):
1 | |
Speculative Execution 的配置:
1 | |
Speculative Execution 的代价是资源浪费——备份 Attempt 占用额外 container。如果集群资源紧张,Speculative 抢占其他作业的资源。生产环境通常在 Map 阶段开启(Map 输出中间结果,备份代价小),Reduce 阶段关闭(Reduce 输出大文件,备份代价大)。
注意 Speculative 与 Combiner、Pipeline 容错的差异。Combiner 是数据优化(减少 Shuffle 流量),Pipeline 容错是传输层重试,Speculative 是 Task 层重试。三者独立,可以同时启用。
Counter:性能观测的核心
MapReduce 提供 Counter 机制让 Task 上报指标。Counter 是 (group, name, value) 三元组,Task 通过 context.getCounter(group, name).increment(1) 上报。AM 聚合所有 Task 的 Counter 在 Web UI 展示。
Hadoop 内置的 Counter 分几类:
1 | |
用户也可以定义自定义 Counter。在 Mapper 或 Reducer 里调用 context.getCounter("MyGroup", "MyCounter").increment(1)。这些 Counter 自动上报给 AM,在 Web UI 可见。
Counter 是 MapReduce 性能调优的核心观测点。每个 Counter 反映一个具体维度,通过对比 baseline 和调优后的 Counter 数值,可以判断调优是否有效。
例如调 io.sort.mb:调优前 Spilled records = 5M,调优后 = 1M,说明 Spill 减少 80%。配合 Filesystem Counter 的本地磁盘字节变化,可以确认 I/O 优化效果。
AM 重启时的 Task Recovery
MRAppMaster 自己是一个 YARN ApplicationAttempt。如果 AM 进程崩溃,RM 会启动新 Attempt(默认 max-attempts = 2,可配置更高)。新 Attempt 必须恢复作业状态,否则已完成的 Task 全部丢失。
MRv2 的 Recovery Manager 负责这个过程:
1 | |
Recovery 的具体行为取决于 yarn.mapreduce.am.job.recovery.enable(默认 true)和 yarn.mapreduce.am.job.recovery.task-state-limiting(默认 2000,超过这个数量 Task 的作业不做 recovery,直接重跑所有 Task)。
Recovery 的代价:
OutputFormat 必须是幂等的。Map Task 的输出在本地磁盘,AM 崩溃后这些输出可能丢失(节点没崩但 AM 进程崩了,输出还在;如果节点也崩了,输出丢失)。Recovery 时已 SUCCEEDED 的 Map Task 被认为输出已持久化(OutputCommitter.commitTask 已调用),不重跑。
Reduce Task 的输出在 HDFS。COMMIT_PENDING 但未 COMMITTED 的 Reduce Task 在 Recovery 时被重新调度(输出未提交)。
Recovery 不是完美的——某些 Task 可能重跑,造成额外开销。但相比"AM 崩溃作业失败需要重新提交",Recovery 的代价可以接受。
实验:观察 MRv2 行为
实验状态:UNVERIFIED_RUNTIME。下面是验证步骤,本轮没有连接真实 Hadoop 集群执行。
提交一个 MapReduce 作业,在 ApplicationMaster Web UI(每个作业有独立 Web UI,URL 形如 http://am-host:random-port)观察:
Job 页面展示作业整体进度、Map/Reduce 完成度、作业 Counter。Tasks 标签页列出所有 Task Attempt,每个 Attempt 显示状态、所在节点、进度、运行时间、Counter。
故意 kill 某个 Task 进程(找到 Task 所在 NM,kill -9 YarnChild 进程),观察 AM 是否在几秒内检测到并启动新 Attempt。新 Attempt 的 AttemptId 后缀 _1、_2 等。
提交一个故意慢的作业(Mapper 里加 Thread.sleep),让一个 Task 显著慢于其他 Task,观察 AM 是否启动 Speculative Attempt。Speculative Attempt 在 Web UI 显示为同一 Task ID 的两个 Attempt 并行跑。
调小 mapreduce.task.timeout 到 30 秒,提交一个 Mapper 里 Thread.sleep(60) 的作业,观察 AM 是否把 Task 标记为 FAILED(心跳超时),然后启动新 Attempt。
模式提炼
MRv2 on YARN 体现的设计模式:
1 | |
这个模式不只是 MRv2。Spark on YARN 的 Driver / Executor 结构是同构——Driver 跑在 AM container,Executor 跑在 task container,通过 RPC 通信。Spark 的 Speculative 机制(spark.speculation)与 MapReduce 的 Speculative 几乎相同。Flink on YARN 的 JobManager / TaskManager 也是同构。
差异在持久化策略。MapReduce 用 HDFS 文件做 Recovery,Spark 用 Event Log,Flink 用 Checkpoint。具体机制各不相同,但"AM 重启后从持久化状态恢复"这个模式是一致的。
工程迁移表
| MRv2 概念 | Spark on YARN | Flink on YARN | Tez on YARN | Hadoop 1.x JobTracker |
|---|---|---|---|---|
| MRAppMaster | Driver | JobManager | DAG AppMaster | JobTracker(部分) |
| YarnChild | Executor | TaskManager | Task | TaskTracker + Task |
| Umbilical RPC | Driver ↔ Executor | JM ↔ TM RPC | Tez RPC | TaskUmbilicalProtocol |
| Speculative Execution | spark.speculation | restart strategy | DAG spec | 是 |
| Counter | accumulator | metric | counter | counter |
| Recovery Manager | Event Log | checkpoint | DAG recovery | job history |
| 数据本地化 | preferred location | - | input locality | data-local slot |
注意 Spark on YARN 这一列。Spark 的 Driver 跑在 AM container(cluster 模式)或 Client 进程(client 模式)。Executor 跑在 task container,与 MapReduce 的 YarnChild 同构。差异在 Spark 的 Executor 是长期存活(Driver 与 Executor 保持长连接),MapReduce 的 YarnChild 是 Task 级别(Task 完成 YarnChild 退出)。这让 Spark 的 Executor 启动开销可以分摊到多个 Task,是 Spark 比 MapReduce 快的另一个原因。
常见误解
误解一:“MRAppMaster 是 Hadoop 1.x JobTracker 的简单移植”。MRAppMaster 与 JobTracker 是不同的代码——JobTracker 是单体进程管理集群所有作业,MRAppMaster 是每作业一个进程。共享的只是部分调度逻辑(task 分配、数据本地化)。
误解二:“Speculative Execution 总是让作业变快”。Speculative 启动备份 Attempt 占用资源,如果集群资源紧张,备份抢其他作业资源。Speculative 适合"作业里个别 Task 慢"场景,不适合"所有 Task 都慢"场景(所有 Task 都慢说明作业本身设计有问题,不是慢节点问题)。
误解三:“AM 重启会自动恢复所有 Task”。Recovery 有上限(默认 2000 Task)。超过这个数量 AM 重启时所有 Task 重跑。另外 Recovery 要求 OutputFormat 幂等,自定义 OutputFormat 必须保证多次 commit 不产生不一致。
误解四:“Task 心跳超时就是 Task 死了”。Task 可能因为长时间 GC 或长时间同步 IO(HBase scan)不报告进度,被 AM 误判 FAILED。处理方法是调大 timeout 或者用户代码主动调用 context.progress() 报告进度。
误解五:“MapReduce 的 Counter 是只读的统计信息”。Counter 不仅可以读,还可以用作控制——用户代码可以读 Counter 判断"已处理多少记录",达到阈值时主动跳出 map/reduce 循环。这种"用 Counter 做控制流"在某些场景有用(例如采样作业处理 N 条记录就停)。
练习
-
提交一个 MapReduce 作业,在 ApplicationMaster Web UI 上找到 Tasks 标签页。每个 Task 显示 Attempt 数——故意 kill 某个 Task 进程,观察 AM 是否启动新 Attempt(Attempt 数变 2)。
-
修改作业配置
mapreduce.map.speculative=true,提交一个故意让某个 Mapper 慢的作业(Mapper 里 Thread.sleep),观察 Speculative Attempt 是否启动,最终作业完成时间是否被慢 Mapper 拖累。 -
在
apache/hadoop源码里找到MRAppMaster.java、YarnChild.java、TaskAttemptListener.java、DefaultSpeculator.java(hadoop-mapreduce-client-app 模块),观察 AM 内部组件的实现。 -
思考题:如果让 MapReduce 的 Speculative 不仅对抗慢节点,还对抗"用户代码 bug 导致 Task 偶发失败"(例如某些记录触发 NPE),Speculative 的设计需要怎么改造?这种"备份对抗正确性问题"的思路在分布式系统里是否有先例?
系列导航
参考资料
- Jeffrey Dean, Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters. OSDI 2004. Section 4 “Backup Tasks” 描述了 Speculative Execution 的早期设计。
- Vinod Kumar Vavilapalli et al. Apache Hadoop YARN: Yet Another Resource Negotiator. SOCC 2013. Section 6 “MapReduce on YARN” 描述了 MRv2 的具体实现。
- Apache Hadoop 官方文档:MapReduce on YARN. https://hadoop.apache.org/docs/r3.4.1/hadoop-mapreduce-client/hadoop-mapreduce-client-core/MapReduceTutorial.html
- Apache Hadoop 3.4.1 Javadoc:
InputFormat#getSplits(JobContext)与FileInputFormat#getSplits(JobContext)。https://hadoop.apache.org/docs/r3.4.1/api/org/apache/hadoop/mapreduce/InputFormat.html - Apache Hadoop 源码:
MRAppMaster.java、YarnChild.java、TaskAttemptListener.java、DefaultSpeculator.java. https://github.com/apache/hadoop - Tom White. Hadoop: The Definitive Guide. O’Reilly, 4th Edition 2015. Chapter 7 详细描述了 MRv2 在 YARN 上的运行机制。
- Maysam Yabandeh et al. Speculative Execution in MapReduce.(Microsoft Research 关于 Speculative Execution 算法的改进研究)
