深入 Hadoop 07 - YARN 架构与三方契约
第一阶段把 HDFS 切完了。从这一篇起进入 YARN,先讲架构。
YARN 常被介绍成"Hadoop 的资源管理器"。这个描述对应了职责但漏了关键点。准确的说法是:YARN 是一个把"全局调度"和"作业内调度"拆开的双层调度系统,引入了 ApplicationMaster 这个新角色承担作业级编排,把 Hadoop 1.x 的 JobTracker 单体拆成 ResourceManager + ApplicationMaster + NodeManager 三个独立契约。Container 抽象是这个拆分的基础——它把"一段资源 + 一段进程"打包成统一的调度单位,让同一套基础设施能同时承载 MapReduce、Spark、Flink、Tez、Storm 等不同计算模型。
本篇只抓一个问题:YARN 的三方契约(ResourceManager、NodeManager、ApplicationMaster)是怎么切开的,Container 抽象的本质是什么,为什么 Hadoop 2.x 必须从 JobTracker 演化到 YARN。
Hadoop 1.x JobTracker 的痛点
Hadoop 1.x 的资源管理由 JobTracker 单体承担。JobTracker 同时做两件事:
1 | |
这三件事在 Hadoop 1.x 早期规模下(几百节点)可以接受。但集群规模到 1000 节点以上时暴露了三个结构性问题:
第一,JobTracker 单点瓶颈。集群里所有作业的元数据、所有 TaskTracker 的状态、所有 task 的进度都集中在一个 JVM 里。1000 节点集群、每节点 10 个 slot、平均 1000 个并发作业、每作业 1000 个 task,JobTracker 内存里要维护 1000 万级别的 task 状态,OOME 频发。
第二,作业编排逻辑与资源管理耦合。JobTracker 内部把 MapReduce 的作业编排逻辑(task 失败重试、speculative)硬编码进了资源管理器。任何新计算模型(流式计算、图计算、迭代式机器学习)想用 Hadoop 集群,要么把自己塞进 MapReduce 模型(很别扭),要么自己写一个 JobTracker 替代品(不可能)。
第三,资源模型过于粗粒度。Hadoop 1.x 用 slot 作为资源单位——每个节点固定 N 个 map slot 和 M 个 reduce slot。这个模型对 MapReduce 刚好够用,但其他计算模型(Spark 用 cores + memory、Flink 用 slots)无法精确表达资源需求。slot 配多了浪费,配少了争抢。
这三个问题的本质是"资源管理"和"作业编排"是两个独立关注点,强行耦合在 JobTracker 里让两者都受限。Hadoop 2.x(2013 年)引入 YARN 把这两个关注点拆开。
YARN 三方契约
YARN 的核心抽象是三个独立角色的契约:
1 | |
三个角色的职责分工:
ResourceManager(RM)是全局唯一的中央调度者。它知道集群里每个节点的资源总量和当前空闲,按调度策略把 container 分配给各个 ApplicationMaster。RM 不关心作业内部的 task 编排,只关心"哪个作业申请了多少资源、当前集群能给它多少"。
NodeManager(NM)是每台节点上常驻的进程。它管理本节点的 container 生命周期——启动 container 进程、监控 container 资源使用(CPU、内存)、向 RM 心跳汇报本节点状态、在 container 超出资源限制时杀掉它。NM 不知道 container 里跑的是什么——可以是 MapReduce Task、Spark Task、Flink Task,NM 只关心资源边界。
ApplicationMaster(AM)是每个作业独占的特殊 container,承担作业级编排。AM 自己也是一个 container(通常叫 container-0),RM 调度它启动后,AM 接管后续编排——向 RM 申请更多 container、在新 container 里启动 task、监控 task 进度、处理失败重试。AM 是计算框架特定的——MapReduce 有 MRAppMaster,Spark 在 YARN 上有自己的 ApplicationMaster(入口类为 org.apache.spark.deploy.yarn.ApplicationMaster),Flink 有 YarnJobClusterEntrypoint。
Client 是提交作业的入口。Client 把作业 jar、配置、依赖上传到 HDFS,向 RM 提交 ApplicationSubmissionContext,RM 收到后启动这个作业的 AM(作为 container-0)。后续整个作业的执行由 AM 接管,Client 可以断开。
为什么拆出 ApplicationMaster
拆出 AM 是 YARN 最关键的设计决策。AM 解决了 Hadoop 1.x 的两个结构性问题:
第一,作业编排逻辑解耦。AM 是计算框架特定的——MapReduce 的 AM 实现 MapReduce 的 task 编排逻辑,Spark 的 AM 实现 Spark 的 stage 切分和 task 派发,Flink 的 AM 实现 Flink 的 operator chain 调度。YARN 不需要知道这些细节,只负责给每个 AM 分配 container。这让 YARN 成了真正的"通用资源管理器"——任何遵循 AM 契约的计算框架都能跑在 YARN 上。
第二,JobTracker 单点瓶颈消失。原来 JobTracker 内部维护的"每个作业的 task 状态"现在分散到每个作业的 AM 里。一个 1000 作业并发的集群,原来 JobTracker 要维护 1000 × 1000 = 100 万 task 状态,现在每个 AM 只维护自己作业的 1000 个 task。RM 自身只维护 1000 个 AM 的状态,内存压力下降 1000 倍。
AM 拆分还带来一个隐性收益——作业级故障隔离。Hadoop 1.x 里某个作业的 task 编排逻辑出 bug 让 JobTracker 崩溃,整个集群所有作业都受影响。YARN 里某个作业的 AM 崩溃只影响这个作业,其他作业的 AM 和 RM 不受影响。RM 在 AM 崩溃后可以按策略重启这个 AM(默认最大重试次数由 yarn.resourcemanager.am.max-attempts 控制,默认 2 次)。
双层调度:中央调度 + 二次调度
YARN 的调度模型可以拆成两层:
1 | |
双层调度的核心收益是"中央调度只做粗粒度决策,作业内调度做细粒度决策"。RM 不需要知道每个 task 的依赖关系(Spark 的 stage 边界、Flink 的 operator chain),只关心"给这个 AM 多少 container"。AM 知道自己作业的内部结构,决定 container 怎么用。
这种双层结构与 Kubernetes 的两层调度(中央 scheduler + 每个 Pod 的 controller manager)有相似性。Omega 论文(EuroSys 2013)描述的 “shared state + optimistic concurrency” 是更激进的版本;Borg 论文则发表于 EuroSys 2015。YARN 选了中间路线:中央调度器只做粗粒度资源分配,作业内调度交给各个 AM。
双层调度的一个副作用是调度延迟。MapReduce 作业启动时要先调度 AM,AM 启动后再向 RM 申请 task container。Spark 动态资源分配只在同一个 Spark application 内按负载增减 executor:executor 空闲超过 spark.dynamicAllocation.executorIdleTimeout 后会被移除,不会把 container 留给下一个作业复用,因此它不消除新作业的 AM 启动开销。
Container 抽象的本质
Container 是 YARN 最基础的抽象单位。Container 不是 Docker 容器——YARN 在 Hadoop 2.x 引入 Container 概念时 Docker 还没流行。Container 的本质是"一段资源 + 一段进程"的打包:
1 | |
Resource 描述这个 container 用多少内存和 CPU。NodeId 标识 container 在哪个节点上跑。Priority 让 AM 在多个 container 申请间排序。ContainerToken 是 RM 签名的令牌,AM 拿这个令牌去对应 NM 启动 container,防止 AM 越权申请别人的资源。Command 是 container 启动后执行的命令——通常是计算框架的 task launcher 脚本。
默认的 DefaultContainerExecutor 不创建 cgroup。NodeManager 的 ContainersMonitor 会周期检查进程树的物理内存和虚拟内存,超限后由 NM 终止 container;只有配置 LinuxContainerExecutor 与 cgroups handler 后,CPU 和内存边界才由 cgroups 强制。
即使启用 cgroups,yarn.nodemanager.linux-container-executor.cgroups.strict-resource-usage 默认仍为 false;此时 CPU 主要按 shares 竞争,只有显式开启严格模式才按申请的 vCores 设置硬上限。
Container 的本质抽象是"把作业级资源需求标准化"。MapReduce 的 task 需要 2GB 内存、Spark 的 task 需要 4GB 内存、Flink 的 task slot 需要 1GB 内存,YARN 都用 Container 统一表达。这让同一套 YARN 集群能同时承载 MapReduce、Spark、Flink,不需要为每个计算模型单独配置资源池。
注意 Container 和 Docker 的差别。YARN Container 最初只是资源配额与进程生命周期抽象,没有镜像概念。Hadoop 2.6.0 加入的是后来废弃的 DockerContainerExecutor;后续实现是在 LinuxContainerExecutor 下提供 DockerLinuxContainerRuntime。Hadoop 3.4.1 的 Docker 容器文档把 yarn.nodemanager.runtime.linux.allowed-runtimes 示例写成 default,docker,javasandbox;runC 有单独的 runtime 文档,Podman 不是 3.4.1 文档里的 LinuxContainerExecutor runtime。
实验:观察 YARN 三方契约
UNVERIFIED_RUNTIME:下面界面字段和命令输出形态用于在真实 YARN 集群上核对现象,本轮未连接 live Hadoop 集群运行。
YARN 的 ResourceManager Web UI 默认监听 8088 端口。打开 http://resourcemanager-host:8088 可以看到集群状态。
Cluster 标签页展示:
1 | |
Apps Pending 和 Containers Pending 是观察调度延迟的关键指标。如果 Containers Pending 长期不为 0,说明集群资源紧张,AM 在排队等 container。
Scheduler 标签页展示当前调度器配置(Capacity Scheduler 或 Fair Scheduler),以及每个队列的资源使用。第九篇会展开调度器对比。
Applications 标签页列出所有作业。每个作业有 ApplicationId(形如 application_1721900000000_0001)、作业类型(MapReduce / Spark / Flink)、当前状态、运行用户、队列、AM 所在节点。
点开某个作业可以看到 ApplicationMaster 的详细信息——AM container 在哪个 NM 上、AM 已申请多少 container、当前活跃 task 数。
命令行层面,yarn application -list 列出所有作业,yarn node -list 列出所有 NM,yarn container -list <app-attempt-id> 列出某个作业的所有 container(这个命令需要 Hadoop 3.x)。
模式提炼
YARN 三方契约体现的设计模式:
1 | |
这个模式不只是 YARN。Mesos(Apache 2010)的双层调度(Mesos Master + Framework Scheduler)是同一思路的具体化,差异在 Mesos 用 offer 模型(master 主动推送资源 offer 给 framework)、YARN 用 request 模型(AM 主动向 master 申请)。Borg / Omega(Google 2015 论文)的 shared state 调度是更激进的去中心化版本。
Kubernetes 的两层抽象(API Server + kubelet + 各 controller)也有类似结构,但 Kubernetes 把"作业级编排"留给了 Deployment / StatefulSet / Job 等 controller,而不是把每个作业都拆出独立的编排器进程。YARN 的 AM 模型让"作业"和"编排器"是 1:1 的,Kubernetes 让多个作业共享一个 controller(Deployment controller 管理所有 Deployment)。
差异源于目标场景。YARN 的"作业"是批处理计算任务(运行几分钟到几小时),独占 AM 进程的开销可以接受。Kubernetes 的"作业"是长期服务(Deployment),不需要每个 Pod 一个编排器。
工程迁移表
| YARN 概念 | Mesos | Borg | Kubernetes | Spark Standalone |
|---|---|---|---|---|
| ResourceManager | Mesos Master | Borgmaster | API Server + Scheduler | Master |
| NodeManager | Mesos Agent | Borglet | kubelet | Worker |
| ApplicationMaster | Framework Scheduler | job controller | Deployment / Job controller | Driver |
| Container(资源-进程打包) | Task | alloc | Pod | Executor |
| 调度模型 | offer(master 推送) | optimistic | request(worker 拉) | request |
| cgroups 隔离 | 是 | 是 | cgroups + namespace | JVM 隔离(弱) |
注意 Kubernetes 这一列的差异。Kubernetes 的 Pod 比 YARN 的 Container 概念更丰富——Pod 是一组共享网络和存储 namespace 的 container 集合,Container 是单个进程。这种差异让 Kubernetes 更适合微服务部署(一个 Pod 里多个协作 container),YARN 更适合批处理作业(一个 Container 一个 task 进程)。Hadoop 3.x 引入的 YARN Service 试图让 YARN 也支持长期服务部署,但市场份额已被 Kubernetes 占据。
常见误解
误解一:“YARN 取代了 Hadoop 1.x 的 JobTracker”。YARN 取代的是 JobTracker 的资源管理部分。JobTracker 还承担的作业编排(task 调度、失败重试、speculative)被拆到了 MRAppMaster 里——MRAppMaster 是 MapReduce 框架的 AM 实现。其他计算框架(Spark / Flink)有自己的 AM 实现。所以"JobTracker 被拆了"更准确,"被 YARN 取代"是简化说法。
误解二:“Container 就是 Docker 容器”。YARN Container 是资源与进程生命周期抽象;默认 DefaultContainerExecutor 直接启动本地进程。需要镜像隔离时,管理员必须另行配置 LinuxContainerExecutor 和 Docker 或 runc runtime。把 YARN Container 等同于 Docker Container 是常见误解。
误解三:“ResourceManager 是单点”。非 HA 部署里 RM 是单点,HA 部署里有多个 RM 通过 ZooKeeper 选主。第十一篇会展开。即使是单 RM,AM 重启机制(yarn.resourcemanager.am.max-attempts)能让 AM 在崩溃后自动重启,配合 RM 的 ApplicationAttempt 状态恢复,作业可以续跑。
误解四:“YARN 只能跑 MapReduce”。YARN 的设计目标就是通用资源管理——任何遵循 AM 契约的计算框架都能跑。生产集群上同时跑 MapReduce、Spark、Flink、Tez、Storm 是常态。Hive on Tez、Hive on Spark、Spark SQL、Flink streaming 都跑在 YARN 上。
误解五:“AM 启动是免费的”。AM 自己是一个 container,启动要承担 JVM 启动、初始化、向 RM 注册和申请 container 的开销。长时间运行的 Spark Streaming application 会自然保留自己的 AM;spark.yarn.maxAppAttempts=1 只表示 AM 失败后不重试,与 AM 是否长期运行无关。
练习
-
在 Hadoop 集群运行
yarn application -list,观察当前所有作业。运行yarn node -list,观察所有 NM 上报的资源量。这两个命令是 YARN 集群健康检查的基本工具。 -
提交一个 MapReduce 或 Spark 作业(如 PiEstimator),在 ResourceManager Web UI 上观察作业从 SUBMITTED → ACCEPTED → RUNNING → FINISHED 的状态变化。点开作业详情,观察 AM container 在哪个 NM 上启动。
-
在
apache/hadoop源码里找到ApplicationClientProtocol.java、ApplicationMasterProtocol.java、ContainerManagementProtocol.java(hadoop-yarn-api模块)以及ResourceTracker.java(hadoop-yarn-server-common模块),观察 Client ↔ RM、AM ↔ RM、AM ↔ NM、NM ↔ RM 四个 RPC 契约。 -
思考题:如果让 YARN 支持任意数量的"资源维度"(CPU、内存、GPU、磁盘 IOPS、网络带宽),调度器的复杂度会怎么变化?为什么 Hadoop 3.x 的 DominantResourceCalculator 是为这种场景设计的?
系列导航
参考资料
- Vinod Kumar Vavilapalli et al. Apache Hadoop YARN: Yet Another Resource Negotiator. SOCC 2013.(YARN 引入论文,详细解释从 JobTracker 演化到 RM + AM + NM 三方契约的动机)
- Malte Schwarzkopf et al. Omega: Flexible, Scalable Schedulers for Large Compute Clusters. EuroSys 2013.(Google 的 shared state 调度器,与 YARN 双层调度的对比)
- Abhishek Verma et al. Large-scale Cluster Management at Google with Borg. EuroSys 2015.(Borg 调度器,YARN 双层调度的另一个对照)
- Apache Hadoop 官方文档:YARN Architecture. https://hadoop.apache.org/docs/r3.4.1/hadoop-yarn/hadoop-yarn-site/YARN.html
- Apache Hadoop 源码:
ResourceManager.java、ApplicationMasterService.java、NodeManager.java. https://github.com/apache/hadoop - Tom White. Hadoop: The Definitive Guide. O’Reilly, 4th Edition 2015. Chapter 4 “YARN” 详述了三方契约的具体行为。
