Apache YARN 资源调度深度工程实战:从 RMApp/RMContainer 状态机、层级队列与 DRF 到心跳拉取式分配、抢占与生产故障全链路

执行摘要

在 Kubernetes 几乎成为"调度"代名词的今天,YARN 仍在全球数千个大数据集群里每天调度着数以亿计的 Container。它常被误解为"MapReduce 的附属品",但事实恰恰相反:YARN 是 Hadoop 2.0 里最激进的一次架构手术——把资源管理与作业生命周期管理彻底拆开,让 MapReduce 从"唯一的框架"降级为"运行在 YARN 上的一个普通应用"。Spark、Flink、Tez 后来能共享同一批机器,靠的全是这次解耦。

但 YARN 的工程价值远不止"能跑多框架"。它的调度内核有一系列非常独特、且在很多现代调度器里被重新发明的设计决策:

  • 心跳拉取式分配:ResourceManager(RM)从不主动把 Container "推送"给 NodeManager(NM),而是 NM 心跳上来"领"任务。这个看似别扭的反转,是 YARN 能撑到 5000+ 节点规模的关键。
  • 三层状态机 + 事件驱动:RMApp / RMAppAttempt / RMContainer 各自独立演进,靠 AsyncDispatcher 单线程串行化事件,用"状态机 + 事件"替代锁。
  • DRF 主导资源公平:在 CPU 与内存二维资源下,"公平"不再有自然定义,YARN 用 Dominant Resource Fairness 给出了一个可计算的答案。

本文拆开这四层,并给出生产环境的故障模式与排查清单。

一、核心抽象:把「资源」从「计算框架」里剥离

YARN 只有三个角色:

组件职责容错语义
ResourceManager全局资源仲裁、应用与 Container 状态机可 HA(基于 ZK 的 ActiveStandbyElector)
NodeManager单机代理:Container 启停、本地化、资源隔离无状态,失联则其上 Container 被标记失败
ApplicationMaster每个应用一个,向 RM 申请资源、与 NM 通信启动任务由 RM 重启,重启次数可配

关键的革命点在于 ApplicationMaster 是用户代码。RM 只认"应用"这个抽象,完全不知道里面跑的是 MapReduce 的 Map Task 还是 Flink 的 TaskManager。这带来一个重要的工程后果:AM 可以撒谎、可以低效、可以泄漏资源,而 RM 只能靠超时与心跳去兜底——这也是 YARN 生产问题的一大来源(下文第六节会展开)。

资源模型本身极度简单:Resource = <memory-mb, vcores>。在 3.x 里扩展为可插拔的 ResourceInformation,支持 GPU、FPGA 甚至自定义资源类型,但底层分配算法仍然只在"向量"上工作。

二、ResourceManager 内部:三层状态机与事件驱动

打开 RM 的源码,你会看到它几乎不用锁。取而代之的是 状态机 + 单线程事件队列:

RMStateStore (ZK/FS/LevelDB)  ←  持久化
        ↑
AsyncDispatcher  ──►  RMApp (RMAppState: NEW → SUBMITTED → ACCEPTED → RUNNING → FINISHED/FAILED/KILLED)
                  ──►  RMAppAttempt (RMAppAttemptState: SCHEDULED → ALLOCATED → LAUNCHED → RUNNING)
                  ──►  RMContainer (RMContainerState: NEW → ALLOCATED → ACQUIRED → RUNNING → COMPLETED)

这三层是刻意分离的:

  • RMApp 管"用户提交的一次作业",它的失败即作业失败;
  • RMAppAttempt 管"这次作业的第 N 次尝试",AM 挂了只让 Attempt 失败,RMApp 可以另起一个 Attempt(由 yarn.resourcemanager.am.max-attempts 控制,默认 2);
  • RMContainer 管单个容器的租约。

为什么要分这么细?因为失败域不同。一个 Task 失败不该杀掉 AM,一个 AM 崩溃不该杀掉整个作业。把这三个生命周期塞进一个对象里,就会得到一堆 if (attempt == currentAttempt) 的条件判断,以及与之相伴的竞态。

