Volcano 源码阅读:一次调度怎样满足 Gang 约束

ray-demo-job-pg 已经存在,head 和两个 worker Pod 都在等待节点。假设 head 和第一个 worker 能找到位置,第二个 worker 放不下,Volcano 会把前两个 Pod 先启动吗?

这篇从 Scheduler 的一次调度循环回答这个问题。沿用 Volcano d8984501e4ad,固定一个调度主实例、一个 Queue ray-batch,不启用拓扑、分片或子组策略。三个 Ray Pod 都显式声明非零资源请求,最低成员数为三;submitter 的需求如何提前记账,留到第九篇。以下过程按源码推演,不是运行日志。

1. 一轮调度从缓存快照开始

Scheduler 不会每比较两个作业,就向 API Server 重新查询整个集群。它通过 informer 观察资源变化,维护长期存在的 SchedulerCache,包括节点、作业、成员和队列等信息。周期性调用的 runOnce 则打开一轮 Session。

Session 是本轮调度的工作视图。framework/session.gocache.Snapshot() 取得数据,在这份视图中组织后续计算。缓存快照在锁保护下构造,节点等对象被复制,避免一边试分配、一边直接改动 informer 更新中的长期状态。

scheduler.go 中打开 Session 的连续节选如下:

1
2
3
4
5
6
ssn := framework.OpenSession(pc.cache, plugins, configurations)
ssn.SetSchGateManager(pc.schGateManager)
defer func() {
framework.CloseSession(ssn)
metrics.UpdateE2eDuration(metrics.Duration(scheduleStartTime))
}()

随后的循环依次调用配置中的 action.Execute(ssn)。每个 Action 看到前面 Action 在同一 Session 内造成的变化。关闭 Session 时再执行插件收尾和相应状态更新。

长期缓存提供本轮快照。Action 在同一个 Session 中依次执行,Plugin 通过回调参与判断;提交后进入缓存与绑定路径。
长期缓存提供本轮快照。Action 在同一个 Session 中依次执行,Plugin 通过回调参与判断;提交后进入缓存与绑定路径。 查看原图

这份快照有明确的时间边界。调度器计算期间,节点可能故障,Pod 也可能被删除。Session 帮助本轮计算保持一致的工作视图,后续提交仍然需要错误处理,不能把快照里的“可以放下”视为现实已经完成。

2. 配置决定执行哪些阶段

当前基线在 pkg/scheduler/util.go 中定义的 DefaultSchedulerConf 是:

1
2
3
4
5
6
7
8
9
10
11
12
actions: "enqueue, allocate, backfill"
tiers:
- plugins:
- name: priority
- name: gang
- name: conformance
- plugins:
- name: overcommit
- name: drf
- name: predicates
- name: proportion
- name: nodeorder

本篇以这份配置作为阅读条件。部署时可以加载另一份配置,应以进程实际加载的内容为准。这里没有 preemptreclaim,所以不能因为源码中有抢占实现,就认定这个配置会主动驱逐其他作业。

Action 决定工作阶段。enqueue 检查作业能否进入调度,allocate 为有资源请求的成员寻找放置位置,backfill 处理 BestEffort 成员。在此调度器的 TaskInfo 中,BestEffort 根据启动资源请求是否为空计算,含义不是“优先级较低”。主案例的 Pod 都显式请求资源,主要跟踪前两个阶段。

Plugin 注册决策函数。先把配置中这些名字与用途对上:

插件 在相关判断点做什么
priority 按优先级比较作业或成员
gang 判断一组成员是否满足最低运行要求
conformance 需要结束已有 Pod 以腾出资源时,排除系统关键 Pod 等受保护对象;配置了插件也不表示本轮会发生驱逐
overcommit 结合集群资源总量、系数和已准入需求,判断能否再接纳作业
drf 比较作业占用各类资源的最大比例,即主导资源份额;第八篇用数字展开
predicates 排除资源不足或不满足放置规则的节点
proportion 计算和约束 Queue 的资源份额
nodeorder 给通过检查的候选节点打分

一次 allocate 会在多个位置调用这些函数,插件不会简单地各执行一次就结束。表中先列出当前阅读所需的作用,各插件还可能注册其他判断函数。

