0%

Yarn工作机制

YARN 工作机制:一个 MapReduce 作业从提交到完成,经历了什么?

你写了一个 MapReduce 作业,打包成 JAR,执行 hadoop jar 提交到集群。

然后呢?这个作业怎么变成了集群上运行的 Map 和 Reduce 任务?YARN 在中间做了什么?

参与的五个角色

在开始之前,先认识一下参与这场”演出”的五个角色:

角色 简称 职责
客户端 Client 提交作业的人
资源管理器 RM 管集群总资源,谁用多少它说了算
节点管理器 NM 管单台机器的资源,启动容器
应用主控 AM 管单个作业的调度,每个作业一个
HDFS 存作业的代码、配置、输入输出数据

它们的关系:RM 管资源,AM 管作业,NM 管节点,Client 提交任务,HDFS 存数据。

分为四个阶段

yarn的工作流程

阶段 1:作业提交(步骤 1-4)

步骤 1:客户端初始化

用户执行 hadoop jar 命令后,客户端创建 YarnRunnerJobSubmitter,检查作业配置(输入输出路径是否存在、Mapper/Reducer 类对不对)。

步骤 2:申请 Application ID

客户端向 RM 发请求,RM 返回一个唯一的 Application ID(形如 application_1620000000000_0001)。这个 ID 是作业在整个生命周期里的唯一标识。

同时,RM 在 HDFS 上为这个作业创建目录:/user/<用户名>/staging/<Application ID>/

步骤 3:上传作业资源到 HDFS

客户端把三样东西传到 HDFS 的 staging 目录:

  • Job.split:输入数据的分片信息——决定起多少个 Map 任务
  • Job.xml:作业的所有配置
  • JAR 包:你写的 Map/Reduce 代码

步骤 4:正式提交作业

客户端告诉 RM:”东西都准备好了,可以开始了。”RM 收到请求后,把作业放入调度队列。

阶段 2:AM 启动(步骤 5-6)

步骤 5:RM 调度并启动 AM

RM 的调度器从集群里找一个空闲节点,分配一个 Container(通常 1-2GB 内存 + 1 核 CPU),通知该节点的 NM 启动 AM。

对于 MapReduce 作业,AM 就是 MRAppMaster 进程。

步骤 6:AM 初始化

AM 启动后:

  • 从 HDFS 加载 Job.split 和 Job.xml
  • 初始化作业状态跟踪器(记录每个任务的进度)
  • 向 RM 注册自己,建立心跳(默认 3 秒一次)

此时,AM 已经拿到了”指挥权”——它负责管理这个作业的一生。

阶段 3:资源申请与任务启动(步骤 7-9)

步骤 7:AM 解析任务需求

AM 读取输入分片信息,算出需要多少个 Map 任务(= 分片数)。再读取 mapreduce.job.reduces 配置,算出需要多少个 Reduce 任务。

步骤 8:向 RM 申请 Container

AM 向 RM 申请资源:

  • Map 任务:优先申请数据所在节点的容器(数据本地性),减少网络传输
  • Reduce 任务:没有本地性要求,哪个节点有空位就行

申请的时候带”本地性级别”:NODE(同一个节点)> RACK(同一个机架)> ANY(任何节点)。RM 尽量满足 NODE,但如果那个节点没资源,就降级到 RACK 或 ANY。

步骤 9:启动任务

RM 分配 Container 后,AM 通知对应节点的 NM 启动任务。NM 在 Container 里启动 YarnChild 进程,执行 Map 或 Reduce 任务。

Map 和 Reduce 任务的区别:Map 任务优先本地读数据,Reduce 任务从各个 Map 拉数据。这俩一个是”本地读”,一个是”网络拉”——所以 Reduce 的资源配置通常要比 Map 大。

阶段 4:任务执行与完成(步骤 10-11)

步骤 10:资源本地化

YarnChild 进程启动后,从 HDFS 下载作业 JAR 包和配置文件到本地,把依赖的资源拉到节点上。Map 任务还要读取 HDFS 上的输入数据块。

步骤 11:执行任务,持续汇报

  • Map 任务:读数据 → 执行 map() → 输出中间结果到本地磁盘
  • Reduce 任务:拉取 Map 输出 → 执行 reduce() → 输出结果到 HDFS

持续汇报:每个任务通过 YarnChild 向 AM 汇报进度(Map 完成了 50%、Reduce 完成了 30%)。AM 汇总后通过心跳告诉 RM。

故障处理:如果任务失败,AM 会申请新 Container 重试(默认 4 次)。如果 AM 自己挂了,RM 会重新启动它。

作业完成:所有任务完成后,AM 向 RM 汇报”作业结束”,释放所有 Container,AM 进程退出。

两个核心交互机制

心跳——“我还活着”

  • AM → RM:每 3 秒一次,汇报进度 + 申请资源
  • NM → RM:每 3 秒一次,汇报节点资源使用情况

心跳断了,RM 就知道对方可能挂了,触发故障处理。

Container——资源的容器

Container 是 YARN 的资源抽象,包含:内存大小、CPU 核心数、所在节点、启动命令。任务就运行在 Container 里,任务结束了 Container 就被回收。

几个关键配置

参数 默认值 说明
mapreduce.map.memory.mb 1024MB 每个 Map 任务的内存
mapreduce.reduce.memory.mb 1024MB 每个 Reduce 任务的内存(通常要比 Map 大)
mapreduce.map.cpu.vcores 1 每个 Map 任务的 CPU 核数
mapreduce.reduce.cpu.vcores 1 每个 Reduce 任务的 CPU 核数
mapreduce.map.maxattempts 4 Map 任务最大重试次数
mapreduce.task.timeout 600000ms (10min) 任务超时时间

Reduce 内存为什么通常比 Map 大? Reduce 要拉取所有 Map 的输出,在内存里做合并排序,内存压力比 Map 大。

常见问题

1. 作业一直卡在 ACCEPTED

提交后状态一直是 ACCEPTED,说明在等资源。去 RM Web UI 看队列资源是不是满了,或者单个 Container 规格太大了(比如申请 8GB,但所有节点都只有 6GB 可用)。

2. Container 被 NodeManager 杀掉

日志里看到 Container is running beyond memory limits。关掉这两个开关通常能解决:

xml

1
2
3
4
5
6
7
8
<property>
<name>yarn.nodemanager.pmem-check-enabled</name>
<value>false</value>
</property>
<property>
<name>yarn.nodemanager.vmem-check-enabled</name>
<value>false</value>
</property>

3. 任务失败重试太多

mapreduce.map.maxattempts 设太大了,导致一个坏任务反复重试浪费资源。把重试次数调到 2-4 次,同时排查为什么会失败。

欢迎关注我的其它发布渠道