事件驱动带来的性能特征是双刃剑:AsyncDispatcher 是单线程消费,所以任何在事件处理里做的阻塞 I/O 都会让整个 RM 停摆。这就是为什么 RM 的日志里出现 SchedulerEvent 处理耗时过长时,集群会整体"卡住"——不是死锁,是一个线程被拖慢了。

// ApplicationMaster 与 RM 交互的最小骨架:注册 + 心跳式资源协商
public class MinimalAM {
    public static void main(String[] args) throws Exception {
        Configuration conf = new YarnConfiguration();
        AMRMClient<ContainerRequest> rmClient = AMRMClient.createAMRMClient();
        rmClient.init(conf);
        rmClient.start();

        // 1) 注册:告诉 RM "我还活着",并拿到本机的资源上限
        RegisterApplicationMasterResponse resp =
            rmClient.registerApplicationMaster(NetUtils.getHostname(), 0, "");
        Resource maxCap = resp.getMaximumResourceCapability();

        // 2) 申请资源:优先级 + 资源量 + 数据本地性偏好
        Priority pri = Priority.newInstance(1);
        Resource cap = Resource.newInstance(2048, 2); // 2GB / 2 vcore
        rmClient.addContainerRequest(
            new ContainerRequest(cap, new String[]{"hdfs-host-a"}, new String[]{"/rack1"}, pri));

        // 3) 心跳拉取:allocate(progress) 既上报进度,也"领取"已分配的 Container
        while (true) {
            AllocateResponse alloc = rmClient.allocate(0.5f);   // progress = 0.5
            for (Container c : alloc.getAllocatedContainers()) {
                // 拿到 Container 后自己去找对应 NM 启动进程
                launchOnNodeManager(conf, c);
            }
            for (ContainerStatus s : alloc.getCompletedContainersStatuses()) {
                // 处理失败/被抢占的容器:被抢占时 exitStatus = -100
            }
            Thread.sleep(1000);
        }
    }
}

注意 allocate() 这个 RPC 的双重语义:它既是心跳,也是分配结果的拉取通道,还是资源请求的增量提交口。一次 RPC 干三件事,是 YARN 为了减少 RM 端 RPC 压力做的刻意的胖接口设计。

三、调度器:层级队列与 DRF 主导资源公平

YARN 提供三种调度器,但生产上只有 Capacity Scheduler 真正被大规模使用。

<!-- capacity-scheduler.xml:层级队列是容量保证的骨架 -->
<property>
  <name>yarn.scheduler.capacity.root.queues</name>
  <value>prod,ad-hoc</value>
</property>
<property>
  <name>yarn.scheduler.capacity.root.prod.capacity</name>
  <value>70</value>
</property>
<property>
  <name>yarn.scheduler.capacity.root.ad-hoc.capacity</name>
  <value>30</value>
</property>
<!-- user-limit-factor > 1 才允许单用户借用超出自身份额的资源 -->
<property>
  <name>yarn.scheduler.capacity.root.ad-hoc.user-limit-factor</name>
  <value>4</value>
</property>
<!-- 队列最大容量:限制 ad-hoc 在集群空闲时能膨胀到的上限 -->
<property>
  <name>yarn.scheduler.capacity.root.ad-hoc.maximum-capacity</name>
  <value>60</value>
</property>

这里有个常被误解的语义:capacity 是"保证的容量",不是"限制的容量"。集群空闲时 ad-hoc 可以涨到 maximum-capacity(60%)而不是被卡在 30%。把 capacity 当成硬限制来配,是集群利用率长期只有 30% 的典型原因。

DRF:二维资源下的"公平"怎么算

当资源是标量(只有内存)时,公平 = 份额相等。但在 (CPU, 内存) 二维下没有自然序。DRF 的定义是:最大化所有用户中最小的"主导份额"(dominant share)。