tiers 把插件组织成有先后次序的层级;框架执行到一个判断点时,会调用插件为这个位置注册的函数,这些函数就是回调。各类回调与开关有不同的组合规则,阅读时应从 Session.JobOrderFnJobEnqueueableJobReady 等入口分别追踪,不能把一份插件列表理解为所有算法结果取平均。第八篇会给出排序回调遇到首个非零比较结果就返回的例子。

3. 进入调度与满足 Gang 是两个判断

先看 actions/enqueue/enqueue.go 的核心分支:

1
2
3
4
5
if job.PodGroup.Spec.MinResources == nil || ssn.JobEnqueueable(job) {
ssn.JobEnqueued(job)
job.PodGroup.Status.Phase = scheduling.PodGroupInqueue
ssn.Jobs[job.UID] = job
}

minResources 时,enqueue 调用资源准入回调。选定配置中的 overcommit 从集群资源和已准入需求计算可接受量,proportion 还会检查 Queue 状态及相应资源上限。它们处理的是资源账目;并没有在这个分支里给三个 Pod 分别找到节点。

源码中的另一个分支也不能略过:未设置 minResources 时,这个条件直接通过,不调用此处的 JobEnqueueable。后续成员分配仍有自己的检查,因此“经过 enqueue”不能一律改述成“所有资源准入回调都已通过”。

通过后,当前 Session 中的 PodGroup phase 变成 Inqueue。这表示组进入调度流程。图里的状态更新还要经过后续写回才反映到 Kubernetes API,观察者不一定立即看到。

allocate 随后选择 Queue、JobInfo、成员 Pod 和节点。这里至少有两类不同的问题:队列还能不能拿这份资源;某一台节点能不能满足这个 Pod 的资源和放置条件。集群总共有四张空闲 GPU,如果分散在四台机器上,一个需要同机两张 GPU 的 Pod 仍然可能放不下。

资源准入、节点试分配和组就绪是不同检查。图中只展开本例的 enqueue 与 allocate,省略 BestEffort、拓扑和其他扩展分支。
资源准入、节点试分配和组就绪是不同检查。图中只展开本例的 enqueue 与 allocate,省略 BestEffort、拓扑和其他扩展分支。 查看原图

Gang 插件在 OnSessionOpen 注册就绪函数,以下为连续源码节选:

1
2
3
4
5
6
7
ssn.AddJobReadyFn(gp.Name(), func(obj interface{}) bool {
ji := obj.(*api.JobInfo)
if ji.CheckTaskReady() && ji.CheckSubJobReady() && ji.IsReady() {
return true
}
return false
})

minMember 进入内存模型后对应 JobInfo.MinAvailable。除作业总成员数外,代码还保留了 task、子组层面的检查;主案例没有额外设置这些策略。

需要留意源码里的 Ready 含义。ReadyTaskNum() 会统计下面几种调度器内部状态:

内部状态 在这条路径中表示什么
Allocated 已在调度视图中试分配资源,还没有完成 API 绑定
Binding 已进入绑定流程
Bound 已分配 Node,Pod 尚未进入运行阶段
RunningSucceeded 已观察到 Pod 运行或成功结束

IsReady() 还考虑 Pending 的 BestEffort 成员。因此这里的 JobReady 是调度器的组门槛判断,不等于 Kubernetes Pod 的 Ready 条件,更不等于 RayCluster 已经就绪。本例全部是有资源请求的新 Pod,试分配成功后处于 Allocated 的三个成员,就可以参与这个判断。

4. 先试分配,再决定是否提交

allocate 使用 Statement 记录一次尝试中的分配、驱逐等操作,供后续提交或撤销。一次 Statement.Allocate 会调整 Session 中的成员状态、节点资源账目,并通知相关插件更新本轮分配数据。这时尚未完成对 Kubernetes 的 Pod 绑定。

假定作业已经通过 enqueue,三份资源总量也足够,但第三个 Pod 要求的节点都没有余量。这样的放置规则可以用节点亲和性表达,例如要求节点带有某个 GPU 型号标签;此处指必须满足的条件。试分配过程是:head 找到节点,Session 扣除对应资源;第一个 worker 找到节点,再扣一份;第二个 worker 找不到合法位置,组仍不满足门槛。