def drf_next(user_demands, user_alloc, total):
    """返回下一个应该拿到资源的用户。demand/alloc 均为 [cpu, mem] 向量"""
    best, best_share = None, None
    for u, demand in user_demands.items():
        alloc = user_alloc[u]
        # 主导资源 = 该用户各维度占比中最大的那一维
        shares = [alloc[d] / total[d] for d in range(len(demand))]
        dom = max(shares)
        # 若主导资源已达需求,说明该用户已饱和,跳过
        if any(alloc[i] >= demand[i] for i in range(len(demand))):
            continue
        if best_share is None or dom < best_share or (dom == best_share and tie_break(u, best)):
            best, best_share = u, dom
    return best

DRF 有几个在真实集群里非常重要的性质:

  • 策略证明性(strategy-proof):虚报需求不会让你拿到更多,只会拉低自己的主导份额增速。
  • 激励共享:一个只用内存的用户和一个只用 CPU 的用户可以完美互补,两者主导份额都能接近 100%。
  • 与"资产公平"的区别:你无法把内存换算成 CPU 再比较,DRF 不做这种换算,避免了人为定价带来的扭曲。

代价是 DRF 的计算复杂度高于简单的 max-min fairness,所以在超大规模队列下,FairScheduler 的 drf 策略比 fair 策略更吃 CPU。Capacity Scheduler 在 3.x 里默认用 绝对资源配置的 max-min,只有显式开启才会走 DRF。

四、心跳拉取式分配:为什么 RM 从不"推送"

这是 YARN 与 Kubernetes 调度最本质的差别之一。

Kubernetes 的 kubelet 通过 Watch 监听 Pod 绑定事件,调度器是主动写入的。YARN 反过来:RM 的调度器把分配结果放进 SchedulerApplicationAttempt 的待领取队列,等到 NM 下一次心跳才被取走。

NM ──(1) nodeHeartbeat: 汇报本节点资源与 Container 状态──► RM
RM ──(2) 响应里携带: 可启动的 Container 列表 / 需清理的 Container──► NM

这个反转的好处是:

  1. RM 不需要维护到每个 NM 的连接状态,NM 失联就是心跳超时(默认 10 分钟),处理极其简单;
  2. 天然的背压:NM 正在启动 20 个 Container 时可以不领新的,避免单机过载;
  3. 可重放:心跳响应丢了也没关系,下一次心跳会重新带上未确认的分配。

代价是分配延迟直接被心跳周期绑架。默认心跳 1 秒,理论上一个 Container 从申请到启动至少 1~2 秒。对于短作业(比如几十个 Task 的 Spark SQL 查询),这个延迟占比很高。生产上常见的优化是把 yarn.nodemanager.heartbeat.interval-ms 调到 100~300ms,但心跳频率与 RM 的 CPU 消耗成正比:5000 个 NM × 每秒 10 次心跳 = 每秒 5 万次 RPC,RM 的 AsyncDispatcher 会成为瓶颈。这就是为什么 YARN 在超大规模下要走 Federation + Router:把集群切成多个子集群,每个 RM 只管一部分 NM。

五、Container 生命周期与本地化

一个 Container 从"分配"到"运行"要经历:ALLOCATED → ACQUIRED → (LOCALIZING) → RUNNING → COMPLETED。

中间最容易出问题的是 本地化(Localization):Container 依赖的 jar、配置文件要先从 HDFS 下载到 NM 本地磁盘,由 ResourceLocalizationService 用独立的线程池完成。这里有两层缓存语义:

  • PUBLIC 资源:所有用户共享,下载到 ${yarn.nodemanager.local-dirs}/filecache;
  • PRIVATE 资源:按用户隔离,下载到 usercache/<user>/filecache。

生产事故高发点:local-dirs 磁盘被打满。NM 有 DiskHealthChecker 定期检查,一旦判定磁盘损坏会把整个节点标记为 unhealthy 并从 RM 摘除。但如果只是"接近满"而未触发阈值,你会看到大量 Container 在本地化阶段超时失败,报错却是含糊的 Container exited with a non-zero exit code 1。排查时第一件事应该是 df -h 看 local-dirs。

资源隔离方面,Linux 上默认用 cgroups(内存硬性限制 + CPU shares)。一个经典陷阱:yarn.nodemanager.resource.memory-mb 配得比物理内存大,加上系统预留不足,会触发内核 OOM Killer 把整个 NM 进程杀掉——表现为一批节点同时失联,而不是单个 Container 失败。

六、抢占、预留与节点分区

抢占(Preemption) 是容量保证的最后一道防线。当 prod 队列的保证容量被 ad-hoc 占用时,调度器会计算"该从哪些 Container 手里收回资源",并通过 RMContainer 状态机触发 KILL。AM 会收到退出码为 -100 的完成状态——AM 必须正确处理这个信号并重新申请,否则作业会直接失败。这是自研 AM 最常见的 bug。

预留(Reservation) 解决的是"大作业饥饿":一个需要 32GB 的 Container 在碎片化集群里可能永远等不到连续资源。调度器会在某个节点上为它"预留",让该节点后续的小 Container 分配跳过。代价是预留会显著降低集群吞吐,所以默认开启但可通过 yarn.resourcemanager.scheduler.monitor.enable 相关参数调节。

节点标签(Node Label / Partition) 用于把集群切成逻辑分区(如 GPU 区、高内存区),队列可以声明自己能访问哪些分区。配置要点是:队列的 accessible-node-labels 与 capacity 要成对配置,只在队列上加了标签却没给该标签分配容量,应用会永远 pending。

七、生产故障模式与排查清单

现象根因排查与处置
应用长期 ACCEPTED,从不 RUNNING队列容量满 / AM 资源请求超过 maxAllocationyarn application -appStates ACCEPTED -list;检查 maximum-allocation-mb
AM 反复重启最终 FAILEDAM 自身 OOM 或依赖下载失败yarn logs -applicationId <id>;看 am.max-attempts 是否过小
大量 Container exit code -100被抢占确认抢占开关与队列容量;AM 需实现重新申请逻辑
RM 响应变慢、UI 卡顿AsyncDispatcher 单线程被阻塞,常见于超大 FBR 式的全量心跳或队列数过多抓取 jstack 看 AsyncDispatcher 线程栈;控制队列层级深度
节点批量失联local-dirs 满 / 磁盘损坏 / cgroup OOMyarn node -list -all;df -h;dmesg 查内核 OOM 记录
改了队列配置不生效未刷新yarn rmadmin -refreshQueues(注意:部分参数仍需重启 RM)

一条经验:80% 的"YARN 慢"问题不在 YARN 本身,而在 AM 的资源请求模式。一个每次只申请 1 个 Container、等它跑完再申请下一个的 AM,会被心跳延迟放大成几十倍的串行等待。正确的做法是一次性申请一批、用 allocate() 流水线化。

八、结论

YARN 的设计取舍可以用一句话概括:用"分配延迟"换"规模与简单性"。

  • 用心跳拉取换掉了 RM 对节点的连接状态管理,代价是秒级分配延迟;
  • 用三层状态机 + 单线程事件换掉了复杂的锁,代价是单点阻塞风险;
  • 用AM 是用户代码换掉了框架无关性,代价是 AM 的质量直接决定集群体验。

它与 Kubernetes 的分野也在这里:K8s 优化的是"声明式收敛 + 秒级以下的调度延迟",YARN 优化的是"万级节点下的吞吐与容错"。当你的工作负载是长生命周期、大批量、可容忍秒级启动的数据处理任务时,YARN 的心跳模型反而比 Watch 模型更抗抖动、更容易水平扩展。

真正要在生产里用好 YARN,记住三件事:capacity 是保底不是上限;AM 必须处理抢占信号;分配延迟是心跳周期的函数,调它之前先算 RM 的 RPC 预算。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部