这里还会遇到 Pipelined:如果某些资源将被释放,调度器可以先记录成员未来可能放置的位置,将这类成员计入组的等待判断。本篇没有这类未来资源可用,分配函数既无法得到 Ready,也无法得到 Pipelined 的组,便执行 stmt.Discard() 并返回。Discard 按与记录相反的顺序撤销操作,恢复这次尝试改动的内存状态和资源账目,其他作业随后可以使用这些资源。

如果三个成员都试分配成功,普通分配路径中的下面这两行会进入提交分支。这是 allocate.go 内连续代码的开头节选,后面还会更新记录并处理剩余成员:

1
2
if stmt != nil && ssn.JobReady(job) { // do not commit stmt when job is pipelined
stmt.Commit()
相同的最低成员数下,成功与失败的两次独立推演。左侧提交三份分配;右侧第三个 Pod 放置失败,撤销此前两份内存分配。
相同的最低成员数下,成功与失败的两次独立推演。左侧提交三份分配;右侧第三个 Pod 放置失败,撤销此前两份内存分配。 查看原图

IsPipelined() 会把等待未来资源的成员计入判断,上面的提交条件却仍要求 JobReady。注释强调这一点,是因为计划中的资源还没有成为当前可用资源。第八篇讨论驱逐时会再用到这种等待状态。

如果最低成员数小于总副本数,达到最低门槛后可以先提交满足要求的部分,剩余成员继续参与后续分配。Gang 检查的范围由最低要求决定,不能一概写成“整个作业的全部 Pod 一次性运行”。

5. Commit 后为什么仍可能只绑定成功一部分

Statement.Commit 遍历已记录的操作。对于 Allocate,它调用内部 allocate,出错时对该操作执行 unallocate 并记录错误。这个循环没有把多个 Pod 的 API 写入包成一个事务。

提交继续进入 SchedulerCache 的 AddBindTask。缓存会查找成员和节点,把成员更新到内部 Binding 状态、计入节点占用,然后送入 BindFlowChannel。这样后续快照能够看见正在绑定的占用,不会只因 API 里的 Pod 还没更新就再次使用同一份资源。

channel 是 Go 中用于传递数据的通道,这里把待绑定项交给后台流程。BindTask 在协程里执行 PreBind 和 Binder:PreBind 是正式绑定前的准备步骤,Binder 负责实际绑定,具体工作随实现而异。绑定结果按成员记录,可能有的成功、有的失败。

提交把分配送入长期缓存,再异步执行绑定。失败成员需要重新同步;已经成功绑定的其他成员不会因本次失败自动回到未绑定状态。
提交把分配送入长期缓存,再异步执行绑定。失败成员需要重新同步;已经成功绑定的其他成员不会因本次失败自动回到未绑定状态。 查看原图

SchedulerCache.Bind 中,失败成员会记录调度错误,执行相关 PreBind 回滚,并调用 resyncTask 重新同步。假设 head 和第一个 worker 已绑定成功,第二个 worker 因节点变化失败,前两个成功结果仍然存在;后续需要结合真实 Pod、Node 状态重新处理失败成员。

所以 Gang 能在提交前避免明显不足的分配,却不保证所有容器同一时刻启动,也不保证多个绑定请求全部成功或全部回滚。镜像拉取、卷挂载、容器启动以及 Ray 节点加入集群,发生在更后面的阶段,各自可能失败。

排查时应分别观察 PodGroup 是否准入、成员是否找到节点、绑定有没有错误、容器是否启动,以及 RayCluster 是否就绪。只看到 PodGroup 的阶段变化,还不足以说明程序已经开始执行。

本篇暂时只有一个 Queue。加入第二个作业后,即使两组都满足 Gang 条件,也可能只有一组拿得到资源。下一篇沿 Queue 和资源策略插件,分析调度器怎样做这个选择。

参考资料


Volcano 源码阅读:一次调度怎样满足 Gang 约束
https://tanxinyu.work/kuberay-volcano-07-volcano-gang/
作者
谭新宇
发布于
2026年9月15日
许可协议