Scheduler:Pod 的节点选择¶
Kubernetes 调度器是一个控制面进程,负责将 Pod 指派到节点上。调度器根据约束和可用资源,为调度队列中的每个 Pod 找出所有可以运行它的节点,再对这些节点排序,将 Pod 绑定到一个合适的节点。同一集群可以运行多个调度器;kube-scheduler 是 Kubernetes 的参考实现。
Scheduling Framework¶
Scheduling Framework 是 Kubernetes 调度器的可插拔架构。它由一组直接编译进调度器的插件 API 组成,使大多数调度功能可以通过插件实现,同时让调度核心保持轻量、易于维护。更多设计细节见 Scheduling Framework 设计提案。
整体架构¶
下图展示了 Pod 的调度上下文以及 Scheduling Framework 暴露的接口。一个插件可以实现多个接口,用于处理更复杂或需要保存状态的调度逻辑。部分接口对应可以通过调度器配置设置的扩展点。

图:Scheduling Framework 扩展点。来源:Kubernetes Scheduling Framework,CC BY 4.0。
一次 Pod 调度包含三个阶段:
- 进入调度队列:
PreEnqueuePlugin决定 Pod 能否进入活动队列,QueueSortPlugin决定活动队列中的先后顺序。 - 调度阶段(Scheduling cycle):Scheduler 串行处理 Pod,从队列中取出一个 Pod,过滤不符合条件的 Node,再对剩余 Node 评分,选出目标 Node。
- 绑定阶段(Binding cycle):Scheduler 可以并发执行不同 Pod 的绑定阶段,完成卷绑定等准备工作,再将 Pod 绑定到选定的 Node,并通过 API Server 保存结果。
Kubernetes v1.36.4 的默认 Profile 通过 MultiPoint 启用插件,同一个插件会进入它所实现的所有扩展点。下表列出普通 Pod 调度路径中的扩展点、辅助接口及默认插件。
| 名称 | 类型 | 所属阶段 | 作用 | 默认插件 |
|---|---|---|---|---|
PreEnqueue |
扩展点 | 进入调度队列 | 检查 Pod 是否具备入队条件 | SchedulingGates、DynamicResources、DefaultPreemption |
QueueSort |
扩展点 | 进入调度队列 | 决定 activeQ 中 Pod 的取出顺序 |
PrioritySort |
PreFilter |
扩展点 | 调度阶段 | 为当前 Pod 准备共享计算结果,也可以缩小候选 Node 范围 | NodeAffinity、NodePorts、NodeResourcesFit、VolumeRestrictions、NodeVolumeLimits、VolumeBinding、VolumeZone、PodTopologySpread、InterPodAffinity、DynamicResources、NodeDeclaredFeatures |
Filter |
扩展点 | 调度阶段 | 逐个检查 Node 是否满足硬约束 | NodeUnschedulable、NodeName、TaintToleration、NodeAffinity、NodePorts、NodeResourcesFit、VolumeRestrictions、NodeVolumeLimits、VolumeBinding、VolumeZone、PodTopologySpread、InterPodAffinity、DynamicResources、NodeDeclaredFeatures |
PostFilter |
扩展点 | 调度阶段 | 没有可行 Node 时尝试释放资源或执行抢占 | DynamicResources、DefaultPreemption |
PreScore |
扩展点 | 调度阶段 | 准备全部候选 Node 共用的评分数据 | TaintToleration、NodeAffinity、NodeResourcesFit、VolumeBinding、PodTopologySpread、InterPodAffinity、NodeResourcesBalancedAllocation |
Score |
扩展点 | 调度阶段 | 为每台可行 Node 计算分数 | TaintToleration、NodeAffinity、NodeResourcesFit、VolumeBinding、PodTopologySpread、InterPodAffinity、DynamicResources、NodeResourcesBalancedAllocation、ImageLocality |
Reserve |
扩展点 | 调度阶段 | 记录可回滚的临时状态,例如卷预绑定和 DRA 设备分配 | VolumeBinding、DynamicResources |
Permit |
扩展点 | 调度阶段 | 允许、等待或拒绝进入绑定阶段 | 无 |
PreBind |
扩展点 | 绑定阶段 | 在提交绑定前完成卷或设备相关操作 | VolumeBinding、DynamicResources |
Bind |
扩展点 | 绑定阶段 | 向 API Server 提交 Pod 与 Node 的绑定关系 | DefaultBinder |
PostBind |
扩展点 | 绑定阶段 | 在绑定成功后运行通知类操作 | 无 |
EnqueueExtensions |
辅助接口 | 入队与重试 | 注册可能改变插件失败结果的集群事件,并通过 QueueingHint 判断 Pod 是否需要重新入队 |
SchedulingGates、NodeUnschedulable、NodeName、TaintToleration、NodeAffinity、NodePorts、NodeResourcesFit、VolumeRestrictions、NodeVolumeLimits、VolumeBinding、VolumeZone、PodTopologySpread、InterPodAffinity、DynamicResources、DefaultPreemption、NodeDeclaredFeatures |
PreFilterExtensions |
辅助接口 | 调度阶段 | 同步更新 PreFilter 的预计算状态 | InterPodAffinity、PodTopologySpread、VolumeRestrictions |
ScoreExtensions |
辅助接口 | 调度阶段 | 所有 Node 完成 Score 后,通过 NormalizeScore 把同一插件的结果调整到 0–100 |
TaintToleration、NodeAffinity、PodTopologySpread、InterPodAffinity、DynamicResources |
Unreserve |
辅助接口 | 调度或绑定失败 | 撤销 Reserve 阶段记录的临时状态 | VolumeBinding、DynamicResources |
默认插件列表见 getDefaultPlugins,in-tree 插件构造函数注册在 NewInTreeRegistry。v1.36.4 默认启用 DynamicResourceAllocation 和 NodeDeclaredFeatures;Alpha 的 GangScheduling 与 TopologyAwareWorkloadScheduling 默认关闭,因此未列入普通 Pod 的默认路径。
Pod 调度过程¶
下面分别说明这三个阶段的处理过程。
进入调度队列
- 观察 Pod:Scheduler 观察 API Server 中的 Pod 变化。已经设置
spec.nodeName的 Pod 表示它已被分配到某个 Node,Scheduler 将它记录到本地 Scheduler Cache,用于计算 Node 的资源占用;尚未指定 Node 且由当前 Scheduler 负责的 Pod 进入 SchedulingQueue,等待调度。 - 检查入队:PreEnqueue 检查 Pod 能否进入可调度队列。首次通过检查的 Pod 进入
activeQ;检查未通过时进入unschedulablePods。 - 等待重试:调度未能继续的 Pod 进入
unschedulablePods,等待可能改变结果的集群事件。Unschedulable表示当前条件下无法完成调度;Pending表示已经选出目标 Node,但插件需要等待外部组件根据调度结果完成操作,Pod 暂时不能进入绑定阶段。相关事件触发重新入队后,PendingPod 直接进入activeQ;UnschedulablePod 在退避尚未结束时进入podBackoffQ,退避结束后进入activeQ。
调度阶段(Scheduling cycle)
- 取出 Pod:Scheduler 从 SchedulingQueue 取出 Pod,读取当前的 Node 和资源状态。
- 筛选 Node:PreFilter 准备筛选所需的状态,Filter 排除不满足硬约束的 Node。没有可行 Node 时,PostFilter 尝试抢占或其他补救措施;仍无法选出 Node 时,由 FailureHandler 记录失败原因,并将 Pod 交回 SchedulingQueue 等待重试。
- 选择 Node:存在多个可行 Node 时,PreScore 准备评分所需的状态,Score 选择得分最高的 Node;只有一个可行 Node 时直接使用该 Node。
- 记录选择:Scheduler 在本地 Cache 中记录目标 Node 和资源占用。Reserve 保存可回滚的临时状态,例如卷的预绑定信息;Permit 决定继续绑定、等待其他条件满足,或者拒绝本次选择。Reserve 或 Permit 失败时,Scheduler 先运行
Unreserve撤销插件状态,再调用ForgetPod删除 Cache 中的 Assumed Pod。
绑定阶段(Binding cycle)
- 检查 PreBind 插件:Scheduler 逐个运行
PreBindPreFlight,记录需要跳过的插件,以及哪些相邻插件可以在后续PreBind中并行执行。预检本身按顺序执行。 - 等待 Permit 放行:Scheduler 调用 WaitOnPermit。只有 Permit 插件在调度阶段返回
Wait时,这一步才会阻塞;否则立即进入PreBind。等待期间保留目标 Node 的本地资源占用,插件拒绝或等待超时会撤销本次选择。 - 准备绑定:PreBind 完成卷绑定等准备工作。
- 提交绑定:Bind 负责提交绑定结果。默认
DefaultBinder请求 Pod 的binding子资源,API Server 将目标 Node 写入spec.nodeName。 - 确认结果:PostBind 在绑定成功后运行。Scheduler 随后观察到已绑定的 Pod,并用新的 API 对象确认本地 Cache 中的 Assumed 状态。
进入调度队列¶
Pod 进入调度队列前,PreEnqueue 检查它是否已具备开始调度的条件。Pod 入队后,QueueSort 决定 activeQ 中的取出顺序;调度失败的 Pod 重新入队时,EnqueueExtensions 和 QueueingHint 决定哪些集群事件可以触发重试。
PreEnqueue¶
PreEnqueue 检查 Pod 是否已经具备开始调度的前置条件。例如,Pod 的 spec.schedulingGates 仍有条目时,内置 SchedulingGates 插件会阻止它进入可调度队列;控制器移除全部 gate 后,Pod 才会重新参加调度。
| 默认插件 | 具体作用 |
|---|---|
SchedulingGates |
Pod 仍包含 spec.schedulingGates 时阻止入队,全部 gate 被移除后允许重新检查 |
DynamicResources |
检查 Pod 引用的 DRA ResourceClaim 是否存在并具备开始调度的条件 |
DefaultPreemption |
启用异步抢占时,阻止正在执行抢占操作的 Pod 或 PodGroup 重复进入调度队列 |
Kueue 通过 schedulingGates 控制 Pod 的调度时机
schedulingGates 适合由外部控制器决定 Pod 何时具备调度条件的场景。该字段只能在 Pod 创建时由客户端设置,或由 Admission Webhook 写入;Pod 创建后只能移除已有 gate,不能再添加。具体约束见 Kubernetes 的 Pod Scheduling Readiness。
下面的 Pod 设置了两个 scheduling gate。创建后,它保持 SchedulingGated 状态;两个 gate 全部移除后,Scheduler 才会把它放入可调度队列。
apiVersion: v1
kind: Pod
metadata:
name: test-pod
spec:
schedulingGates:
- name: example.com/foo
- name: example.com/bar
containers:
- name: pause
image: registry.k8s.io/pause:3.6
Kueue 的 Pod 集成通过 Mutating Webhook 注入 kueue.x-k8s.io/admission gate。对应的 Workload 等待配额、ResourceFlavor 或 AdmissionCheck 时,Pod 保持 SchedulingGated,不会进入 activeQ。Workload 完成准入后,Kueue 移除自己管理的 gate;Pod 没有其他 gate 时,SchedulingGates 插件允许它进入调度队列。
Kueue 直接管理 Pod 时使用 schedulingGates 控制调度时机。Kueue 管理 Kubernetes Job时,通常通过 Job 的 spec.suspend 控制启动时机。
SchedulingQueue 准备将 Pod 放入 activeQ 或 backoffQ 时,依次调用当前 Profile 中的 PreEnqueuePlugin.PreEnqueue。全部插件返回 Success 后,Pod 才会进入目标队列;任一插件返回失败状态时,调用立即停止,队列记录拒绝它的插件,并将 Pod 转入 unschedulablePods。这条分流可以在 moveToActiveQ 和 moveToBackoffQ 中看到。此后只有该插件关注的事件或通配事件才会触发重新检查。
QueueSort¶
activeQ 中可能同时存在多个 Pod,QueueSort 决定 Scheduler 下一次取出哪一个。
| 默认插件 | 具体作用 |
|---|---|
PrioritySort |
优先级较高的 Pod 先取出;优先级相同时,入队更早的 Pod 先取出 |
| Pod | Priority | 入队时间 | 取出次序 |
|---|---|---|---|
| A | 100 | 10:02 | 2 |
| B | 50 | 10:00 | 3 |
| C | 100 | 10:01 | 1 |
Pod C 和 Pod A 的 Priority 相同,入队更早的 Pod C 先被取出;Pod B 的 Priority 较低,最后取出。因此,取出顺序为 C、A、B。
QueueSortPlugin.Less 每次比较两个 QueuedPodInfo,activeQ 的堆结构使用比较结果维护上述顺序。newQueuedPodInfo 创建队列对象时将 Timestamp 设置为当前时间;一次调度失败后,AddUnschedulableIfNotPresent 在 Pod 重新入队前刷新该字段。InitialAttemptTimestamp 在 Pod 第一次进入 activeQ 时设置,之后保持不变,用于计算端到端调度延迟。
一个 Scheduler 进程中的所有 Profile 共享同一条调度队列,因此必须使用相同的 QueueSort 插件和配置。profile.Map 的校验逻辑会拒绝不一致的配置。
EnqueueExtensions¶
插件拒绝 Pod 后,SchedulingQueue 需要知道哪些集群变化可能改变上次的调度结果。EnqueueExtensions 让插件注册自己关注的事件,避免无关的 Node、Pod 或 PVC 变化反复触发调度。
| 默认插件 | 具体作用 |
|---|---|
SchedulingGates |
关注 Pod 的 scheduling gate 是否已全部移除 |
NodeUnschedulable |
关注新 Node、Node 的 unschedulable taint 和目标 Pod 的 Toleration 变化 |
NodeName |
关注与 spec.nodeName 匹配的新 Node |
TaintToleration |
关注 Node Taint 和目标 Pod Toleration 的变化 |
NodeAffinity |
关注新 Node,以及 Node label 是否从不匹配变为匹配 |
NodePorts |
关注新 Node,以及占用冲突 hostPort 的 Pod 是否被删除 |
NodeResourcesFit |
关注 Node 可分配资源增加、已调度 Pod 删除或原地缩容 |
VolumeRestrictions |
关注 Pod 删除或 PVC 创建后,原有卷冲突是否消失 |
NodeVolumeLimits |
关注 CSINode、PVC、Pod 和 VolumeAttachment 的变化是否释放卷连接额度 |
VolumeBinding |
关注 Node、PV、PVC、StorageClass 和 CSI 容量信息是否出现可用的卷绑定方案 |
VolumeZone |
关注 Node、PV、PVC 和 StorageClass 的变化是否满足卷拓扑要求 |
PodTopologySpread |
关注 Pod 或 Node 的变化是否改善拓扑分布 |
InterPodAffinity |
关注 Pod 或 Node 的变化是否满足 Pod Affinity 或 Anti-Affinity |
DynamicResources |
关注 ResourceClaim、ResourceSlice、DeviceClass 等 DRA 对象是否提供可用设备 |
DefaultPreemption |
启用异步抢占时,关注被抢占 Pod 删除后抢占操作是否完成 |
NodeDeclaredFeatures |
关注 Node 声明的 Feature 或 Pod 所需 Feature 是否发生变化 |
Scheduler 启动时调用一次 EnqueueExtensions.EventsToRegister,并在运行期间沿用这份注册结果。该方法返回一组 ClusterEventWithHint:
Event:指定插件关注的资源和动作,例如 Node label 更新。QueueingHintFn:接收上次调度失败的 Pod,以及事件涉及的新旧对象,判断这次变化是否值得让 Pod 重新入队。
例如,Pod 要求 Node 带有 storage=ssd label,但当前没有 Node 满足条件。Node label 更新后,NodeAffinity 的 QueueingHint 回调会重新检查该 Node:
- Node 从不匹配变为匹配时返回
Queue,Pod 可以重新入队。 - 更新后的 Node 仍不匹配时返回
QueueSkip,Pod 继续留在unschedulablePods。
Queue 只说明这次事件可能改变调度结果,不保证 Pod 已经具备调度条件,也不直接决定它进入哪个队列。
收到 Queue 后,SchedulingQueue 再根据上次的插件状态选择入队方式:
- 返回
Queue的插件此前将 Pod 标记为Pending时,Pod 跳过退避并进入activeQ。 - 返回
Queue的插件此前将 Pod 标记为Unschedulable时,Pod 根据剩余退避时间进入podBackoffQ或activeQ。
Kubernetes v1.36.4 默认启用 SchedulerPopFromBackoffQ。Pod 从 unschedulablePods 写入 activeQ 或 podBackoffQ 前会重新运行 PreEnqueue;检查未通过时,Pod 继续留在 unschedulablePods。
PreEnqueue、PreFilter、Filter、Reserve 或 Permit 插件没有实现 EnqueueExtensions 时,Scheduler 会为它注册默认的全部已知集群事件。这样可以避免 Pod 因缺少事件注册而一直留在 unschedulablePods,代价是部分无关事件也会触发重新入队。EventsToRegister 返回错误时,Scheduler 无法启动。
调度阶段(scheduling cycle)¶
调度阶段为一个 Pod 选择目标 Node,并把选择结果写入 Scheduler 本地状态。CycleState 是本轮调度中插件共享计算结果的键值存储;每轮创建一份,轮次结束后不再复用。
PreFilter¶
多个 PreFilter 插件按照配置顺序串行执行,每个插件在当前 Pod 的调度周期中运行一次。插件接收当前 Pod 和全部候选 Node,可以把供后续 Filter 使用的计算结果写入 CycleState。PreFilterPlugin.PreFilter 还可以返回 PreFilterResult,把后续检查限制在一组 Node 中。
| 默认插件 | 具体作用 |
|---|---|
NodeAffinity |
解析 nodeSelector、必需 Node Affinity 和 Scheduler 额外配置的 Node Affinity |
NodePorts |
收集 Pod 请求的 hostPort,供 Filter 检查端口冲突 |
NodeResourcesFit |
计算 Pod 的 CPU、内存、临时存储、Pod 数量和扩展资源请求 |
VolumeRestrictions |
收集 ReadWriteOncePod PVC 等卷信息,准备卷冲突检查 |
NodeVolumeLimits |
识别 Pod 使用的 CSI 卷,准备计算每个 Driver 的卷连接数量 |
VolumeBinding |
解析 Pod 的 PVC,并区分已绑定卷和需要延迟绑定的卷 |
VolumeZone |
读取已绑定 PV 的 zone 和 region 信息 |
PodTopologySpread |
计算硬性拓扑分布约束及各拓扑域中的匹配 Pod 数量 |
InterPodAffinity |
计算必需 Pod Affinity 和 Anti-Affinity 使用的拓扑计数 |
DynamicResources |
解析 DRA ResourceClaim,准备每个 Node 的设备分配状态 |
NodeDeclaredFeatures |
推导 Pod 所需的 Node Feature 集合 |
PreFilter 返回 Unschedulable 时,本轮不再检查 Node;返回 Skip 时,与它配套的 Filter 和 PreFilterExtensions 在本轮一起跳过。
CycleState 中的预计算结果反映了执行 PreFilter 时的 Node 状态。Scheduler 后续可能修改 NodeInfo 副本中的 Pod 集合,再使用修改后的状态运行 Filter。有些 PreFilter 结果取决于 Node 上已有的 Pod,例如 Pod 亲和性、拓扑分布和卷冲突的计数;如果只修改 NodeInfo,Filter 读取的 Node 状态和预计算结果就会不一致。
PreFilterExtensions 提供 AddPod 和 RemovePod 两个方法。Scheduler 在 NodeInfo 副本中加入 Pod 后调用 AddPod,移除 Pod 后调用 RemovePod,由相应的 PreFilter 插件同步更新 CycleState。
| 默认插件 | 具体作用 |
|---|---|
InterPodAffinity |
Pod 集合变化时更新 Affinity 和 Anti-Affinity 的拓扑计数 |
PodTopologySpread |
Pod 集合变化时更新各拓扑域中的匹配 Pod 数量 |
VolumeRestrictions |
Pod 集合变化时更新 ReadWriteOncePod PVC 的使用计数 |
这两个方法主要用于以下两类调度计算:
- Nominated Pod 检查:对某个 Node 运行
Filter前,Scheduler 会把已经提名该 Node、且优先级不低于当前 Pod 的 Pod 加入NodeInfo副本,并调用AddPod。 - 抢占模拟:没有可行 Node 时,
PostFilter阶段的DefaultPreemption会修改NodeInfo副本并重新运行Filter。它先移除可能被抢占的低优先级 Pod,并调用RemovePod更新预计算结果;当前 Pod 可以放入后,再逐个加回低优先级 Pod 并调用AddPod,判断哪些 Pod 可以保留。这种增量更新避免了每次修改后重新执行完整的PreFilter。DefaultPreemption 的模拟过程展示了这些调用。
这些操作只更新 Scheduler 的本地模拟状态,不会直接增删集群中的 Pod。
Filter¶
Filter 检查候选 Node 是否满足 Pod 的硬约束。FilterPlugin.Filter 每次接收 Pod 和一台候选 Node 的 NodeInfo;同一 Node 通过全部 Filter 插件后,才会进入可行节点集合。
| 默认插件 | 具体作用 |
|---|---|
NodeUnschedulable |
检查 node.spec.unschedulable;只有 Pod 容忍对应状态时才允许使用该 Node |
NodeName |
Pod 指定了 spec.nodeName 时,只允许名称相同的 Node |
TaintToleration |
排除带有 Pod 无法容忍的 NoSchedule 或 NoExecute Taint 的 Node |
NodeAffinity |
检查 Node label 和名称是否满足 nodeSelector、必需 Node Affinity 以及 Scheduler 额外配置的 Node Affinity |
NodePorts |
检查 Pod 请求的 hostPort 是否与 Node 上已有 Pod 冲突 |
NodeResourcesFit |
检查 Node 的可分配 CPU、内存、临时存储、Pod 数量和扩展资源是否满足请求 |
VolumeRestrictions |
检查 ReadWriteOncePod PVC 和部分块存储类型是否与 Node 上已有 Pod 的卷冲突 |
NodeVolumeLimits |
按 CSI Driver 统计 Node 已连接和本次需要连接的卷,检查是否超过 CSINode 声明的上限 |
VolumeBinding |
检查已绑定 PV 的 Node Affinity;对于延迟绑定 PVC,寻找能够在当前 Node 上绑定或动态供应的卷 |
VolumeZone |
检查已绑定 PV 的 zone 和 region 是否与 Node 标签一致 |
PodTopologySpread |
检查 DoNotSchedule 约束的 maxSkew 和 minDomains,排除会破坏硬性拓扑分布的 Node |
InterPodAffinity |
检查必需的 Pod Affinity 和 Anti-Affinity,确认当前 Node 的拓扑域满足约束 |
DynamicResources |
检查已分配设备能否在当前 Node 使用,并尝试为尚未分配的 ResourceClaim 找到可用设备 |
NodeDeclaredFeatures |
检查 Node 声明的 Feature 是否包含 Pod 所需的全部 Feature |
常见 Filter 包括资源容量、Node affinity、taint、端口冲突和卷拓扑检查。插件的返回状态决定当前调度尝试如何继续:
- Success:当前插件接受该 Node,继续运行后续 Filter。
- Unschedulable:当前 Node 不满足条件,但条件可能通过集群状态变化或抢占得到解决。
- UnschedulableAndUnresolvable:当前 Node 的条件无法通过抢占解决,例如必需 label 不匹配。
- Error:插件执行失败,终止本次调度,并让 Pod 进入错误退避路径。
Filter 接口接收 NodeInfo。NodeInfo 聚合 Node 对象、已经分配和 Assumed 的 Pod、资源请求、端口、PVC 与镜像状态,插件使用传入的 NodeInfo 完成当前节点判断。
PostFilter¶
当 PreFilter 或 Filter 没有留下可行 Node 时,Scheduler 调用 PostFilter,尝试通过抢占等方式为当前 Pod 腾出资源。
| 默认插件 | 具体作用 |
|---|---|
DynamicResources |
尝试清理尚未被其他 Pod 使用的 DRA ResourceClaim 分配,让后续调度重新选择设备 |
DefaultPreemption |
选择候选 Node 和需要移除的低优先级 Pod,并为当前 Pod 设置 nominated Node |
默认的抢占插件会跳过 UnschedulableAndUnresolvable 的 Node,并检查移除哪些低优先级 Pod 后能够容纳当前 Pod。找到候选方案后,它发起抢占并返回 nominated Node。当前调度周期仍然结束,Pod 会在被抢占 Pod 退出后重新调度,不会从 PostFilter 直接进入绑定。
PostFilterPlugin.PostFilter 只在 PreFilter 或 Filter 返回不可调度结果时运行。没有插件能够改善结果时,Scheduler 记录 FailedScheduling Event,并将 Pod 放回不可调度队列或退避队列。
PreScore¶
通过 Filter 的 Node 多于一台时,PreScore 会先处理整组可行 Node,为后续评分准备共享状态。例如,拓扑相关插件可以先统计各拓扑域中的 Pod 数量,再由 Score 使用这些结果逐台打分。PreScorePlugin.PreScore 的计算结果保存在本轮 CycleState 中。
| 默认插件 | 具体作用 |
|---|---|
TaintToleration |
收集 Pod 对 PreferNoSchedule Taint 的 Toleration,供 Score 统计无法容忍的 Taint |
NodeAffinity |
解析首选 Node Affinity 条件 |
NodeResourcesFit |
计算参与资源利用率评分的 Pod 资源请求 |
VolumeBinding |
判断本轮是否需要根据可动态供应卷的容量进行评分 |
PodTopologySpread |
统计软性拓扑约束下各拓扑域中的匹配 Pod 数量 |
InterPodAffinity |
计算首选 Pod Affinity 和 Anti-Affinity 的拓扑权重 |
NodeResourcesBalancedAllocation |
计算参与资源均衡评分的 Pod 资源请求 |
PreScore 返回 Skip 时,与它配套的 Score 插件在本轮一起跳过。Filter 只剩一台可行 Node 时,schedulePod 直接选择该 Node,不运行 PreScore 和 Score。
Score¶
Filter 留下多台可行 Node 时,Scheduler 运行 Score 插件为这些 Node 打分,再根据加权总分选择目标 Node。ScorePlugin.Score 每次接收 Pod 和一台 Node 的 NodeInfo,并返回该插件为这台 Node 计算的分数;分数越高,Node 越符合该插件的偏好。
| 默认插件 | 默认权重 | 具体作用 |
|---|---|---|
TaintToleration |
3 | 无法容忍的 PreferNoSchedule Taint 越少,Node 得分越高 |
NodeAffinity |
2 | 累加 Node 命中的首选 Node Affinity 条件权重 |
NodeResourcesFit |
1 | 按配置的 LeastAllocated、MostAllocated 或 RequestedToCapacityRatio 策略计算资源利用率得分 |
VolumeBinding |
1 | 根据可绑定卷的容量与 PVC 请求量计算容量匹配度;不需要容量评分时跳过 |
PodTopologySpread |
2 | 根据各拓扑域已有的匹配 Pod 数量评分,倾向于降低分布偏差 |
InterPodAffinity |
2 | 根据首选 Pod Affinity 和 Anti-Affinity 的权重为 Node 评分 |
DynamicResources |
2 | 启用 DRA prioritized list 时,根据候选设备分配结果为 Node 评分 |
NodeResourcesBalancedAllocation |
1 | 比较调度后 CPU、内存等资源的利用率差异,倾向于资源使用更均衡的 Node |
ImageLocality |
1 | 根据 Node 已缓存镜像的大小和镜像在集群中的分布评分,减少镜像拉取和启动等待 |
NormalizeScore¶
Score 每次只为一台 Node 打分。某些插件得到的是原始值,需要看到全部候选 Node 后才能换算成可比较的分数。NormalizeScore 接收该插件产生的完整分数列表,并将每台 Node 的分数调整到 Framework 规定的 0–100 范围;分数越高,Node 越符合该插件的偏好。
| 默认插件 | 具体作用 |
|---|---|
TaintToleration |
反向归一化无法容忍的 PreferNoSchedule Taint 数量,使数量更少的 Node 得分更高 |
NodeAffinity |
将首选 Node Affinity 的原始权重映射到 0–100 |
PodTopologySpread |
将拓扑域计数产生的原始分数映射到 0–100 |
InterPodAffinity |
根据候选 Node 的最小值和最大值线性归一化 Affinity 分数 |
DynamicResources |
将 DRA prioritized list 产生的设备分配分数映射到 0–100 |
假设一个插件按照 Node 上已经缓存的数据块数量打分,三台 Node 的原始值分别是 2、6 和 10,NormalizeScore 可以将它们映射为 0、50 和 100。
插件通过 ScoreExtensions() 提供可选的 NormalizeScore。返回 nil 表示 Score() 已经生成可以直接使用的分数,不需要再次归一化。Framework 校验最终分数是否位于 0–100,分数越界时本次调度失败。校验通过后,Framework 将每个插件的分数乘以 Profile 中配置的 weight,再把加权结果相加,选择总分最高的 Node。
Reserve¶
Scheduler 选出目标 Node 并在 Scheduler Cache 中 Assume Pod 后,Reserve 让插件把自己的临时状态同步到这次选择。它接收 Pod 和目标 Node 名称,不负责再次选择 Node。例如 VolumeBinding 会在本地缓存中 Assume 本轮选定的 PV 与 PVC 组合,避免其他调度周期同时使用这组卷绑定结果。
| 默认插件 | 具体作用 |
|---|---|
VolumeBinding |
在 Scheduler 本地缓存中假定选定的 PV/PVC 绑定成立,避免其他 Pod 同时使用相同绑定结果 |
DynamicResources |
在本轮 CycleState 中保存目标 Node 对应的 DRA Claim 和设备分配结果 |
如果某个 Reserve 插件失败,或者后续 Permit、PreBind、Bind 失败,Framework 会按 Reserve 的相反顺序调用 Unreserve。VolumeBinding 会在这里撤销之前 Assume 的 PV 与 PVC 状态;Scheduler 随后还会 Forget Pod,释放 Scheduler Cache 中为目标 Node 记录的假定资源。
ReservePlugin 要求 Unreserve 保持幂等。Framework 可能在对应 Reserve 尚未运行时调用它,因此实现必须能够安全处理不存在的临时状态。
Permit¶
Permit 是进入绑定阶段前的最后一个决策点。它接收 Pod 和已经选定的 Node,在 Scheduler 完成 Assume 与 Reserve 后决定立即放行、保留当前结果继续等待,还是结束本次尝试。
| 默认插件 | 具体作用 |
|---|---|
| 无 | 默认 Profile 不在 Permit 阶段等待或拒绝 Pod |
例如,一组分布式训练 Pod 需要凑齐全部成员才能启动。Permit 插件可以让已经选到 Node 的成员返回 Wait;等全部成员都有位置后,插件再逐个调用 WaitingPod.Allow。等待期间,这些 Pod 仍占用 Scheduler Cache 中的 Assumed 资源,避免同一批资源被后续 Pod 选走。
PermitPlugin.Permit 可以返回以下结果:
- Success:当前插件接受这次选择,Framework 继续运行后面的 Permit 插件。全部插件执行完且没有等待、拒绝或错误时,Pod 进入绑定阶段。
- Wait:当前插件给出等待时长,Framework 继续运行后面的 Permit 插件。全部插件执行完且没有拒绝或错误时,只要存在一个 Wait,Scheduler 就会登记 WaitingPod;全部等待插件调用
Allow后继续绑定,任一插件 Reject 或超时都会撤销本次选择。 - Pending:适用于外部组件需要根据本次节点选择完成后续操作、Pod 暂时还不能绑定的情况。例如,自定义 Permit 插件将目标 Node 交给外部设备分配控制器,等待它更新资源状态。插件返回
Pending后,当前调度尝试结束。Scheduler 随后运行Unreserve和ForgetPod,撤销本地预留;关联事件到达且 QueueingHint 返回Queue后,Pod 直接进入activeQ,无需等待退避。 - Unschedulable:插件认为当前条件下无法继续调度,但抢占或集群状态变化可能解决问题。例如,移除低优先级 Pod 后,目标 Node 可能释放出足够的资源。Scheduler 撤销本地状态,并根据插件注册的事件等待重新调度。
- UnschedulableAndUnresolvable:插件确认抢占无法解决当前问题。例如,Pod 要求的 Node label 或卷拓扑与目标 Node 不匹配,移除其他 Pod 也不会改变结果。Scheduler 撤销本地状态,等待相关配置或对象发生变化后重新调度。
- Error:插件执行失败或收到意外输入。Scheduler 撤销本地状态,并按错误退避路径重试。
两者都会结束当前调度尝试并进入失败退避。区别在于 Unschedulable 表示 PostFilter 仍可能改善结果,而 UnschedulableAndUnresolvable 表示当前条件下无需尝试抢占。Permit 阶段已经选出目标 Node,两种状态使用相同的 Unreserve、ForgetPod 和重新入队流程。
每个 Permit 插件给出的 Wait 时长最多为 15 分钟;上限由 maxTimeout 定义。RunPermitPlugins 收集等待时长,Scheduler 登记 WaitingPod 时再由 newWaitingPod 为各插件启动计时器。任一计时器先到期,整个 Pod 的 Permit 等待都会失败。
绑定阶段(binding cycle)¶
绑定阶段把调度阶段选出的目标 Node 写入 Kubernetes API。在此之前,Scheduler 会等待 Permit 放行,并完成卷绑定等必须成功的准备工作;绑定成功后,再运行通知和清理类操作。
不同 Pod 可以并发执行绑定阶段。同一个 Pod 依次经过 PreBindPreFlight、WaitOnPermit、PreBind、Bind 和 PostBind。scheduleOnePod 为每个完成调度阶段的 Pod 启动一个 goroutine,并在其中调用 runBindingCycle。
WaitOnPermit¶
WaitOnPermit 确保所有请求等待的 Permit 插件都已放行,Pod 才能继续绑定。没有 Permit 插件返回 Wait 时,这一步立即成功;存在等待请求时,Scheduler 会等到所有相关插件调用 Allow、任一插件调用 Reject,或任一等待时限到期。
等待期间,Pod 仍以 Assumed 状态占用 Scheduler Cache 中的目标 Node 资源。等待失败后,Scheduler 运行 Unreserve 并从 Cache 中 ForgetPod,再由 FailureHandler 处理重新入队。
Framework.WaitOnPermit 先查询当前 Pod 的 WaitingPod 记录。没有记录时直接返回 nil,绑定流程继续执行 PreBind;存在记录时才会阻塞当前 Pod 的绑定 goroutine。它是 Framework 的内部步骤,不是插件扩展点。
PreBind¶
PreBind 在写入 Pod 与 Node 的绑定关系前,完成绑定所依赖的准备工作。例如,VolumeBinding 插件会把调度阶段确定的卷绑定结果提交给 API Server,并等待卷控制器完成绑定。
| 默认插件 | 具体作用 |
|---|---|
VolumeBinding |
提交调度阶段假定的 PV/PVC 绑定,并等待卷控制器完成绑定 |
DynamicResources |
写入 ResourceClaim 的设备分配和 reservedFor,必要时等待设备满足绑定条件 |
PreBind 插件接收 Pod、目标 Node 名称和本轮 CycleState。WaitOnPermit 成功后,所有需要运行的 PreBind 插件都必须返回 Success;任一插件失败都会终止本次绑定。
PreBindPlugin 同时包含 PreBindPreFlight 和 PreBind。PreBindPreFlight 只判断插件是否需要处理当前 Pod:Success 表示随后运行 PreBind,Skip 表示跳过,Error 表示终止绑定;AllowParallel 标记该插件的 PreBind 能否与相邻插件并行执行。
RunPreBindPreFlights 使用普通 for 循环逐个执行预检,并将 Skip 和 AllowParallel 结果写入 CycleState,预检本身不并行。RunPreBindPlugins 随后把连续允许并行的插件组成一组,并通过 Parallelizer().Until 并行执行组内插件;不允许并行的插件会截断分组并单独执行。
Bind¶
Bind 把已经选定的目标 Node 写入 Kubernetes API。绑定请求成功后,Pod 的 spec.nodeName 已经保存;Scheduler 的 Informer 随后收到 Pod 更新,并用它确认本地的 Assumed 状态。
| 默认插件 | 具体作用 |
|---|---|
DefaultBinder |
请求 Pod 的 binding 子资源,将目标 Node 写入 spec.nodeName |
BindPlugin.Bind 接收 Pod、目标 Node 名称和 CycleState,并按 Profile 中的配置顺序执行。Skip 表示当前插件不处理该 Pod,Framework 继续调用下一个插件;Success 表示绑定已经完成,其余 Bind 插件不再运行。所有插件都返回 Skip,或者任一插件失败时,本次绑定失败。
默认 DefaultBinder 构造一个 v1.Binding 请求体,并调用 Pod 的 binding 子资源。API Server 根据该请求更新现有 Pod 的 spec.nodeName;v1.Binding 在这里是请求对象,不是另行持久化的资源。具体调用见 DefaultBinder.Bind 和 API Server 的 BindingREST。
PostBind¶
PostBind 在 Bind 成功后执行,用于记录指标、发送通知或清理插件内部状态;这些操作不参与绑定是否成功的判断。
| 默认插件 | 具体作用 |
|---|---|
| 无 | 默认 Profile 不执行 PostBind 回调 |
PostBind 插件接收 Pod、目标 Node 名称和 CycleState,并按配置顺序执行。PostBindPlugin.PostBind 没有返回状态,因此不能撤销已经完成的绑定;插件需要自行处理内部错误。
核心源码¶
下图展示 kube-scheduler 的核心组件与主要调用关系,涵盖初始化、插件管理、调度数据维护,以及 Pod 的调度、绑定和失败处理。
- Scheduler 结构:保存
SchedulingQueue、Cache、Profiles和nodeInfoSnapshot,并通过NextPod、SchedulePod、FailureHandler分别连接取出 Pod、选择 Node 和失败处理。 - Scheduler 初始化与启动:
scheduler.New创建 Framework、队列和 Cache,并注册事件处理函数。完成 Informer 数据同步且满足选主条件后,Scheduler.Run启动队列后台任务和调度主循环。 - Scheduling Profile:将调度器名称、插件及其参数组合成一套调度策略。Pod 通过
spec.schedulerName选择 Profile,Profiles按该名称保存对应的 Framework。 - Framework 插件管理和调用:每个 Profile 对应一个 Framework,负责创建插件实例、维护各扩展点的插件列表,并通过
RunFilterPlugins、RunScorePlugins等方法执行插件。 - Informer 与事件处理:通过 LIST/WATCH 接收 API Server 中的对象变化。事件处理函数将待调度 Pod 及相关事件送入 SchedulingQueue,将 Node 和已绑定 Pod 的变化写入 Cache。
- SchedulingQueue:管理待调度和等待重试的 Pod,使用
activeQ、backoffQ和unschedulablePods保存不同状态的 Pod。NextPod默认调用队列的Pop,取出下一次调度要处理的 Pod。 - Scheduler Cache 与 Snapshot:Cache 汇总 Node、已绑定 Pod 和 Assumed Pod 的状态;每轮调度开始时,通过
UpdateSnapshot刷新节点视图,供调度算法和插件读取。 - 调度阶段:
ScheduleOne取出 Pod 后调用scheduleOnePod,由schedulingCycle同步完成节点选择和绑定前准备。选出 Node 后,先记录 Assumed 状态,再运行 Reserve 和 Permit。 - 绑定阶段:调度阶段成功后,
scheduleOnePod启动独立 goroutine 执行runBindingCycle,完成等待 Permit、PreBind、Bind 和 PostBind。默认 Bind 插件向 API Server 提交 Pod 与 Node 的绑定请求。 - 失败处理:调度或绑定失败时,按需通过 Unreserve 和
ForgetPod撤销插件状态与本地资源占用,再由FailureHandler记录诊断,并将仍需重试的 Pod 交回 SchedulingQueue。
Scheduler 结构¶
Scheduler 结构体集中保存调度器运行所需的本地数据和流程入口。SchedulingQueue、Cache、Profiles 和 nodeInfoSnapshot 提供调度数据;NextPod、SchedulePod 和 FailureHandler 分别连接取出 Pod、选择 Node 和失败处理流程。
- 保存本地数据:
SchedulingQueue保存待处理 Pod,Cache汇总 Node 与 Pod 状态,Profiles保存各调度策略对应的 Framework,nodeInfoSnapshot提供本轮调度使用的节点视图。 - 取出 Pod:
NextPod默认指向SchedulingQueue.Pop,队列为空时阻塞等待。 - 选择 Node:
SchedulePod默认指向schedulePod,使用选定的 Framework 和 NodeInfo Snapshot 完成筛选与评分。 - 处理失败:
FailureHandler默认指向handleSchedulingFailure,负责记录诊断信息并把 Pod 交回 SchedulingQueue。
下面的源码保留普通 Pod 调度主路径使用的字段。
type Scheduler struct {
// 保存 Node、已绑定 Pod 和 Assumed Pod 的本地状态。
Cache internalcache.Cache
...
// 从调度队列中取出下一项;队列为空时阻塞等待。
NextPod func(logger klog.Logger) (*framework.QueuedPodInfo, error)
// 处理调度失败,记录诊断信息并把 Pod 交回队列。
FailureHandler FailureHandlerFn
// 执行节点筛选、评分和目标 Node 选择。
SchedulePod func(
ctx context.Context,
fwk framework.Framework,
state fwk.CycleState,
podInfo *framework.QueuedPodInfo,
) (ScheduleResult, error)
...
// 保存等待调度、退避和暂时不可调度的 Pod。
SchedulingQueue internalqueue.SchedulingQueue
// 按 schedulerName 保存各 Profile 对应的 Framework。
Profiles profile.Map
// 为当前调度周期提供稳定的 NodeInfo 视图。
nodeInfoSnapshot *internalcache.Snapshot
...
}
Scheduler 初始化与启动¶
kube-scheduler 命令读取配置并创建 Scheduler,再启动 Informer、等待初始数据处理完成,最后进入 ScheduleOne 主循环。Leader Election 开启时,只有持有 Lease 的实例会运行主循环。
相关代码位置如下:
cmd/kube-scheduler/app/
├── server.go # 命令入口、Setup 与进程级 Run
└── config/config.go # Config、CompletedConfig:保存进程运行配置
pkg/scheduler/
├── scheduler.go # New 创建组件,Run 启动队列与调度循环
├── eventhandlers.go # 注册事件回调,等待初始对象处理完成
├── backend/queue/scheduling_queue.go # 队列 Run、退避刷新与停留超时检查
└── schedule_one.go # ScheduleOne:主循环调用的调度入口
下图展示了从命令入口到调度主循环的启动顺序。
命令与对象创建
- 创建命令:
NewSchedulerCommand创建 Cobra 命令,在RunE中调用runCommand。runCommand建立进程 Context,再调用Setup。 - 创建 Scheduler:
Setup补齐配置默认值、校验配置,并收集 out-of-tree 插件,再调用scheduler.New创建 Scheduler。 - 启动进程服务:
Setup返回后,runCommand调用app.Run,启动事件广播、安全端口和健康检查,并按配置初始化 Informer、参与选主。
Informer 初始化与选主
-
初始化时机:
- 未启用选主:直接初始化 Informer,完成后开始调度,
DelayCacheUntilActive的取值不影响这条路径。 - 启用选主且
DelayCacheUntilActive = false:先初始化 Informer,再参与选主;成为 Leader 后开始调度。 - 启用选主且
DelayCacheUntilActive = true:成为 Leader 后才初始化 Informer,完成后开始调度。
- 未启用选主:直接初始化 Informer,完成后开始调度,
-
Informer 初始化:
app.Run调用InformerFactory.Start启动 Informer,通过WaitForCacheSync等待初始缓存同步,再通过WaitForHandlersSync等待初始对象经过 Scheduler 注册的事件处理函数。 - 选主回调:启用选主时,实例成为 Leader 后执行
OnStartedLeading回调。回调完成必要的 Informer 初始化后,调用Scheduler.Run。
启动调度循环
-
启动队列后台任务:
Scheduler.Run先调用SchedulingQueue.Run,分别在独立 goroutine 中运行两个任务:- 处理退避到期的 Pod:
flushBackoffQCompleted处理podBackoffQ和podErrorBackoffQ中退避已经结束的 Pod,将可以重新调度的 Pod 移入activeQ,并唤醒等待取 Pod 的调度循环。 - 检查长期未重试的 Pod:
flushUnschedulablePodsLeftover每隔 30 秒检查unschedulablePods,对停留超过阈值的 Pod 触发重新入队检查。检查通过后,按剩余退避时间进入activeQ或backoffQ;未通过则继续留在unschedulablePods。
- 处理退避到期的 Pod:
-
启动调度 goroutine:
SchedulingQueue.Run启动后台任务并返回后,Scheduler.Run再创建独立 goroutine,串行循环调用ScheduleOne,从队列取出 Pod 并执行调度。 - 等待停止信号:原 goroutine 阻塞在
<-ctx.Done()。Context 取消后,它关闭队列并清理插件;等待期间,队列后台任务和调度 goroutine 继续运行。
命令入口¶
NewSchedulerCommand 创建 Cobra 命令和配置选项,并在 RunE 中调用 runCommand。
runCommand 建立进程 Context 并调用 Setup;Setup 内部调用 scheduler.New,再把 CompletedConfig 和 Scheduler 返回给 runCommand。CompletedConfig 保存 Scheduler 进程运行所需的完整上下文,包括组件配置、API Client、Informer Factory、服务端和 Leader Election 配置。
func NewSchedulerCommand(registryOptions ...Option) *cobra.Command {
opts := options.NewOptions()
cmd := &cobra.Command{
Use: "kube-scheduler",
RunE: func(cmd *cobra.Command, args []string) error {
// 命令参数解析完成后进入 Scheduler 的创建与运行流程。
return runCommand(cmd, opts, registryOptions...)
},
...
}
...
return cmd
}
func runCommand(
cmd *cobra.Command,
opts *options.Options,
registryOptions ...Option,
) error {
...
cc, sched, err := Setup(ctx, opts, registryOptions...)
if err != nil {
return err
}
// cc 保存进程运行配置和依赖,sched 是 scheduler.New 创建的调度器。
// Run 使用二者启动 Informer、Leader Election 和调度主循环。
return Run(ctx, cc, sched)
}
Setup 补齐配置默认值、校验配置并收集 out-of-tree 插件,再调用 scheduler.New 创建 Scheduler:
func Setup(
ctx context.Context,
opts *options.Options,
outOfTreeRegistryOptions ...Option,
) (*config.CompletedConfig, *scheduler.Scheduler, error) {
...
// 将 app.WithPlugin 注册的 out-of-tree 插件写入 Registry。
outOfTreeRegistry := make(runtime.Registry)
for _, option := range outOfTreeRegistryOptions {
if err := option(outOfTreeRegistry); err != nil {
return nil, nil, err
}
}
// 传入配置、InformerFactory 和插件 Registry,创建 Scheduler。
sched, err := scheduler.New(
ctx,
cc.Client,
cc.InformerFactory,
cc.DynInformerFactory,
recorderFactory,
scheduler.WithProfiles(cc.ComponentConfig.Profiles...),
scheduler.WithFrameworkOutOfTreeRegistry(outOfTreeRegistry),
...,
)
...
return &cc, sched, nil
}
Scheduler 创建¶
scheduler.New 负责初始化 Scheduler。它合并 in-tree 插件 Registry 与 out-of-tree 插件 Registry,并完成以下初始化工作:
- Snapshot:保存本轮调度读取的节点视图,每轮调度开始时由 Scheduler Cache 增量刷新。
- Scheduler Cache:保存 Node、已绑定 Pod 和 Assumed Pod 的最新聚合状态。
- Scheduling Profile:
profile.NewMap为每个 Profile 创建 Framework,并按schedulerName写入profile.Map。调度 Pod 时,Scheduler 根据spec.schedulerName取得对应的 Framework,运行该 Profile 配置的插件。 - SchedulingQueue:保存等待调度或等待重试的 Pod,所有 Profile 共用一个队列。
- 注册事件处理器:为 Pod、Node 和存储资源注册 Event Handler,用对象变化更新 SchedulingQueue 与 Scheduler Cache。
以下源码展示这些数据结构的创建、字段赋值和 Event Handler 注册位置:
func New(..., opts ...Option) (*Scheduler, error) {
...
// 创建 Kubernetes 源码内置的 in-tree 插件 Registry。
registry := frameworkplugins.NewInTreeRegistry()
// 合并通过 app.WithPlugin 注册的 out-of-tree 插件。
// 后续 Profile 根据插件名称从同一个 Registry 创建插件实例。
if err := registry.Merge(options.frameworkOutOfTreeRegistry); err != nil {
return nil, err
}
// Cache 持续接收 Informer 事件和 Assume、Forget 更新,保存最新的集群状态。
// Snapshot 在每轮调度开始时由 Cache.UpdateSnapshot 增量刷新,供本轮调度读取。
snapshot := internalcache.NewEmptySnapshot()
schedulerCache := internalcache.New(
ctx,
apiDispatcher,
feature.DefaultFeatureGate.Enabled(features.GenericWorkload),
)
// 每个 Profile 创建一个 Framework,并按 schedulerName 保存。
profiles, err := profile.NewMap(
ctx, options.profiles, registry, recorderFactory,
frameworkruntime.WithSnapshotSharedLister(snapshot),
...,
)
if err != nil {
return nil, fmt.Errorf("initializing profiles: %v", err)
}
// 所有 Profile 共用一个 SchedulingQueue。
podQueue := internalqueue.NewSchedulingQueue(
profiles[options.profiles[0].SchedulerName].QueueSortFunc(),
informerFactory,
...,
)
sched := &Scheduler{
// 保存持续变化的 Node 和 Pod 聚合状态。
Cache: schedulerCache,
// 保存等待调度或等待重试的 Pod。
SchedulingQueue: podQueue,
// 按 schedulerName 保存各 Profile 对应的 Framework,供调度 Pod 时查询。
Profiles: profiles,
// 每轮调度开始时由 Cache.UpdateSnapshot 增量刷新。
nodeInfoSnapshot: snapshot,
...,
}
// 主循环通过 NextPod 从 SchedulingQueue 取出下一个 Pod。
sched.NextPod = podQueue.Pop
// 设置默认的 SchedulePod 和 FailureHandler。
sched.applyDefaultHandlers()
// 从 SharedInformerFactory 取得 Pod、Node、PVC 等资源各自的 Shared Informer,
// 并为它们注册 Add、Update 和 Delete 回调。
// Informer 启动后,回调更新 SchedulingQueue 与 Scheduler Cache,并触发 Pod 重新入队。
if err := addAllEventHandlers(sched, informerFactory, ...); err != nil {
return nil, fmt.Errorf("adding event handlers: %w", err)
}
return sched, nil
}
数据同步与主循环¶
app.Run 是进程级运行入口,负责健康检查、事件广播、Shared Informer 同步和 Leader Election。WaitForCacheSync 等待各 Informer 完成第一次 LIST 并填充本地 Store,WaitForHandlersSync 等待这些对象经过 Scheduler 注册的 Event Handler。数据同步完成后调用 Scheduler.Run,启动 SchedulingQueue 后台任务和 ScheduleOne 主循环。
func Run(
ctx context.Context,
cc *schedulerserverconfig.CompletedConfig,
sched *scheduler.Scheduler,
) error {
logger := klog.FromContext(ctx)
...
// 这里只定义启动 Informer 并等待同步的函数,尚未执行。
// 后面根据 DelayCacheUntilActive 和 Leader Election 的状态调用它。
startInformersAndWaitForSync := func(ctx context.Context) {
// 启动 SharedInformerFactory 中已经创建的 Informer。
cc.InformerFactory.Start(ctx.Done())
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.Start(ctx.Done())
}
// 等待初始 LIST 完成,确保各 Informer 的 Store 已经填充。
cc.InformerFactory.WaitForCacheSync(ctx.Done())
if cc.DynInformerFactory != nil {
cc.DynInformerFactory.WaitForCacheSync(ctx.Done())
}
// 等待初始对象经过 Scheduler 注册的 Event Handler。
if err := sched.WaitForHandlersSync(ctx); err != nil {
logger.Error(err, "handlers are not fully synchronized")
}
...
}
// 未延迟缓存,或未启用 Leader Election 时,在这里调用并完成同步。
if !cc.ComponentConfig.DelayCacheUntilActive || cc.LeaderElection == nil {
startInformersAndWaitForSync(ctx)
}
if cc.LeaderElection != nil {
...
cc.LeaderElection.Callbacks = leaderelection.LeaderCallbacks{
OnStartedLeading: func(ctx context.Context) {
...
// 延迟缓存时,当前实例成为 Leader 后才在这里调用该函数。
if cc.ComponentConfig.DelayCacheUntilActive {
startInformersAndWaitForSync(ctx)
}
// 只有 Leader 进入 Scheduler 主循环。
sched.Run(ctx)
},
OnStoppedLeading: func() {
...
},
}
leaderElector, err := leaderelection.NewLeaderElector(*cc.LeaderElection)
if err != nil {
return fmt.Errorf("couldn't create leader elector: %v", err)
}
leaderElector.Run(ctx)
return fmt.Errorf("lost lease")
}
// 未启用 Leader Election,当前实例直接运行 Scheduler。
sched.Run(ctx)
...
}
delayCacheUntilActive 默认为 false,所有实例都会提前启动 Informer 并同步本地调度数据,以缩短 Leader 切换时间。设置为 true 后,备用实例成为 Leader 时才同步数据,可以减少内存和 API Server 开销,但会增加切换后的等待时间。当未启用 Leader Election 时,Informer 会立即启动。
Scheduler.Run 启动队列后台任务,并在独立 goroutine 中循环调用 ScheduleOne:
func (sched *Scheduler) Run(ctx context.Context) {
// 启动退避刷新、不可调度 Pod 超时检查等队列任务。
sched.SchedulingQueue.Run(klog.FromContext(ctx))
...
// ScheduleOne 从 SchedulingQueue 取出下一个 Pod;队列为空时会阻塞等待。
// 如果在当前 goroutine 中运行,Context 取消后将无法继续执行
// SchedulingQueue.Close 来解除阻塞,导致关闭过程死锁,因此使用独立 goroutine。
go wait.UntilWithContext(ctx, sched.ScheduleOne, 0)
<-ctx.Done()
sched.SchedulingQueue.Close()
_ = sched.Profiles.Close()
}
Scheduling Profile¶
Scheduling Profile 将调度器名称、插件及其参数组合成一套调度策略。一个 kube-scheduler 进程可以配置多个 Profile,Pod 通过 spec.schedulerName 选择其中一套策略。kube-scheduler 随后按照该 Profile 的配置,在各个扩展点运行对应插件。
相关代码位置如下:
staging/src/k8s.io/kube-scheduler/config/v1/
└── types.go # KubeSchedulerProfile、插件列表与参数配置字段
pkg/scheduler/
├── apis/config/v1/
│ ├── defaults.go # 补齐 Profile 默认值与内置插件参数
│ └── default_plugins.go # 默认插件集合,以及与用户插件配置的合并
├── profile/profile.go # NewMap 创建并校验 Profile,保存 Framework
├── framework/runtime/framework.go # NewFramework 按单份 Profile 创建插件运行环境
└── schedule_one.go # frameworkForPod 选择 Framework;计算可行节点数量
下图展示 Profile 的创建与选择:启动阶段将每份配置转换为 Framework,并以 schedulerName 为 key 写入 profile.Map;调度 Pod 时,Scheduler 根据 spec.schedulerName 取得对应的 Framework。
- 读取 Profile 配置:
KubeSchedulerProfile保存调度器名称、启用的插件及插件参数。 - 合并默认插件并补齐参数:Scheduler 将用户配置与内置插件集合合并,并补齐缺少的内置插件参数。
- 遍历 Profile:
profile.NewMap依次处理每份 Profile 配置。 - 创建 Framework:
newProfile为当前 Profile 调用runtime.NewFramework。 - 校验 Profile:
schedulerName必须非空且唯一;多个 Profile 的 QueueSort 插件名称和参数必须一致。校验失败时,Scheduler 无法启动。 - 建立索引:校验通过后,Framework 以
schedulerName为 key 写入profile.Map。 - 读取调度器名称:调度 Pod 时,Scheduler 读取
spec.schedulerName。 - 查找 Framework:
frameworkForPod使用该名称查询profile.Map;没有对应 Profile 时,本次调度返回错误。 - 运行插件:查询成功后,Scheduler 使用对应的 Framework 运行该 Profile 配置的插件。
Profile 配置¶
KubeSchedulerProfile 包含以下字段:
schedulerName:Profile 的名称,与 Pod 的spec.schedulerName匹配。只有一个 Profile 且未填写名称时,默认值为default-scheduler;配置多个 Profile 时,每个名称都必须非空且唯一。percentageOfNodesToScore:控制找到多少个可行 Node 后停止筛选。默认值0表示自动计算,公式为max(5, 50 - Node 总数 / 125)%,同时至少查找 100 个可行 Node。例如,集群有 1,000 个 Node 时比例为 42%,找到 420 个可行 Node 后停止。numFeasibleNodesToFind实现了这项计算。plugins:按扩展点启用、禁用或重新配置插件。enabled可以添加插件,也可以覆盖默认插件在该扩展点上的顺序和 Score 权重;disabled删除指定的默认插件,*表示删除该扩展点的全部默认插件。pluginConfig:按插件名称提供初始化参数。一个插件实现多个扩展点时只创建一个实例,这组参数由该实例共同使用。
setDefaults_KubeSchedulerProfile 先把用户配置与默认插件集合合并。对于已经启用但没有配置参数的内置插件,Scheduler 会生成默认参数,供创建插件实例时使用。
func setDefaults_KubeSchedulerProfile(
logger klog.Logger,
prof *configv1.KubeSchedulerProfile,
) {
// enabled、disabled 和 MultiPoint 配置在这里与默认插件合并。
prof.Plugins = mergePlugins(logger, getDefaultPlugins(), prof.Plugins)
scheme := GetPluginArgConversionScheme()
existingConfigs := sets.New[string]()
// 遍历用户显式提供的 pluginConfig,并为 Args 中未设置的字段补上默认值。
for j := range prof.PluginConfig {
// 记录已经配置过参数的插件,后续不再为它生成整份默认 pluginConfig。
existingConfigs.Insert(prof.PluginConfig[j].Name)
args := prof.PluginConfig[j].Args.Object
// out-of-tree 插件参数由插件工厂自行解析并补齐默认值。
if _, isUnknown := args.(*runtime.Unknown); isUnknown {
continue
}
scheme.Default(args)
}
// 为其余内置插件补上对应的默认 Args。
for _, name := range pluginsNames(prof.Plugins) {
if existingConfigs.Has(name) {
continue
}
...
prof.PluginConfig = append(prof.PluginConfig, configv1.PluginConfig{
Name: name,
Args: runtime.RawExtension{Object: args},
})
}
}
插件配置¶
这份配置包含两个 Profile。default-scheduler 使用默认插件;gpu-binpack-scheduler 调整插件权重、默认抢占和插件初始化参数。
gpu-binpack-scheduler 是本例定义的 Profile 名称。它继续使用 kube-scheduler 的内置插件,只通过 KubeSchedulerProfile 修改插件配置,不需要编写或编译新的调度插件。当内置插件无法满足调度规则时,可以使用 Scheduling Framework 编写自定义插件;后面的自定义调度插件通过完整示例介绍插件接口、注册、配置和验证过程。
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
profiles:
# 未配置插件时,Profile 使用 kube-scheduler 的默认插件和参数。
- schedulerName: default-scheduler
# 这个 Profile 倾向于集中放置 CPU、内存和 GPU 工作负载。
- schedulerName: gpu-binpack-scheduler
plugins:
score:
# 不再根据镜像是否已经缓存在 Node 上增加分数。
disabled:
- name: ImageLocality
# 覆盖默认 Score 插件的权重;其他扩展点保持默认配置。
enabled:
- name: NodeResourcesFit
weight: 5
- name: PodTopologySpread
weight: 3
postFilter:
# 没有可行 Node 时不运行默认抢占。
disabled:
- name: DefaultPreemption
pluginConfig:
- name: NodeResourcesFit
args:
# 优先选择资源利用率较高的 Node,集中放置工作负载。
scoringStrategy:
type: MostAllocated
resources:
- name: cpu
weight: 1
- name: memory
weight: 1
- name: nvidia.com/gpu
weight: 5
- 多个 Score 插件:
NodeResourcesFit和PodTopologySpread仍在默认扩展点运行,这里只把它们的 Score 权重改为5和3。 - 禁用单个插件:
ImageLocality只从 Score 扩展点移除;其他默认 Score 插件继续参与评分。 - 禁用默认抢占:
DefaultPreemption从 PostFilter 移除。该 Profile 找不到可行 Node 时,不会通过默认抢占移除低优先级 Pod。 - 配置插件参数:
NodeResourcesFit的scoringStrategy.type设置为MostAllocated,资源已请求比例越高的 Node 得分越高。resources[].name指定 CPU、内存和 GPU 参与评分,resources[].weight将nvidia.com/gpu的比重设为 CPU 和内存的 5 倍。
Pod 设置 spec.schedulerName: gpu-binpack-scheduler 时,Scheduler 使用第二个 Profile;没有显式设置该字段的 Pod 由 API Server 默认填写 default-scheduler。
Profile 创建¶
profile.NewMap 为每份 Profile 创建 Framework,并以 schedulerName 建立索引;Framework 保存该 Profile 的插件实例和扩展点执行顺序,具体结构见 Framework 插件管理和调用。
profile.NewMap 遍历所有 Profile 配置。newProfile 将一份配置交给 frameworkruntime.NewFramework,创建对应的 Framework;校验通过后,再以 schedulerName 为 key 保存到 profile.Map:
func newProfile(
ctx context.Context,
cfg config.KubeSchedulerProfile,
r frameworkruntime.Registry,
recorderFact RecorderFactory,
opts ...frameworkruntime.Option,
) (framework.Framework, error) {
recorder := recorderFact(cfg.SchedulerName)
opts = append(opts, frameworkruntime.WithEventRecorder(recorder))
// 根据 Profile 的 plugins 和 pluginConfig 创建插件实例与执行链。
return frameworkruntime.NewFramework(ctx, r, &cfg, opts...)
}
// Map 使用 schedulerName 索引各个 Profile 对应的 Framework。
type Map map[string]framework.Framework
func NewMap(...) (Map, error) {
m := make(Map)
v := cfgValidator{m: m}
for _, cfg := range cfgs {
// 每份 Profile 配置创建一个 Framework。
p, err := newProfile(ctx, cfg, r, recorderFact, opts...)
if err != nil {
...
}
// 检查 schedulerName 和多个 Profile 的 QueueSort 配置。
if err := v.validate(cfg, p); err != nil {
return nil, err
}
m[cfg.SchedulerName] = p
}
return m, nil
}
cfgValidator.validate 检查 Profile 名称是否为空或重复,并比较不同 Profile 的 QueueSort 插件及参数。
Profile 选择¶
调度一个 Pod 时,frameworkForPod 使用 spec.schedulerName 查询这个 Map:
func (sched *Scheduler) frameworkForPod(pod *v1.Pod) (framework.Framework, error) {
// schedulerName 决定本次调度使用哪一组插件和插件参数。
fwk, ok := sched.Profiles[pod.Spec.SchedulerName]
if !ok {
return nil, fmt.Errorf(
"profile not found for scheduler name %q",
pod.Spec.SchedulerName,
)
}
return fwk, nil
}
选择同一个 Profile 的 Pod 复用同一个 Framework;每个 Pod 会创建独立的 CycleState,保存当前调度周期中的临时计算结果。
Framework 插件管理和调用¶
Framework 是 Scheduler 根据 Profile 创建的插件管理与调用对象。Scheduler 通过 Framework 运行当前 Profile 配置的扩展点插件。
相关代码位置如下,plugins/ 中列出部分内置插件实现:
pkg/scheduler/framework/
├── interface.go # Scheduler 内部使用的 Framework 接口
├── runtime/
│ ├── registry.go # Registry 与 PluginFactory:保存插件工厂
│ └── framework.go # frameworkImpl、NewFramework 与扩展点插件调用
└── plugins/
├── registry.go # NewInTreeRegistry:注册内置插件工厂
├── schedulinggates/scheduling_gates.go # SchedulingGates:PreEnqueue 检查调度门控
├── queuesort/priority_sort.go # PrioritySort:按优先级和入队时间排序
├── nodeunschedulable/node_unschedulable.go # NodeUnschedulable:Filter 检查节点禁调度标记
├── nodeaffinity/node_affinity.go # NodeAffinity:节点亲和性筛选与评分
├── noderesources/fit.go # NodeResourcesFit:资源请求检查与评分
├── defaultpreemption/default_preemption.go # DefaultPreemption:PostFilter 抢占,复用 Filter
├── volumebinding/volume_binding.go # VolumeBinding:卷匹配与预留,在 PreBind 中绑定卷
├── defaultbinder/default_binder.go # DefaultBinder:Bind 提交 Pod 绑定请求
└── ... # 其他内置插件未在此展开
staging/src/k8s.io/kube-scheduler/framework/
└── interface.go # 公开的 Handle、PluginsRunner 与扩展点插件接口
插件来源与配置
-
读取插件工厂:Registry 按插件名称保存
PluginFactory。NewFramework根据 Profile 中启用的插件名称,从 Registry 取得对应工厂并创建插件实例。 -
读取插件列表:Profile 使用
plugins指定各扩展点启用或禁用哪些插件,以及同一扩展点内的执行顺序。MultiPoint可以把一个插件展开到它实现的多个扩展点。 -
读取插件参数:Profile 使用
pluginConfig按插件名称提供初始化参数。同一个插件只读取一份参数,即使它同时实现 Filter、Score 等多个扩展点。
Framework 构建
-
创建 Framework:
runtime.NewFramework根据 Registry 和当前 Profile 创建 Framework,只初始化当前 Profile 实际启用的插件。 -
保存插件实例:
frameworkImpl是Framework的实现。每个插件名称在pluginsMap中只对应一个实例,多个扩展点可以引用同一个对象。 -
建立扩展点列表:
frameworkImpl按 Profile 配置建立 PreFilter、Filter、Score、Bind 等列表。列表保存插件实例和调用顺序,Run*Plugins运行时直接读取对应列表。 -
保存 Score 权重:
scorePluginWeight按插件名称保存经过校验的 Score 权重。未设置或配置为0时使用1,该权重只在RunScorePlugins汇总 Node 分数时使用。
接口分工
-
提供调度调用接口:Scheduler 通过
Framework运行 PreEnqueue、PreFilter、PostFilter、Reserve、Permit、PreBind、Bind 和 PostBind 等插件。frameworkImpl实现该接口,因此 Scheduler 不需要依赖具体插件类型。 -
提供插件共享能力:
fwk.Handle提供 Snapshot、Informer 和 ClientSet 等能力。Framework嵌入fwk.Handle;NewFramework调用PluginFactory时,将当前frameworkImpl作为Handle参数传入,插件可以保存并使用它。 -
提供可复用运行方法:内部
Framework嵌入fwk.Handle,后者又嵌入PluginsRunner,因此两者都包含RunFilterPlugins、RunPreScorePlugins和RunScorePlugins等方法。Scheduler 通过Framework调用,插件通过工厂传入的Handle调用;实际执行的都是当前 Profile 对应的frameworkImpl中的方法。
插件调用
-
调用扩展点插件:
frameworkImpl实现各扩展点的Run*Plugins方法。Scheduler 通过Framework接口调用这些方法;方法读取对应的扩展点列表,再按配置顺序调用插件实现。Filter 可以拒绝当前 Node,Permit 可以让 Pod 等待,Score 结果需要校验并乘以插件权重。 -
Filter 调用示例:图中的
RunFilterPlugins由frameworkImpl实现。它读取filterPlugins,跳过本轮无需执行的插件,再按顺序调用各插件的Filter()方法;遇到首个非Success结果即返回。这里展示的是一台 Node 上的插件调用顺序,多台 Node 的检查由调度主流程并行发起。 -
实现插件逻辑:
PluginFactory创建具体插件实例,NewFramework将实例保存到pluginsMap。NodeResourcesFit、NodeAffinity 和 PodTopologySpread 等内置插件实现PreScorePlugin、ScorePlugin或其他扩展点接口,由frameworkImpl调用其PreScore()、Score()、Filter()等方法。同一个实例可以实现并加入多个扩展点列表。
Framework 接口¶
Framework 是 Scheduler 运行当前 Profile 插件的统一接口,由 frameworkImpl 实现。它嵌入 fwk.Handle,获得 NodeInfo Snapshot、Informer、Kubernetes Client 等共享能力,以及 PluginsRunner 定义的部分插件运行方法;Framework 本身再声明调度和绑定阶段需要的其他入口:
type Framework interface {
// Handle 提供 Filter、PreScore、Score 和 SharedLister 等能力。
fwk.Handle
// 控制 Pod 入队条件、重试事件和 activeQ 的排序方式。
PreEnqueuePlugins() []fwk.PreEnqueuePlugin
EnqueueExtensions() []fwk.EnqueueExtensions
QueueSortFunc() fwk.LessFunc
// 运行筛选前的准备和无可行 Node 时的补救插件。
RunPreFilterPlugins(
ctx context.Context, state fwk.CycleState, pod *v1.Pod,
) (*fwk.PreFilterResult, *fwk.Status, sets.Set[string])
RunPostFilterPlugins(
ctx context.Context, state fwk.CycleState,
pod *v1.Pod, filteredNodeStatusMap fwk.NodeToStatusReader,
) (*fwk.PostFilterResult, *fwk.Status)
// 在写入 Binding 前后运行对应插件。
RunPreBindPlugins(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
) *fwk.Status
RunPreBindPreFlights(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
) *fwk.Status
RunPostBindPlugins(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
)
// Reserve 先占用插件内部资源;后续失败时按逆序运行 Unreserve。
RunReservePluginsReserve(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
) *fwk.Status
RunReservePluginsUnreserve(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
)
// Permit 可以放行、拒绝或让 Pod 等待外部条件。
RunPermitPlugins(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
) (pluginsWaitTime map[string]time.Duration, status *fwk.Status)
AddWaitingPod(pod *v1.Pod, pluginsWaitTime map[string]time.Duration)
WillWaitOnPermit(ctx context.Context, pod *v1.Pod) bool
WaitOnPermit(ctx context.Context, pod *v1.Pod) *fwk.Status
// Bind 插件提交 Pod 与 Node 的绑定关系。
RunBindPlugins(
ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string,
) *fwk.Status
// ...
Close() error
}
RunFilterPlugins、RunPreScorePlugins 和 RunScorePlugins 没有直接声明在上面的 Framework 接口中,而是定义在 PluginsRunner 中。Framework 嵌入 fwk.Handle,fwk.Handle 再嵌入 PluginsRunner,因此 Scheduler 仍然可以通过 Framework 调用这些方法:
type Handle interface {
...
// 允许插件复用 Framework 的 Filter、PreScore 和 Score 执行逻辑。
PluginsRunner
...
}
type PluginsRunner interface {
RunPreScorePlugins(
context.Context, CycleState, *v1.Pod, []NodeInfo,
) *Status
RunScorePlugins(
context.Context, CycleState, *v1.Pod, []NodeInfo,
) ([]NodePluginScores, *Status)
RunFilterPlugins(
context.Context, CycleState, *v1.Pod, NodeInfo,
) *Status
...
}
RunFilterPlugins、RunPreScorePlugins 和 RunScorePlugins 最初直接声明在 Framework 接口中,后来为了让获得 Handle 的插件复用这些运行方法,被提取到独立的 PluginsRunner 接口。例如,DefaultPreemption 实现了 PostFilterPlugin 接口,是 PostFilter 扩展点的默认抢占插件。它在结构体中保存 fwk.Handle;抢占模拟从 NodeInfo 副本中移除低优先级 Pod 后,SelectVictimsOnNode 通过 Handle 调用 RunFilterPluginsWithNominatedPods。该方法最终调用 RunFilterPlugins,复用当前 Profile 配置的 Filter 插件:
type DefaultPreemption struct {
fh fwk.Handle
...
}
// 编译期确认 DefaultPreemption 实现 PostFilterPlugin 接口。
var _ fwk.PostFilterPlugin = &DefaultPreemption{}
func (pl *DefaultPreemption) SelectVictimsOnNode(
ctx context.Context,
state fwk.CycleState,
pod *v1.Pod,
nodeInfo fwk.NodeInfo,
pdbs []*policy.PodDisruptionBudget,
) ([]*v1.Pod, int, *fwk.Status) {
...
// 使用当前 Profile 的 Filter 插件检查移除候选 Pod 后的 Node。
if status := pl.fh.RunFilterPluginsWithNominatedPods(
ctx, state, pod, nodeInfo,
); !status.IsSuccess() {
return nil, 0, status
}
...
}
Framework 创建¶
Framework 创建时同时读取 Registry 和 Profile。Registry 保存插件名称与 PluginFactory 的对应关系;Profile 决定启用哪些插件,并通过 pluginConfig 提供初始化参数。PluginFactory 和 Registry 的定义如下:
// PluginFactory 使用插件配置和 Framework Handle 创建一个插件实例。
type PluginFactory = func(
ctx context.Context,
// configuration 对应当前插件的 pluginConfig.args。
configuration runtime.Object,
// f 向插件提供 Snapshot、Informer 和 ClientSet 等共享能力。
f fwk.Handle,
) (fwk.Plugin, error)
// Registry 按插件名称保存对应的插件工厂。
type Registry map[string]PluginFactory
frameworkImpl 保存当前 Profile 启用的插件实例,并按扩展点组织执行列表。pluginsMap 按插件名称索引实例;一个插件即使实现多个扩展点,也只创建一个实例,各扩展点列表引用 pluginsMap 中的同一个对象。
每个列表的元素类型对应一个公开的扩展点接口。例如,filterPlugins []fwk.FilterPlugin 保存实现了 FilterPlugin 的插件实例,顺序来自 Profile 配置。RunFilterPlugins 遍历该列表时,循环变量 pl 表示其中一个插件实例,随后通过 pl.Filter(...) 调用它实现的 Filter 方法:
type frameworkImpl struct {
registry Registry
// 按插件名称保存经过校验的 Score 权重,仅 RunScorePlugins 使用。
scorePluginWeight map[string]int
// 各扩展点按 Profile 配置保存插件调用顺序。
preEnqueuePlugins []fwk.PreEnqueuePlugin
enqueueExtensions []fwk.EnqueueExtensions
queueSortPlugins []fwk.QueueSortPlugin
preFilterPlugins []fwk.PreFilterPlugin
filterPlugins []fwk.FilterPlugin
postFilterPlugins []fwk.PostFilterPlugin
preScorePlugins []fwk.PreScorePlugin
scorePlugins []fwk.ScorePlugin
reservePlugins []fwk.ReservePlugin
preBindPlugins []fwk.PreBindPlugin
bindPlugins []fwk.BindPlugin
postBindPlugins []fwk.PostBindPlugin
permitPlugins []fwk.PermitPlugin
// 每个插件名称只对应一个实例,各扩展点列表引用这里的对象。
pluginsMap map[string]fwk.Plugin
...
}
NewFramework 按以下顺序创建插件并建立扩展点列表:
-
选择插件:读取当前 Profile 的
plugins和pluginConfig,再从 Registry 找到对应的插件工厂。 -
创建插件实例:插件工厂收到配置参数和
Handle。Handle是 Framework 提供给插件的共享能力接口,可以访问 NodeInfo Snapshot、Informer 和 Kubernetes Client 等能力。创建完成的实例按名称写入pluginsMap。 -
建立扩展点列表:
updatePluginList按 Profile 中的配置顺序,将pluginsMap中的实例加入 PreFilter、Filter、Score、Bind 等执行列表。同一个插件实现多个扩展点时,这些列表引用同一个实例。 -
展开 MultiPoint:
expandMultiPointPlugins将 MultiPoint 中启用的插件加入它所实现的全部扩展点。某个扩展点存在显式的enabled或disabled配置时,以显式配置为准。 -
保存 Score 权重:
getValidScoreWeights校验已启用 Score 插件的权重,并按插件名称写入scorePluginWeight。权重没有设置或配置为0时使用1。该权重只作用于 Score 扩展点。
func NewFramework(
ctx context.Context,
r Registry,
profile *config.KubeSchedulerProfile,
opts ...Option,
) (framework.Framework, error) {
...
f := &frameworkImpl{
registry: r,
...,
}
f.profileName = profile.SchedulerName
// pg 只保留当前 Profile 启用的插件名称。
pg := f.pluginsNeeded(profile.Plugins)
pluginConfig := make(map[string]runtime.Object, len(profile.PluginConfig))
for i := range profile.PluginConfig {
name := profile.PluginConfig[i].Name
if _, ok := pluginConfig[name]; ok {
return nil, fmt.Errorf("repeated config for plugin %s", name)
}
pluginConfig[name] = profile.PluginConfig[i].Args
}
// pluginsMap 按名称保存当前 Profile 使用的唯一插件实例。
f.pluginsMap = make(map[string]fwk.Plugin)
for name, factory := range r {
// Registry 包含全部已注册工厂,只创建当前 Profile 启用的插件。
if !pg.Has(name) {
continue
}
// PluginFactory 的第三个参数类型是 fwk.Handle。
// f 是当前 frameworkImpl,向插件提供 Snapshot、Informer 和 ClientSet 等能力。
args := pluginConfig[name]
p, err := factory(ctx, args, f)
if err != nil {
return nil, fmt.Errorf("initializing plugin %q: %w", name, err)
}
// 扩展点列表随后通过插件名称复用这个实例。
f.pluginsMap[name] = p
// 收集插件关注的重新入队事件,未实现时使用默认事件集合。
f.fillEnqueueExtensions(p)
}
// 将各扩展点显式启用的插件按配置顺序写入对应的执行列表。
// updatePluginList 从 pluginsMap 取得实例,并检查它是否实现该扩展点接口。
for _, point := range f.getExtensionPoints(profile.Plugins) {
if err := updatePluginList(point.slicePtr, *point.plugins, f.pluginsMap); err != nil {
return nil, err
}
}
// MultiPoint 只需启用一次插件,再按插件实现的接口加入所有适用的扩展点。
// 某个扩展点显式配置的 enabled 或 disabled 优先于 MultiPoint。
if len(profile.Plugins.MultiPoint.Enabled) > 0 {
if err := f.expandMultiPointPlugins(logger, profile); err != nil {
return nil, err
}
}
// 校验 Score 插件的权重;未设置或为 0 时按 1 处理。
podScoreWeights, err := getValidScoreWeights(
f,
reflect.TypeFor[fwk.ScorePlugin](),
append(profile.Plugins.Score.Enabled, profile.Plugins.MultiPoint.Enabled...),
)
if err != nil {
return nil, fmt.Errorf("score plugins: %w", err)
}
// RunScorePlugins 按插件名称读取权重并计算加权分数。
f.scorePluginWeight = podScoreWeights
...
return f, nil
}
插件调用¶
Framework 按 Profile 配置的顺序调用同一扩展点的插件。以 RunFilterPlugins 为例:它逐个检查 Filter 插件,遇到第一个非 Success 结果便停止,并在返回状态中记录拒绝 Pod 的插件名。
func (f *frameworkImpl) RunFilterPlugins(
ctx context.Context,
state fwk.CycleState,
pod *v1.Pod,
nodeInfo fwk.NodeInfo,
) *fwk.Status {
for _, pl := range f.filterPlugins {
// PreFilter 返回 Skip 时,同名 Filter 插件在本轮一起跳过。
if state.GetSkipFilterPlugins().Has(pl.Name()) {
continue
}
if status := f.runFilterPlugin(ctx, pl, state, pod, nodeInfo); !status.IsSuccess() {
// Filter 只应返回 Success 或拒绝状态;其他状态按插件错误处理。
if !status.IsRejected() {
status = fwk.AsStatus(fmt.Errorf(
"running %q filter plugin: %w", pl.Name(), status.AsError(),
))
}
// 失败状态携带插件名,供诊断和 QueueingHint 使用。
status.SetPlugin(pl.Name())
return status
}
}
return nil
}
RunFilterPlugins 通过 runFilterPlugin 调用当前插件的 Filter 方法。CycleState 未要求记录插件指标时直接调用;需要记录时,runFilterPlugin 还会统计插件的执行时间和返回状态:
func (f *frameworkImpl) runFilterPlugin(
ctx context.Context,
pl fwk.FilterPlugin,
state fwk.CycleState,
pod *v1.Pod,
nodeInfo fwk.NodeInfo,
) *fwk.Status {
// 不记录当前插件的指标时,直接执行 Filter。
if !state.ShouldRecordPluginMetrics() {
return pl.Filter(ctx, state, pod, nodeInfo)
}
startTime := time.Now()
status := pl.Filter(ctx, state, pod, nodeInfo)
// 按扩展点、插件名和返回状态记录本次调用耗时。
f.metricsRecorder.ObservePluginDurationAsync(
metrics.Filter,
pl.Name(),
status.Code().String(),
metrics.SinceInSeconds(startTime),
)
return status
}
pl 的类型是 fwk.FilterPlugin,实际调用的方法由当前插件实现。例如,内置 NodeUnschedulable.Filter 检查当前 Pod 能否放到被标记为不可调度的 Node:
func (pl *NodeUnschedulable) Filter(
ctx context.Context,
_ fwk.CycleState,
pod *v1.Pod,
nodeInfo fwk.NodeInfo,
) *fwk.Status {
node := nodeInfo.Node()
// Node 可以调度,返回 nil 表示检查通过。
if !node.Spec.Unschedulable {
return nil
}
logger := klog.FromContext(ctx)
// Node 已被标记为不可调度,检查 Pod 是否容忍对应的 NoSchedule Taint。
podToleratesUnschedulable := v1helper.TolerationsTolerateTaint(
logger,
pod.Spec.Tolerations,
&v1.Taint{
Key: v1.TaintNodeUnschedulable,
Effect: v1.TaintEffectNoSchedule,
},
utilfeature.DefaultFeatureGate.Enabled(
features.TaintTolerationComparisonOperators,
),
)
if !podToleratesUnschedulable {
// Pod 不容忍该 Taint,拒绝状态无法通过抢占等方式解决。
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
ErrReasonUnschedulable,
)
}
// Pod 容忍 node.kubernetes.io/unschedulable:NoSchedule,检查通过。
return nil
}
SchedulingQueue¶
SchedulingQueue 是 Scheduler 使用的队列接口,定义等待调度 Pod 的加入、更新、取出和重新入队操作。
相关代码位置如下:
pkg/scheduler/
├── backend/queue/
│ ├── scheduling_queue.go # SchedulingQueue、PriorityQueue 与重新入队策略
│ ├── active_queue.go # activeQueue、取出 Pod 与调度期间的事件跟踪
│ ├── backoff_queue.go # 普通退避与错误退避队列
│ └── unschedulable_pods.go # 保存等待相关事件的不可调度 Pod
├── scheduler.go # 创建 SchedulingQueue 并启动后台任务
├── eventhandlers.go # 将待调度 Pod 变化与相关集群事件交给队列
└── schedule_one.go # 取出 Pod,处理失败后的重新入队
下图展示 Pod 在 activeQ、backoffQ 和 unschedulablePods 之间的主要流转路径。
- 接收待调度 Pod:新增或更新、尚未指定
spec.nodeName且由当前 Scheduler 负责的 Pod 被加入 SchedulingQueue。 - 检查首次入队:Pod 先运行 PreEnqueue。全部插件通过后进入
activeQ;任一插件拒绝后进入unschedulablePods,并记录拒绝它的插件。 - 取出 Pod:Scheduler 按 QueueSort 的顺序从
activeQ取出 Pod。启用SchedulerPopFromBackoffQ后,activeQ为空时也可以直接从podBackoffQ取出一个 Pod。 -
执行调度与绑定:
- 调度阶段成功:Scheduler 启动异步的绑定阶段。
- 调度阶段失败:FailureHandler 调用
AddUnschedulableIfNotPresent处理重新入队。- 插件拒绝且调度期间没有相关事件时,Pod 进入
unschedulablePods。 - 插件拒绝且调度期间已经发生相关事件时,Scheduler 调用相应插件的 QueueingHint,判断该事件是否可能改变上次的调度结果,以及是否需要重新入队。
- Scheduler Error 进入
podErrorBackoffQ的重试路径。
- 插件拒绝且调度期间没有相关事件时,Pod 进入
- 绑定阶段成功:API Server 保存绑定结果,Pod 进入 Bound 状态。
- 绑定阶段失败:Scheduler 回滚本地状态,再由 FailureHandler 处理重新入队。
-
重新入队:SchedulingQueue 根据失败原因和相关集群事件安排 Pod 重试,先选择重试方式,再检查 Pod 当前是否满足入队条件。
- 响应集群事件:Node、Pod、PVC 等对象变化时,SchedulingQueue 调用此前拒绝该 Pod 的插件为该事件注册的 QueueingHint。相关回调均返回
QueueSkip时,Pod 继续留在unschedulablePods;有回调返回Queue时,继续选择重试方式。 -
选择重试方式:
- 跳过退避:
PendingPlugins中有插件的 QueueingHint 返回Queue时,优先跳过退避。 - 普通退避:仅其他拒绝插件的 QueueingHint 返回
Queue时,根据剩余退避时间安排重试。 - 错误退避:发生 Scheduler Error 时,失败处理直接选择错误退避方式,无需等待集群事件。
- 跳过退避:
-
处理停留超时:Pod 在
unschedulablePods中停留超时时,定时任务 会触发重新入队检查,并按普通退避方式处理。 - 检查并写入队列:启用
SchedulerPopFromBackoffQ后,activeQ为空时,Scheduler 可以直接从podBackoffQ取出 Pod,不必等待退避结束。因此,Pod 在进入退避队列前就要完成 PreEnqueue 检查,确保被取出时已经满足入队条件;检查规则与首次入队相同。检查未通过则转入unschedulablePods;检查通过后,无需继续退避的 Pod 进入activeQ,仍需普通退避或错误退避的 Pod 分别进入podBackoffQ、podErrorBackoffQ。
- 响应集群事件:Node、Pod、PVC 等对象变化时,SchedulingQueue 调用此前拒绝该 Pod 的插件为该事件注册的 QueueingHint。相关回调均返回
PriorityQueue¶
Scheduler 创建过程中会调用 NewSchedulingQueue,得到默认实现 PriorityQueue。PriorityQueue 使用三个字段保存不同状态的 Pod:
activeQ:现在可以尝试调度的 Pod。backoffQ:保存已经满足重新入队条件、但仍处于退避期的 Pod。podBackoffQ保存插件拒绝后等待普通退避结束的 Pod;podErrorBackoffQ保存发生 Scheduler Error 后等待错误退避结束的 Pod。unschedulablePods:保存尚未满足重新入队条件的 Pod,等待可能改变上次调度结果的集群事件。
type PriorityQueue struct {
// ...
// activeQ 保存当前可以调度的 Pod。
activeQ activeQueuer
// backoffQ 保存仍处于退避时间内的 Pod。
backoffQ backoffQueuer
// unschedulablePods 等待可能改变调度结果的集群事件。
unschedulablePods *unschedulablePods
}
activeQ¶
activeQueue 使用 Heap 保存可立即调度的 Pod,并记录正在调度的 Pod、调度期间发生的事件和调度周期序号:
type activeQueue struct {
// 保护 activeQ、inFlightPods、inFlightEvents 和 schedCycle 等字段。
lock sync.RWMutex
// 保存可以立即调度的 Pod,Heap 队首是下一次取出的 Pod。
queue *heap.Heap[*framework.QueuedPodInfo]
// 包装同一个 Heap,供已经持有 lock 的代码调用。
unlockedQueue *unlockedActiveQueue
// 新 Pod 进入 activeQ 时唤醒正在等待的 Pop()。
cond sync.Cond
// 记录已经取出、但尚未结束本次调度处理的 Pod。
inFlightPods map[types.UID]*list.Element
// 记录这些 Pod 调度期间到达的集群事件,失败时用于判断是否立即重新入队。
inFlightEvents *list.List
// 每取出一个 Pod 增加一次,用于标识调度周期。
schedCycle int64
// Queue 关闭后,使正在等待的 Pop() 退出。
closed bool
isSchedulingQueueHintEnabled bool
metricsRecorder *metrics.MetricAsyncRecorder
// 启用 SchedulerPopFromBackoffQ 时,允许 activeQ 为空后从 podBackoffQ 取出 Pod。
backoffQPopper backoffQPopper
}
Scheduler 从 Profile 取得 QueueSortFunc,作为 lessFn 传入 NewPriorityQueue;该比较函数决定 Heap 的 Pod 顺序:
func NewPriorityQueue(
lessFn fwk.LessFunc,
informerFactory informers.SharedInformerFactory,
opts ...Option,
) *PriorityQueue {
// QueueSort 提供的 LessFunc 接收 fwk.QueuedPodInfo,
// 这里转换为内部 *framework.QueuedPodInfo 的比较函数。
lessConverted := convertLessFn(lessFn)
// activeQ 使用 Heap 保存 Pod,队首由 QueueSort 的比较结果决定。
pq.activeQ = newActiveQueue(
heap.NewWithRecorder(
podInfoKeyFunc,
heap.LessFunc[*framework.QueuedPodInfo](lessConverted),
metrics.NewActivePodsRecorder(),
),
isSchedulingQueueHintEnabled,
options.metricsRecorder,
backoffQPopper,
)
// ...
}
新 Pod、已完成退避的 Pod,以及被相关事件重新激活的 Pod 都可以进入 activeQ。PriorityQueue.Pop() 调用 activeQueue.pop 取出队首 Pod;随后由 unlockedMovePodToInFlight 增加调度尝试次数、记录 inFlightPods 并递增调度周期序号。
backoffQ¶
backoffQueue 以退避结束时间排序,并使用两个 Heap:
podBackoffQ:保存带有UnschedulablePlugins或PendingPlugins记录、并被安排进入退避的 Pod。普通路径主要来自Unschedulable失败;Pending插件关联的事件返回Queue时,Pod 会跳过退避并进入activeQ。podErrorBackoffQ:保存没有上述插件记录、因 Scheduler Error 进入错误退避的 Pod。
backoffQueue.add 根据失败记录选择内部队列:
type backoffQueue struct {
// 插件相关失败使用普通退避。
podBackoffQ *heap.Heap[*framework.QueuedPodInfo]
// Scheduler Error 使用单独的错误退避。
podErrorBackoffQ *heap.Heap[*framework.QueuedPodInfo]
// ...
}
func (bq *backoffQueue) add(logger klog.Logger, pInfo *framework.QueuedPodInfo, event string) {
// ...
// 没有记录 Unschedulable 或 Pending 插件,表示本次失败来自 Scheduler Error。
if pInfo.UnschedulablePlugins.Len() == 0 && pInfo.PendingPlugins.Len() == 0 {
bq.podErrorBackoffQ.AddOrUpdate(pInfo)
// ...
return
}
// 插件相关失败进入普通退避队列。
bq.podBackoffQ.AddOrUpdate(pInfo)
// ...
}
Kubernetes v1.36 默认启用 Beta 的 SchedulerPopFromBackoffQ。activeQ 为空时,activeQueue.unlockedPop 可以直接取出 podBackoffQ 的队首,即使它还没有到退避截止时间;podErrorBackoffQ 不会沿这条路径提前取出:
func (aq *activeQueue) unlockedPop(logger klog.Logger) (*framework.QueuedPodInfo, error) {
// activeQ 为空但 podBackoffQ 非空时,不再阻塞等待新 Pod。
for aq.queue.Len() == 0 {
if aq.backoffQPopper != nil && aq.backoffQPopper.lenBackoff() != 0 {
break
}
aq.cond.Wait()
}
pInfo, err := aq.queue.Pop()
if err != nil {
// activeQ 没有 Pod,直接从 podBackoffQ 取出队首。
pInfo, err = aq.backoffQPopper.popBackoff()
}
// ...
}
flushBackoffQCompleted 会把 podBackoffQ 和 podErrorBackoffQ 中已经完成退避的 Pod 移入 activeQ。podBackoffQ 和 podErrorBackoffQ 共用初始退避时间、最大退避时间和指数退避公式,但使用不同的失败计数。podBackoffQ 使用 UnschedulableCount,podErrorBackoffQ 使用 ConsecutiveErrorsCount。
calculateBackoffDuration 按 min(initialBackoff × 2^(失败次数-1), maxBackoff) 计算退避时间。默认初始值为 1 秒、最大值为 10 秒,因此两个队列的默认退避时间均依次为 1、2、4、8、10 秒,之后保持 10 秒。
unschedulablePods¶
unschedulablePods 使用 Map 保存当前无法调度、且尚未满足重试条件的 Pod:
type unschedulablePods struct {
// Pod 的 namespace/name 是 key,QueuedPodInfo 保存 Pod 及其失败信息。
podInfoMap map[string]*framework.QueuedPodInfo
keyFunc func(*v1.Pod) string
// 分别记录普通不可调度 Pod 和被 Scheduling Gate 阻塞的 Pod 数量。
unschedulableRecorder, gatedRecorder metrics.MetricRecorder
}
QueuedPodInfo 中的 UnschedulablePlugins 和 PendingPlugins 保存上次阻止调度继续的插件名。isPodWorthRequeuing 只调用这些插件关联的 QueueingHint,判断这次事件是否可能让该 Pod 重新具备调度条件:
func (p *PriorityQueue) isPodWorthRequeuing(
logger klog.Logger,
pInfo *framework.QueuedPodInfo,
event fwk.ClusterEvent,
oldObj, newObj interface{},
) queueingStrategy {
// 只检查上次阻止当前 Pod 继续调度的插件。
rejectorPlugins := pInfo.UnschedulablePlugins.Union(pInfo.PendingPlugins)
// ...
hintMap := p.queueingHintMap[pInfo.Pod.Spec.SchedulerName]
pod := pInfo.Pod
for eventToMatch, hintfns := range hintMap {
if !framework.MatchClusterEvents(eventToMatch, event) {
continue
}
for _, hintfn := range hintfns {
if !rejectorPlugins.Has(hintfn.PluginName) {
continue
}
hint, err := hintfn.QueueingHintFn(logger, pod, oldObj, newObj)
// 根据 QueueSkip 或 Queue 决定是否进入重新入队流程。
// ...
}
}
}
Pod 重新入队¶
对象变化到达后,MoveAllToActiveOrBackoffQueue 进入事件处理流程,movePodsToActiveOrBackoffQueue 使用 QueueingHint 判断 unschedulablePods 中的哪些 Pod 可以重新入队。requeuePodWithQueueingStrategy 再根据判断结果和剩余退避时间选择目标队列:
func (p *PriorityQueue) movePodsToActiveOrBackoffQueue(
logger klog.Logger,
podInfoList []*framework.QueuedPodInfo,
event fwk.ClusterEvent,
oldObj, newObj interface{},
) {
// 没有插件关注该事件时,不遍历 unschedulablePods。
if !p.isEventOfInterest(logger, event) {
return
}
activated := false
for _, pInfo := range podInfoList {
// 只运行上次阻止当前 Pod 调度的插件所关联的 QueueingHint。
schedulingHint := p.isPodWorthRequeuing(logger, pInfo, event, oldObj, newObj)
if schedulingHint == queueSkip {
// QueueingHint 未建议重新入队,Pod 继续留在 unschedulablePods。
continue
}
p.unschedulablePods.delete(pInfo.Pod, pInfo.Gated())
queue := p.requeuePodWithQueueingStrategy(logger, pInfo, schedulingHint, event.Label())
if queue == activeQ || (p.isPopFromBackoffQEnabled && queue == backoffQ) {
activated = true
}
}
if p.isSchedulingQueueHintEnabled {
// 同一事件也要记录给仍在调度中的 Pod,供其失败处理时检查。
p.activeQ.addEventIfAnyInFlight(oldObj, newObj, event)
}
if activated {
p.activeQ.broadcast()
}
}
func (p *PriorityQueue) requeuePodWithQueueingStrategy(
logger klog.Logger,
pInfo *framework.QueuedPodInfo,
strategy queueingStrategy,
event string,
) string {
// 没有事件可能改变结果,继续留在 unschedulablePods。
if strategy == queueSkip {
p.unschedulablePods.addOrUpdate(pInfo, pInfo.Gated(), event)
return unschedulableQ
}
// 需要遵守退避且退避尚未结束,进入 backoffQ。
if strategy == queueAfterBackoff && p.backoffQ.isPodBackingoff(pInfo) {
if added := p.moveToBackoffQ(logger, pInfo, event); added {
return backoffQ
}
return unschedulableQ
}
// Pending 插件要求立即重试,或者普通退避已经结束,进入 activeQ。
if added := p.moveToActiveQ(logger, pInfo, event, false); added {
return activeQ
}
// PreEnqueue 未通过时,moveToActiveQ 已将 Pod 放回 unschedulablePods。
return unschedulableQ
}
Pod 调度期间也可能发生相关集群事件。AddUnschedulableIfNotPresent 处理调度失败时,determineSchedulingHintForInFlightPod 会检查这些已经记录的事件,再使用同一个 requeuePodWithQueueingStrategy 选择目标队列:
func (p *PriorityQueue) AddUnschedulableIfNotPresent(
logger klog.Logger,
pInfo *framework.QueuedPodInfo,
podSchedulingCycle int64,
) error {
// ...
// 检查 Pod 从 Pop 到调度失败之间发生的事件,
// 避免错过已经到达的重新入队机会。
schedulingHint := p.determineSchedulingHintForInFlightPod(logger, pInfo)
queue := p.requeuePodWithQueueingStrategy(
logger, pInfo, schedulingHint, framework.ScheduleAttemptFailure,
)
// ...
return nil
}
Scheduler Cache¶
Cache 是 kube-scheduler 保存在内存中的本地集群状态。它把已经分配和暂时假定分配的 Pod 汇总到各个 Node 的 NodeInfo 中;每轮调度开始时,Scheduler 再从 Cache 增量刷新 Snapshot,节点筛选、评分和相关插件读取这份本轮视图。
相关代码位置如下:
pkg/scheduler/backend/cache/
├── interface.go # Cache 接口
├── cache.go # cacheImpl、内部索引和 UpdateSnapshot
├── node_tree.go # 按 Zone 组织 Node,并交错生成顺序,使优先检查的节点来自不同 Zone
└── snapshot.go # 本轮调度读取的 Snapshot
pkg/scheduler/framework/types.go # NodeInfo 的具体实现
pkg/scheduler/eventhandlers.go # 将已分配 Pod 和 Node 的变化写入 Cache
pkg/scheduler/schedule_one.go # 刷新 Snapshot,记录和撤销 Assumed Pod
pkg/scheduler/framework/runtime/
└── framework.go # 保存 Snapshot 的 SharedLister,运行扩展点插件
staging/src/k8s.io/kube-scheduler/framework/
├── listers.go # 定义 SharedLister 与 NodeInfoLister 查询接口
└── interface.go # 在 Handle 中暴露 SnapshotSharedLister()
下图展示 Scheduler Cache 的三条主要数据路径:更新 Cache、增量刷新 Snapshot,以及在调度周期内读取 Snapshot。
-
更新 Cache:Cache 有三类更新来源:
- 同步已分配 Pod 状态:Pod Informer 将 API Server 中已分配 Pod 的新增、变化和删除同步到 Cache。
Cache.AddPod()、UpdatePod()和RemovePod()更新podStates以及所在 Node 的NodeInfo,使资源、端口和 PVC 等占用与集群状态一致;AddPod()观察到此前 Assumed 的 Pod 时,会用已确认的对象替换临时记录,而不会重复计算占用。 - 同步 Node 状态:Node Informer 把 Node 的新增、属性变化和删除同步到 Cache。
Cache.AddNode()、UpdateNode()和RemoveNode()更新NodeInfo中的资源容量、标签、污点和拓扑等信息,并维护节点集合与遍历顺序,使节点筛选使用最新状态。 - 记录调度中的临时占用:Scheduler 选出目标 Node 后、异步绑定完成前,
Cache.AssumePod()先把 Pod 计入podStates、assumedPods和目标NodeInfo,避免后续 Pod 再次使用同一份资源。后续阶段失败时,Cache.ForgetPod()撤销这些记录;绑定成功后,Pod Informer 再通过AddPod()将临时记录更新为已确认状态。
- 同步已分配 Pod 状态:Pod Informer 将 API Server 中已分配 Pod 的新增、变化和删除同步到 Cache。
-
增量刷新 Snapshot:上述操作只更新 Cache。每轮调度开始时,
schedulingCycle()调用Cache.UpdateSnapshot(),把变化的NodeInfo和相关索引增量刷新到 Snapshot。 - 调度期间读取 Snapshot:
UpdateSnapshot()完成后,schedulingCycle()调用schedulingAlgorithm(),为当前 Pod 选择目标 Node。schedulingAlgorithm()通过schedulePod()从 Snapshot 读取本轮NodeInfo,完成节点筛选与评分并返回选中的 Node。调度期间的插件也可以通过Handle.SnapshotSharedLister()查询同一份 Snapshot。
Cache 数据结构¶
默认实现 cacheImpl 使用 Pod 索引、Node 索引和更新时间链表维护持续变化的集群状态。下面只保留普通 Pod 调度路径直接使用的字段:
// 每个链表项保存一个 NodeInfo;越靠近 headNode,更新时间越晚。
type nodeInfoListItem struct {
info *framework.NodeInfo
next *nodeInfoListItem
prev *nodeInfoListItem
}
type cacheImpl struct {
mu sync.RWMutex
// Pod UID → Pod 状态;assumedPods 是其中 Assumed Pod 的 UID 集合。
podStates map[string]*podState
assumedPods sets.Set[string]
// Node 名称索引,以及按更新时间排列的链表入口。
nodes map[string]*nodeInfoListItem
headNode *nodeInfoListItem
// 按 Zone 保存 Node 名称;需要重建 Snapshot 节点列表时,交错生成跨 Zone 的遍历顺序。
nodeTree *nodeTree
// 汇总镜像所在的 Node,供各 NodeInfo 建立 ImageStates。
imageStates map[string]*fwk.ImageStateSummary
// ...
}
type podState struct {
pod *v1.Pod
}
Pod 缓存状态¶
Pod 缓存状态有两个更新来源:
- Pod Informer:Pod Informer 接收所有 Pod 变化,注册的
addPod()、updatePod()和deletePod()根据spec.nodeName选择处理路径。由当前 Scheduler 负责的未分配 Pod 进入 SchedulingQueue;已分配 Pod 通过Cache.AddPod()、Cache.UpdatePod()或Cache.RemovePod()写入 Cache。 - Scheduler 调度流程:Scheduler 选出目标 Node 后调用
Cache.AssumePod(),先记录尚未写入 API Server 的本地占用;后续阶段失败时调用Cache.ForgetPod()撤销这次占用。
Cache 没有分别定义 Assumed 和 Added 类型。每个缓存中的 Pod 都保存在 podStates:如果它的 UID 还在 assumedPods 中,它就是 Assumed;如果 UID 不在 assumedPods 中,它就是 Added。两种状态下,Pod 都已经计入目标 NodeInfo。
例如,Scheduler 为 Pod train-worker-0(UID 为 pod-123)选中 node-a。下面用简化结构展示绑定结果同步前后的变化:
AssumePod() 后:
podStates
└─ pod-123 → podState{pod: Pod{name: train-worker-0, spec.nodeName: node-a}}
assumedPods
└─ pod-123
Pod Informer 调用 Cache.AddPod() 后:
podStates
└─ pod-123 → podState{pod: API Server 中的 train-worker-0}
assumedPods
└─ (空)
AddPod() 先从 node-a 的 NodeInfo 和 podStates 移除本地假定的对象,再加入 API Server 返回的对象,同时删除 assumedPods 中的 pod-123。转换完成后,NodeInfo 中仍然只有一个 train-worker-0,它的资源只计算一次。
- 进入 Assumed:Scheduler 选出目标 Node 后,
AssumePod()立即把 Pod 写入podStates和assumedPods,并将其资源计入目标NodeInfo。 - 直接进入 Added:Scheduler 启动时观察到的已分配 Pod,以及由其他组件完成绑定的 Pod,会通过
AddPod()直接写入podStates和目标NodeInfo。 - Assumed 转为 Added:API Server 保存绑定结果后,Informer 送达已分配 Pod。
AddPod()使用观察到的对象更新缓存,并删除assumedPods中的 UID。 - 撤销 Assumed:Reserve、Permit、PreBind 或 Bind 等后续阶段失败时,
ForgetPod()删除缓存记录并撤销目标 Node 上的假定占用。 - 更新 Added:
UpdatePod()从NodeInfo移除 old Pod,再加入 new Pod。Assumed Pod 必须先收到 Add 事件转为 Added,才能处理 Update 事件。 - 移除 Added:
RemovePod()使用podStates中保存的对象扣除 Node 聚合值,再删除podStates和assumedPods中的记录。
cacheImpl.addPod() 同时更新目标 NodeInfo、更新时间链表和 Pod 索引。下面省略 PodGroup 相关逻辑;在普通 Pod 缓存路径中,assumePod 决定是否把 UID 额外写入 assumedPods:
func (cache *cacheImpl) addPod(logger klog.Logger, pod *v1.Pod, assumePod bool) error {
key, err := framework.GetPodKey(pod) // 使用 Pod UID 作为 Cache key
if err != nil {
return err
}
// 按目标 Node 名称查找对应的链表项。
n, ok := cache.nodes[pod.Spec.NodeName]
if !ok {
// 目标 Node 尚无缓存项时,先创建空的 NodeInfo;随后写入 Pod,Node 对象由 Node 事件补充。
n = newNodeInfoListItem(framework.NewNodeInfo())
cache.nodes[pod.Spec.NodeName] = n
}
// 将 Pod 及其资源、端口和 PVC 等计入目标 NodeInfo。
n.info.AddPod(pod)
// 把变化的 NodeInfo 移到链表头,供 UpdateSnapshot 增量扫描。
cache.moveNodeInfoToHead(logger, pod.Spec.NodeName)
// 所有 Assumed 和 Added Pod 都保存在 podStates。
cache.podStates[key] = &podState{pod: pod}
if assumePod {
// 只有 Assumed Pod 的 UID 还会写入 assumedPods。
cache.assumedPods.Insert(key)
}
// ...
return nil
}
Node 状态与索引¶
NodeInfo 是 Scheduler Cache 为单个 Node 维护的调度状态汇总。它集中保存 Node 对象、该 Node 上已分配和 Assumed 的 Pod,以及资源请求、可分配资源、端口、PVC、亲和性和镜像等调度信息;Snapshot 复制这些状态,供节点筛选、评分和相关插件读取。
NodeInfo 会被两类数据变化更新:
- Node 状态变化:Node Informer 通过
addNodeToCache()、updateNodeInCache()和deleteNodeFromCache(),分别调用Cache.AddNode()、Cache.UpdateNode()和Cache.RemoveNode(),把 API Server 中 Node 的新增、属性变化和删除同步到 Scheduler Cache。 - Pod 状态变化:Pod Informer 调用
Cache.AddPod()、Cache.UpdatePod()和Cache.RemovePod(),Scheduler 调度流程调用Cache.AssumePod()和Cache.ForgetPod()。这些方法更新目标NodeInfo中的 Pod,以及资源请求、端口、PVC 和亲和性 Pod 子集,使后续调度看到该 Node 当前已经占用的资源。
下图展示 Scheduler Cache 如何维护 Node 的名称索引和更新时间链表,以及如何利用这些结构增量更新 Snapshot;nodeTree 则负责生成按 Zone 交错排列的节点顺序。
-
nodes按名称定位:nodes[nodeName]直接返回对应的nodeInfoListItem。图中nodes["node-b"]指向node-b的链表项。 -
headNode按更新时间定位:headNode指向最近变化的链表项。NodeInfo发生变化后,moveNodeInfoToHead()会把对应的链表项移到表头,使链表从新到旧排列。沿next访问更早变化的项,沿prev返回更新的项。这样排列是为了减少 Snapshot 的刷新范围。UpdateSnapshot()从headNode开始扫描,遇到已经同步过的Generation后便可以停止,不必每轮遍历所有 Node。 -
链表项保存
NodeInfo:每个nodeInfoListItem的info都指向对应 Node 的NodeInfo。NodeInfo汇总 Node 对象、已分配及 Assumed Pod、资源、端口、PVC、亲和性 Pod 子集和Generation等调度状态。 -
增量更新
nodeInfoMap:Cache.UpdateSnapshot(snapshot)主动从headNode开始,沿next按更新时间从新到旧扫描。对于Generation > Snapshot.generation的链表项,它复制其中的NodeInfo并更新Snapshot.nodeInfoMap;遇到第一个不满足条件的链表项后结束扫描。 -
nodeTree生成 Node 名称顺序:Node 新增或删除时,nodeTree.list()先取各 Zone 的第一个 Node,再取各 Zone 的第二个 Node,生成交错排列的 Node 名称。Scheduler 筛选节点时,找到足够数量的可行节点后便可能提前停止,不一定检查所有 Node。交错排列使前面检查的节点尽量来自不同 Zone,减少候选节点集中在同一个 Zone 的情况。 -
重建
nodeInfoList:UpdateSnapshot()按nodeTree.list()返回的顺序查询nodeInfoMap[nodeName],并将同一个NodeInfo依次加入nodeInfoList。nodeInfoList是本轮调度的候选节点列表。Scheduler 逐个读取其中的NodeInfo,交给 Filter 插件判断当前 Pod 能否在对应 Node 上运行;通过检查的NodeInfo组成可行节点列表,存在多个可行节点时再由 Score 插件评分并选出目标 Node。
Snapshot 增量更新与读取¶
Snapshot 保存调度周期读取的 NodeInfo 及其索引。Scheduler 持有同一个 nodeInfoSnapshot,每轮把它传给 Cache.UpdateSnapshot() 增量刷新,而不是重新创建完整副本。
type Snapshot struct {
// NodeInfo 索引:按名称查找,以及按 nodeTree 顺序遍历。
nodeInfoMap map[string]*framework.NodeInfo
nodeInfoList []fwk.NodeInfo
// 亲和性和 PVC 查询使用的辅助索引。
havePodsWithAffinityNodeInfoList []fwk.NodeInfo
havePodsWithRequiredAntiAffinityNodeInfoList []fwk.NodeInfo
usedPVCSet sets.Set[string]
// 上次已同步到的最新 NodeInfo.Generation。
generation int64
// ...
}
下图展示 UpdateSnapshot() 如何根据 generation 扫描变化节点、更新 Snapshot 内部索引,以及调度算法如何读取刷新后的 NodeInfo:
- 进入调度周期:
schedulingCycle()在运行调度算法前调用Cache.UpdateSnapshot(),将 Cache 中上轮之后发生变化的节点状态同步到 Snapshot,供本轮 Filter 和 Score 使用同一份 NodeInfo 视图。 - 扫描变化节点:
UpdateSnapshot()从headNode向后扫描。当前 NodeInfo 的 generation 大于旧Snapshot.generation时继续;遇到不大于该值的链表项后停止。 - 复制最终状态:
NodeInfo.SnapshotConcrete()复制发生变化的 NodeInfo,并更新nodeInfoMap[nodeName]。同一 Node 在两次刷新之间变化多次,也只复制刷新时的最终状态。 - 更新 Snapshot 索引与同步位置:扫描结束后,若
headNode非空,UpdateSnapshot()先将Snapshot.generation更新为当前headNode.info.Generation,作为下一轮判断 NodeInfo 是否已同步的基准值。随后按需清理已删除的 Node;Node 新增或删除时,UpdateSnapshot()按nodeTree.list()返回的名称顺序重建nodeInfoList。如果本轮节点变化涉及 Pod 亲和性或 PVC 使用情况,则重新生成havePodsWithAffinityNodeInfoList、havePodsWithRequiredAntiAffinityNodeInfoList和usedPVCSet。 - 读取节点视图:
schedulingCycle()刷新 Snapshot 后调用schedulingAlgorithm(),后者通过schedulePod()读取本轮nodeInfoList,完成节点筛选与评分。Framework 再按照扩展点需要,将单个NodeInfo或NodeInfo列表传给插件。
调度阶段(scheduling cycle)¶
ScheduleOne 每次从队列取出一个 Pod。普通 Pod 进入 scheduleOnePod(),同步执行 schedulingCycle();该阶段刷新 Snapshot、选择目标 Node,并完成 Assume、Reserve 和 Permit。调度阶段成功后,scheduleOnePod() 再启动异步绑定。
相关代码位置如下:
pkg/scheduler/
├── schedule_one.go # ScheduleOne、节点选择算法与绑定前准备
├── backend/
│ ├── queue/scheduling_queue.go # Pop:取得下一次调度要处理的 Pod
│ └── cache/cache.go # UpdateSnapshot 刷新视图,AssumePod 记录占用
└── framework/runtime/framework.go # 执行筛选、评分、Reserve 和 Permit 插件
下图展示普通 Pod 从出队到启动绑定的路径,包含节点筛选、评分、绑定前准备和失败处理。
- 取出待调度 Pod:
ScheduleOne()调用NextPod,后者默认通过PriorityQueue.Pop取得下一个QueuedPodInfo,开始一次 Pod 调度尝试。 - 建立本轮调度上下文:
scheduleOnePod()根据schedulerName选择 Framework,并创建当前 Pod 独享的CycleState,随后同步调用schedulingCycle()。 - 执行调度周期:
schedulingCycle()先调用Cache.UpdateSnapshot()刷新本轮节点视图,再依次运行节点选择算法和绑定前准备。任一步骤失败都会把状态返回给scheduleOnePod(),由它调用FailureHandler。 - 选择目标 Node:
schedulingAlgorithm()调用SchedulePod,其默认实现schedulePod()运行 PreFilter 和 Filter;存在多个可行 Node 时,再运行 PreScore 和 Score 选出目标 Node。无法选出 Node 或发生错误时,本轮调度结束。 - 准备绑定状态:
prepareForBindingCycle()先通过 Assume 将 Pod 的资源占用写入 Scheduler Cache,再运行 Reserve 和 Permit。Reserve 或 Permit 失败时调用 Unreserve,并通过ForgetPod()撤销本地假定;Permit 返回Wait时,Scheduler 先记录 WaitingPod。 - 启动异步绑定:Permit 返回
Success或Wait后,schedulingCycle()都会成功返回,scheduleOnePod()随即启动runBindingCycle()goroutine。Wait状态会在该 goroutine 内继续等待 Permit 插件放行。
ScheduleOne¶
ScheduleOne 是一次调度的入口。它通过 NextPod 从 SchedulingQueue 取出队首 Pod;NextPod 默认调用 PriorityQueue.Pop,队列为空时会等待。启用 GenericWorkload 且 Pod 设置了 spec.schedulingGroup 时进入 PodGroup 调度路径,其余 Pod 交给 scheduleOnePod。
func (sched *Scheduler) ScheduleOne(ctx context.Context) {
// NextPod 默认从 SchedulingQueue 取出队首 Pod。
// 队列为空时等待;队列关闭时可能返回 nil。
podInfo, err := sched.NextPod(klog.FromContext(ctx))
if err != nil || podInfo == nil || podInfo.Pod == nil {
return
}
// 启用 GenericWorkload 且 Pod 指定 SchedulingGroup 时,
// 先取得组信息,再进入 PodGroup 调度流程。
if sched.genericWorkloadEnabled && podInfo.Pod.Spec.SchedulingGroup != nil {
podGroupInfo, err := sched.podGroupInfoForPod(ctx, podInfo)
if err != nil {
...
return
}
sched.scheduleOnePodGroup(ctx, podGroupInfo)
} else {
// 其余 Pod 使用单 Pod 调度流程。
sched.scheduleOnePod(ctx, podInfo)
}
}
scheduleOnePod¶
scheduleOnePod 负责组织一个普通 Pod 的完整调度尝试。它根据 spec.schedulerName 取得对应的 Framework,创建本次调度使用的 CycleState,然后同步执行调度阶段。调度失败时交给 FailureHandler;成功时,Pod 已在 Cache 中完成 Assume,因此绑定阶段可以在独立 goroutine 中运行。
func (sched *Scheduler) scheduleOnePod(
ctx context.Context,
podInfo *framework.QueuedPodInfo,
) {
pod := podInfo.Pod
...
fwk, err := sched.frameworkForPod(pod)
if err != nil {
sched.SchedulingQueue.Done(pod.UID)
return
}
if sched.skipPodSchedule(ctx, fwk, pod) {
sched.SchedulingQueue.Done(pod.UID)
return
}
// 保存本次调度中各扩展点插件共享的临时状态;
// 调度阶段和后续绑定阶段使用同一个 CycleState。
state := framework.NewCycleState()
...
// 调度阶段同步执行,完成 Snapshot 刷新、节点选择、
// Assume、Reserve 和 Permit。
schedulingCycleCtx, cancel := context.WithCancel(ctx)
defer cancel()
scheduleResult, assumedPodInfo, status := sched.schedulingCycle(
schedulingCycleCtx, state, fwk, podInfo, start, podsToActivate,
)
if !status.IsSuccess() {
// 记录失败原因,并决定 Pod 后续进入哪个重试队列。
sched.FailureHandler(
schedulingCycleCtx, fwk, assumedPodInfo,
status, scheduleResult.nominatingInfo, start,
)
return
}
// Pod 已通过 Assume 计入 Scheduler Cache。
// 绑定可以异步执行,主调度循环继续处理下一个 Pod。
go sched.runBindingCycle(
ctx, state, fwk, scheduleResult, assumedPodInfo, start, podsToActivate,
)
}
schedulingCycle¶
schedulingCycle 是同步调度阶段。它先刷新 Snapshot,再从这份固定视图中选择目标 Node;选中后调用 prepareForBindingCycle,把 Pod 临时计入目标 Node,并运行 Reserve 和 Permit 插件。
func (sched *Scheduler) schedulingCycle(...) (
ScheduleResult, *framework.QueuedPodInfo, *fwk.Status,
) {
// 把 Cache 中发生变化的 NodeInfo 增量更新到 Snapshot。
// 本轮 Filter 和 Score 都读取刷新后的同一份节点视图。
if err := sched.Cache.UpdateSnapshot(
klog.FromContext(ctx), sched.nodeInfoSnapshot,
); err != nil {
return ScheduleResult{nominatingInfo: clearNominatedNode},
podInfo, fwk.AsStatus(err)
}
// 运行节点筛选和评分,返回选中的目标 Node。
result, status := sched.schedulingAlgorithm(
ctx, state, schedFramework, podInfo, start,
)
if !status.IsSuccess() {
return result, podInfo, status
}
// 在真正绑定前应用调度结果:设置目标 Node、执行
// Cache.AssumePod(),再运行 Reserve 和 Permit。
assumedPodInfo, status := sched.prepareForBindingCycle(
ctx, state, schedFramework, podInfo, podsToActivate, result,
)
if !status.IsSuccess() {
return ScheduleResult{nominatingInfo: clearNominatedNode},
assumedPodInfo, status
}
// 返回已完成 Assume 的 Pod,交给异步绑定阶段继续处理。
return result, assumedPodInfo, nil
}
schedulingAlgorithm¶
schedulingAlgorithm 先调用 SchedulePod 运行 PreFilter 和 Filter;需要对多个可行 Node 评分时,再运行 PreScore 和 Score。没有可行 Node 时,schedulingAlgorithm 会在 SchedulePod 返回后运行 PostFilter。Reserve、Permit 和绑定相关扩展点由后续阶段执行。PostFilter 即使找到候选方案,本次调度仍返回不可调度,Pod 会在后续周期重新尝试。
func (sched *Scheduler) schedulingAlgorithm(...) (
ScheduleResult, *fwk.Status,
) {
// 从本轮 Snapshot 中筛选可行 Node,并通过评分选出目标 Node。
result, err := sched.SchedulePod(
ctx, schedFramework, state, podInfo,
)
if err == nil {
return result, nil
}
// Snapshot 中没有可供调度的 Node。
if err == ErrNoNodesAvailable {
status := fwk.NewStatus(fwk.UnschedulableAndUnresolvable).WithError(err)
return ScheduleResult{nominatingInfo: clearNominatedNode}, status
}
// FitError 包含各 Node 未通过 Filter 的原因。
fitError, ok := err.(*framework.FitError)
if !ok {
return ScheduleResult{nominatingInfo: clearNominatedNode},
fwk.AsStatus(err)
}
if !schedFramework.HasPostFilterPlugins() {
return ScheduleResult{nominatingInfo: clearNominatedNode},
fwk.NewStatus(fwk.Unschedulable).WithError(err)
}
// PostFilter 根据 Filter 的失败结果尝试改善后续调度条件,
// 例如 DefaultPreemption 可以选择一个抢占候选 Node。
postFilterResult, status := schedFramework.RunPostFilterPlugins(
ctx, state, podInfo.Pod, fitError.Diagnosis.NodeToStatus,
)
...
var nominatingInfo *fwk.NominatingInfo
if postFilterResult != nil {
// 保存候选 Node,供 Pod 后续调度周期使用。
nominatingInfo = postFilterResult.NominatingInfo
}
...
// 当前调度仍以不可调度结束,不会进入绑定阶段。
return ScheduleResult{nominatingInfo: nominatingInfo},
fwk.NewStatus(fwk.Unschedulable).WithError(err)
}
Filter 与 Score 的节点并行¶
SchedulePod 的默认实现 schedulePod() 先通过 findNodesThatPassFilters 筛选候选 Node;得到多台可行 Node 时,再由 prioritizeNodes() 依次运行 PreScore 和 Score。Filter 和 Score 都会并行处理 Node,但并行入口不同。
Filter:Node 并行由 findNodesThatPassFilters 发起
Score:Node 并行由 RunScorePlugins 发起
Filter 并行源码
findNodesThatPassFilters 先计算本轮需要找到多少台可行 Node,再把“检查一台 Node”封装成 checkNode。Parallelizer 并行执行这些工作项;每台检查通过的 Node 通过原子计数取得结果数组中的位置,找到足够数量后取消剩余工作:
func (sched *Scheduler) findNodesThatPassFilters(
ctx context.Context,
schedFramework framework.Framework,
state fwk.CycleState,
pod *v1.Pod,
diagnosis *framework.Diagnosis,
nodes []fwk.NodeInfo,
) ([]fwk.NodeInfo, error) {
numAllNodes := len(nodes)
numNodesToFind := sched.numFeasibleNodesToFind(
schedFramework.PercentageOfNodesToScore(),
int32(numAllNodes),
)
feasibleNodes := make([]fwk.NodeInfo, numNodesToFind)
...
errCh := parallelize.NewResultChannel[error]()
var feasibleNodesLen int32
ctx, cancel := context.WithCancelCause(ctx)
defer cancel(errors.New("findNodesThatPassFilters has completed"))
type nodeStatus struct {
node string
status *fwk.Status
}
result := make([]*nodeStatus, numAllNodes)
// checkNode 只负责检查一台 Node。
checkNode := func(i int) {
// 从上一轮停止的位置继续,避免每轮总从相同 Node 开始。
nodeInfo := nodes[(sched.nextStartNodeIndex+i)%numAllNodes]
// 该方法处理 Nominated Pod 后,调用 RunFilterPlugins 检查当前 Node。
status := schedFramework.RunFilterPluginsWithNominatedPods(
ctx, state, pod, nodeInfo,
)
if status.Code() == fwk.Error {
errCh.SendWithCancel(status.AsError(), func() {
cancel(errors.New("some other Filter operation failed"))
})
return
}
if status.IsSuccess() {
// 多个工作项并行写入结果,使用原子计数分配唯一位置。
length := atomic.AddInt32(&feasibleNodesLen, 1)
if length > numNodesToFind {
// 已找到足够的可行 Node,停止尚未完成的检查。
cancel(errors.New("findNodesThatPassFilters has found enough nodes"))
atomic.AddInt32(&feasibleNodesLen, -1)
} else {
feasibleNodes[length-1] = nodeInfo
}
} else {
// 每个工作项只写入自己的下标,不需要额外加锁。
result[i] = &nodeStatus{
node: nodeInfo.Node().Name,
status: status,
}
}
}
// 将全部候选 Node 索引作为工作项交给 Parallelizer。
schedFramework.Parallelizer().Until(
ctx, numAllNodes, checkNode, metrics.Filter,
)
feasibleNodes = feasibleNodes[:feasibleNodesLen]
...
return feasibleNodes, nil
}
Score 并行源码
RunScorePlugins() 接收多台可行 Node,并执行三轮并行处理:先按 Node 计算各插件的原始分数,再按插件运行可选的 NormalizeScore,最后按 Node 应用插件权重并汇总总分。第一轮中,不同 Node 可以并行处理,但同一台 Node 上的 Score 插件仍按 Profile 中的配置顺序运行。
func (f *frameworkImpl) RunScorePlugins(
ctx context.Context,
state fwk.CycleState,
pod *v1.Pod,
nodes []fwk.NodeInfo,
) (ns []fwk.NodePluginScores, status *fwk.Status) {
...
// len(nodes) 个索引是可以并行执行的工作项。
f.Parallelizer().Until(ctx, len(nodes), func(index int) {
nodeInfo := nodes[index]
nodeName := nodeInfo.Node().Name
// 同一台 Node 内,各 Score 插件仍按 Profile 配置顺序运行。
for _, pl := range plugins {
score, status := f.runScorePlugin(
ctx, pl, state, pod, nodeInfo,
)
if !status.IsSuccess() {
errCh.SendWithCancel(
fmt.Errorf("plugin %q failed with: %w", pl.Name(), status.AsError()),
cancel,
)
return
}
// 保存当前插件为当前 Node 生成的原始分数。
pluginToNodeScores[pl.Name()][index] = fwk.NodeScore{
Name: nodeName,
Score: score,
}
}
}, metrics.Score)
...
// 每个 Score 插件独立归一化自己为所有 Node 生成的分数。
f.Parallelizer().Until(ctx, len(plugins), func(index int) {
pl := plugins[index]
if pl.ScoreExtensions() == nil {
return
}
nodeScoreList := pluginToNodeScores[pl.Name()]
// runScoreExtension 调用 pl.ScoreExtensions().NormalizeScore(...),
// 对当前插件为所有候选 Node 生成的分数进行归一化。
status := f.runScoreExtension(ctx, pl, state, pod, nodeScoreList)
if !status.IsSuccess() {
errCh.SendWithCancel(
fmt.Errorf("plugin %q failed with: %w", pl.Name(), status.AsError()),
cancel,
)
}
}, metrics.Score)
...
// 再次并行处理各台 Node,校验分数、应用权重并计算总分。
f.Parallelizer().Until(ctx, len(nodes), func(index int) {
nodePluginScores := fwk.NodePluginScores{
Name: nodes[index].Node().Name,
Scores: make([]fwk.PluginScore, len(plugins)),
}
for i, pl := range plugins {
// scorePluginWeight 只用于 Score 扩展点;
// 同一插件实现的 Filter、Reserve 等扩展点不参与加权。
weight := f.scorePluginWeight[pl.Name()]
score := pluginToNodeScores[pl.Name()][index].Score
// 插件最终提供的分数必须位于 0–100。
if score > fwk.MaxNodeScore || score < fwk.MinNodeScore {
errCh.SendWithCancel(
fmt.Errorf("plugin %q returns an invalid score %v", pl.Name(), score),
cancel,
)
return
}
weightedScore := score * int64(weight)
nodePluginScores.Scores[i] = fwk.PluginScore{
Name: pl.Name(),
Score: weightedScore,
}
nodePluginScores.TotalScore += weightedScore
}
allNodePluginScores[index] = nodePluginScores
}, metrics.Score)
...
return allNodePluginScores, nil
}
prepareForBindingCycle¶
prepareForBindingCycle 位于节点选择和异步绑定之间。它先调用 assumeAndReserve():通过 Cache.AssumePod() 将 Pod 临时计入目标 Node 的资源占用,再运行 Reserve 插件保存卷预绑定等可回滚状态。Scheduler 因此无需等待 API Server 完成绑定,就能继续调度其他 Pod;Reserve 失败时,Scheduler 执行 Unreserve 和 ForgetPod(),撤销已经写入的临时状态。
Assume 和 Reserve 成功后,Scheduler 运行 Permit 插件。Permit 返回 Success 时可以直接进入绑定;返回 Wait 时,Framework 将 Pod、等待它的插件和超时时间保存为 WaitingPod。这两种结果都会启动异步绑定阶段,其中 WaitOnPermit() 负责等待插件放行。Permit 直接拒绝或返回错误时,Scheduler 在返回失败前执行 Unreserve 和 ForgetPod(),清理 Reserve 状态与 Cache 中的临时占用。
func (sched *Scheduler) prepareForBindingCycle(...) (
*framework.QueuedPodInfo, *fwk.Status,
) {
// 先创建 Assumed Pod:在副本中写入目标 Node 并计入 Cache,
// 再让 Reserve 插件保存绑定阶段需要的可回滚状态。
assumedPodInfo, status := sched.assumeAndReserve(
ctx, state, schedFramework, podInfo, scheduleResult,
)
if !status.IsSuccess() {
// Assume 或 Reserve 未完成,停止当前调度周期。
// Reserve 失败时,assumeAndReserve 已执行 Unreserve 和 ForgetPod。
return assumedPodInfo, status
}
// Permit 插件在进入异步绑定阶段前决定立即放行、等待或拒绝。
pluginsWaitTime, runPermitStatus := schedFramework.RunPermitPlugins(
ctx, state, assumedPodInfo.Pod, scheduleResult.SuggestedHost,
)
if runPermitStatus.IsWait() {
// Wait 表示至少一个 Permit 插件需要异步决定。
// WaitingPod 保存等待该 Pod 的插件及超时时间;绑定阶段随后
// 通过 WaitOnPermit 等待全部插件放行,或在拒绝、超时后失败。
schedFramework.AddWaitingPod(assumedPodInfo.Pod, pluginsWaitTime)
} else if !runPermitStatus.IsSuccess() {
// Permit 直接拒绝或返回错误,绑定阶段不会启动。
// 先清理 Reserve 插件状态,再从 Cache 撤销 Assumed Pod。
_ = sched.unreserveAndForget(
ctx, state, schedFramework,
assumedPodInfo, scheduleResult.SuggestedHost,
)
if runPermitStatus.IsRejected() {
// 将所选 Node 上的拒绝结果包装为 FitError,
// 供 FailureHandler 记录失败原因并处理重新入队。
fitErr := &framework.FitError{...}
...
return assumedPodInfo,
fwk.NewStatus(runPermitStatus.Code()).WithError(fitErr)
}
return assumedPodInfo, runPermitStatus
}
...
// Success 可以直接进入绑定;Wait 已登记 WaitingPod,绑定阶段会等待。
// 两种情况都将 Assumed Pod 副本交给异步绑定阶段。
return assumedPodInfo, nil
}
绑定阶段(binding cycle)¶
绑定阶段由独立 goroutine 执行。多个已经完成调度阶段的 Pod 可以并发运行 PreBind、Bind 和 PostBind;Scheduler 的主循环同时继续为后续 Pod 选择 Node。
相关代码位置如下:
pkg/scheduler/
├── schedule_one.go # runBindingCycle、bindingCycle 与 bind
├── backend/queue/active_queue.go # 结束本轮 Pod 跟踪,清理不再需要的事件
└── framework/
├── runtime/framework.go # 等待 Permit,运行绑定阶段的扩展点插件
└── plugins/defaultbinder/default_binder.go # DefaultBinder 提交 pods/binding 请求
下图展示一个 Pod 的绑定流程。
- 启动绑定 goroutine:
scheduleOnePod在调度阶段成功后启动runBindingCycle。 - 执行绑定周期:
bindingCycle负责组织 PreBindPreFlight、WaitOnPermit、PreBind、Bind 和 PostBind;只有其中返回失败时,runBindingCycle才调用handleBindingCycleError。 - 收集 PreBind 计划:对应功能启用时,PreBindPreFlight 判断哪些 PreBind 插件需要运行,并记录可并行执行的插件。
- 等待 Permit 结果:
WaitOnPermit读取上一阶段登记的 WaitingPod。没有 WaitingPod 时立即成功;存在 WaitingPod 时等待插件放行、拒绝或超时。 - 结束队列事件跟踪:Pod 从
SchedulingQueue取出后,队列会暂存其调度期间发生的集群事件,供失败时判断是否立即重试。Permit 通过后,SchedulingQueue.Done()结束对该 Pod 的跟踪,并清理不再被其他调度中 Pod 使用的事件,避免事件列表持续增长。 - 完成绑定前准备:PreBind 插件执行卷绑定等 API Binding 前的准备工作。
- 运行绑定逻辑:
sched.bind()调用 Framework 的 Bind 插件,默认由DefaultBinder向 API Server 提交绑定请求。 - 提交 API Binding:默认的 Bind 插件请求 API Server 的
pods/binding子资源,将目标 Node 保存到spec.nodeName。 - 运行 PostBind 插件:Pod 成功绑定到 Node 后,Framework 依次调用所有已配置的
PostBind插件,用于清理本轮调度的临时状态或记录绑定结果。PostBind()没有返回值,不会改变已经完成的绑定。 - 进入失败处理:PreBindPreFlight、WaitOnPermit、PreBind 或 Bind 失败时,
runBindingCycle调用handleBindingCycleError()。PostBind 不返回状态,不进入这条失败分支。 - 回滚 Reserve 状态:
RunReservePluginsUnreserve撤销 Reserve 插件写入的临时状态。 - 撤销本地占用:
Cache.ForgetPod()从 Scheduler Cache 移除 Assumed Pod 的资源占用。 - 处理调度失败:FailureHandler 记录失败结果,并根据最新的 Pod 状态决定是否重新入队。
runBindingCycle¶
runBindingCycle() 调用 bindingCycle() 执行绑定流程。PreBindPreFlight、WaitOnPermit、PreBind 或 Bind 返回失败时,它会调用 handleBindingCycleError(),撤销 Reserve 状态和 Cache 中的 Assumed Pod,再把 Pod 交给 FailureHandler:
func (sched *Scheduler) runBindingCycle(...) {
bindingCycleCtx, cancel := context.WithCancel(ctx)
defer cancel()
...
status := sched.bindingCycle(
bindingCycleCtx, state, schedFramework,
scheduleResult, assumedPodInfo, start, podsToActivate,
)
if !status.IsSuccess() {
sched.handleBindingCycleError(
bindingCycleCtx, state, schedFramework,
assumedPodInfo, start, scheduleResult, status,
)
return
}
}
bindingCycle¶
bindingCycle 在对应功能启用时先运行 PreBindPreFlight,收集当前 Pod 需要执行的 PreBind 工作及其并行能力;再处理可能存在的 Permit 等待,最后执行 PreBind、Bind 和 PostBind。
func (sched *Scheduler) bindingCycle(...) *fwk.Status {
assumedPod := assumedPodInfo.Pod
var preFlightStatus *fwk.Status
if sched.nominatedNodeNameForExpectationEnabled {
// 收集当前 Pod 的 PreBind 执行计划。
preFlightStatus = schedFramework.RunPreBindPreFlights(
ctx, state, assumedPod, scheduleResult.SuggestedHost,
)
if preFlightStatus.Code() == fwk.Error || preFlightStatus.IsRejected() {
return preFlightStatus
}
...
}
// Permit 没有返回 Wait 时,这里立即通过。
if status := schedFramework.WaitOnPermit(ctx, assumedPod); !status.IsSuccess() {
if status.IsRejected() {
// 保存目标 Node 和拒绝插件,供失败处理与 QueueingHint 使用。
fitErr := &framework.FitError{
NumAllNodes: 1,
Pod: assumedPodInfo.Pod,
Diagnosis: framework.Diagnosis{
NodeToStatus: framework.NewDefaultNodeToStatus(),
UnschedulablePlugins: sets.New(status.Plugin()),
},
}
fitErr.Diagnosis.NodeToStatus.Set(
scheduleResult.SuggestedHost, status,
)
return fwk.NewStatus(status.Code()).WithError(fitErr)
}
return status
}
sched.SchedulingQueue.Done(assumedPod.UID)
// 记录正在执行 PreBind 的 Pod,以便抢占路径能够取消后续绑定。
if preFlightStatus.IsSuccess() {
var podInPreBindCancel context.CancelCauseFunc
ctx, podInPreBindCancel = context.WithCancelCause(ctx)
defer podInPreBindCancel(nil)
defer schedFramework.RemovePodInPreBind(assumedPod.UID)
schedFramework.AddPodInPreBind(assumedPod.UID, podInPreBindCancel)
}
// VolumeBinding 等插件在这里完成 API Binding 前的准备。
if status := schedFramework.RunPreBindPlugins(
ctx, state, assumedPod, scheduleResult.SuggestedHost,
); !status.IsSuccess() {
return status
}
// PreBind 期间 Pod 可能被抢占;只有仍可绑定时才继续。
bindingPod := schedFramework.GetPodInPreBind(assumedPod.UID)
if bindingPod != nil && !bindingPod.MarkPrebound() {
return fwk.AsStatus(context.Cause(ctx))
}
if status := sched.bind(
ctx, schedFramework, assumedPod,
scheduleResult.SuggestedHost, state,
); !status.IsSuccess() {
return status
}
schedFramework.RunPostBindPlugins(
ctx, state, assumedPod, scheduleResult.SuggestedHost,
)
return nil
}
RunPreBindPlugins 根据 PreBindPreFlight 保存的结果组织执行:允许并行的插件可以并发运行,需要串行的插件仍按配置顺序执行。
API Binding¶
bind 先检查是否有 Extender Binder 接管当前 Pod;否则调用 Framework 中的 Bind 插件。默认的 DefaultBinder 向 API Server 提交 pods/binding 子资源,API Server 将目标 Node 保存到 Pod 的 spec.nodeName。
func (b DefaultBinder) Bind(
ctx context.Context,
state fwk.CycleState,
pod *v1.Pod,
nodeName string,
) *fwk.Status {
binding := &v1.Binding{
ObjectMeta: metav1.ObjectMeta{
Namespace: pod.Namespace,
Name: pod.Name,
UID: pod.UID,
},
Target: v1.ObjectReference{Kind: "Node", Name: nodeName},
}
// 启用异步 API Cacher 时,通过它提交并等待写入完成。
if b.handle.APICacher() != nil {
onFinish, err := b.handle.APICacher().BindPod(binding)
if err != nil {
return fwk.AsStatus(err)
}
err = b.handle.APICacher().WaitOnFinish(ctx, onFinish)
if err != nil {
return fwk.AsStatus(err)
}
return nil
}
// 直接请求 API Server 的 pods/binding 子资源。
err := b.handle.ClientSet().CoreV1().Pods(binding.Namespace).
Bind(ctx, binding, metav1.CreateOptions{})
if err != nil {
return fwk.AsStatus(err)
}
return nil
}
失败处理¶
Pod 尚未进入 Assumed 状态时,Scheduler 可以直接调用 handleSchedulingFailure()。Assume 完成后,后续失败需要先回滚 Reserve 插件状态,并释放 Cache 中记录的本地资源占用。
相关代码位置如下:
pkg/scheduler/
├── schedule_one.go # 调度与绑定失败入口、回滚及诊断更新
├── framework/runtime/framework.go # 逆序运行 Reserve 插件的 Unreserve 回调
└── backend/
├── cache/cache.go # ForgetPod:撤销 Assumed Pod 的本地占用
└── queue/scheduling_queue.go # AddUnschedulableIfNotPresent:安排重新入队
下图串联三类失败的处理路径,以及 handleSchedulingFailure() 如何读取最新 Pod 并决定是否重新入队。
- 处理未形成 Assumed 状态的失败:Filter、PostFilter、调度算法或
Cache.AssumePod()失败时,Cache 中没有需要撤销的本地占用,可以直接进入 FailureHandler。 - 处理 Assume 后的同步失败:Reserve 或 Permit 失败时,Pod 已经占用目标 Node 的本地资源,需要先回滚再进入 FailureHandler。
- 处理绑定阶段失败:PreBindPreFlight、WaitOnPermit、PreBind 或 Bind 失败时,
handleBindingCycleError()负责完成回滚,再调用 FailureHandler。 - 撤销本地状态:
unreserveAndForget()按逆序运行 Unreserve,再调用Cache.ForgetPod()释放 Assumed 资源。回滚成功后,Scheduler 会向 SchedulingQueue 发送EventAssignedPodDelete,表示一个已计入 Cache 的 Pod 占用已经移除。队列据此重新评估其他待调度 Pod,并通过 QueueingHint 判断是否重新入队;当前失败的 Pod 仍由 FailureHandler 单独处理。 - 保存失败诊断:
handleSchedulingFailure()记录错误类型以及返回Unschedulable或Pending的插件,供 Event 和 QueueingHint 使用。 - 确认失败 Pod 是否仍需重试:FailureHandler 通过 Pod Lister 读取最新对象。Pod 已删除或已经绑定时,不再重新入队;仍是同一个未绑定 Pod 时,才继续处理。UID 已变化说明同名 Pod 已重建,Scheduler 会结束旧对象的处理,也不会更新新 Pod 的状态。
- 将失败 Pod 交回队列:
AddUnschedulableIfNotPresent()根据失败原因、QueueingHint 和退避状态,将当前失败 Pod 放入activeQ、backoffQ或unschedulablePods,等待后续重试,并结束本次 in-flight 跟踪。 - 更新 API 状态:除 UID 已变化的情况外,Scheduler 记录
FailedSchedulingEvent,并把PodScheduledCondition 更新为False。
失败场景¶
- PreFilter、Filter 或 PostFilter 无法得到 Node:
handleSchedulingFailure保存拒绝插件,从 Informer Cache 取得当前 Pod 并完成重新入队,再记录FailedSchedulingEvent 和 Pod Condition。 - Reserve 或 Permit 同步失败:
assumeAndReserve或prepareForBindingCycle先运行 Unreserve 和 ForgetPod,schedulingCycle返回后再进入handleSchedulingFailure。 - WaitOnPermit、PreBind 或 Bind 在异步绑定阶段失败:
handleBindingCycleError先按逆序运行 Unreserve,再调用 ForgetPod 释放本地假定资源,随后进入统一FailureHandler。
unreserveAndForget¶
unreserveAndForget 先按 Reserve 插件的逆序运行 Unreserve,再调用 Cache.ForgetPod 移除 Assumed Pod。逆序回滚可以让后执行的 Reserve 插件先释放依赖前一个插件建立的临时状态。
func (sched *Scheduler) unreserveAndForget(...) error {
// 撤销 Reserve 插件保存的卷预绑定等临时状态。
schedFramework.RunReservePluginsUnreserve(
ctx, state, assumedPodInfo.Pod, nodeName,
)
...
// 释放 Scheduler Cache 中的本地资源占用。
return sched.Cache.ForgetPod(
klog.FromContext(ctx), assumedPodInfo.Pod,
)
}
handleBindingCycleError¶
handleBindingCycleError 处理异步绑定阶段的失败。它先调用 unreserveAndForget,再将 Assigned Pod 删除事件交给 SchedulingQueue,使其他受该资源占用影响的 Pod 有机会重新入队,最后调用统一的 FailureHandler。
func (sched *Scheduler) handleBindingCycleError(...) {
logger := klog.FromContext(ctx)
assumedPod := podInfo.Pod
if err := sched.unreserveAndForget(
ctx, state, fwk, podInfo, scheduleResult.SuggestedHost,
); err == nil {
if status.IsRejected() {
// 重新评估其他 Pod;当前失败 Pod 由 FailureHandler 处理。
defer sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(
logger, framework.EventAssignedPodDelete,
assumedPod, nil,
func(pod *v1.Pod) bool { return assumedPod.UID != pod.UID },
)
} else {
sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(
logger, framework.EventAssignedPodDelete,
assumedPod, nil, nil,
)
}
}
sched.FailureHandler(
ctx, fwk, podInfo, status, clearNominatedNode, start,
)
}
handleSchedulingFailure¶
handleSchedulingFailure 负责处理一次失败的调度:记录失败原因,确认 Pod 仍是同一个未绑定对象,把它放回 SchedulingQueue,并更新 FailedScheduling Event 和 PodScheduled Condition。下面的源码保留了普通 Pod 的主要处理路径:
func (sched *Scheduler) handleSchedulingFailure(
ctx context.Context,
podFwk framework.Framework,
podInfo *framework.QueuedPodInfo,
status *fwk.Status,
nominatingInfo *fwk.NominatingInfo,
start time.Time,
) {
calledDone := false
defer func() {
// 未重新入队的分支仍要结束本轮 in-flight 记录。
if !calledDone {
sched.SchedulingQueue.Done(podInfo.Pod.UID)
}
}()
logger := klog.FromContext(ctx)
reason := v1.PodReasonSchedulerError
if status.IsRejected() {
reason = v1.PodReasonUnschedulable
}
pod := podInfo.Pod
err := status.AsError()
errMsg := status.Message()
// 使用本次调度结果替换上次尝试留下的失败插件记录。
podInfo.ClearRejectorPlugins()
if fitError, ok := err.(*framework.FitError); ok {
podInfo.UnschedulablePlugins = fitError.Diagnosis.UnschedulablePlugins
podInfo.PendingPlugins = fitError.Diagnosis.PendingPlugins
}
// 从 Informer Cache 取得当前 Pod,避免把旧对象重新入队。
podLister := podFwk.SharedInformerFactory().Core().V1().Pods().Lister()
cachedPod, e := podLister.Pods(pod.Namespace).Get(pod.Name)
if e != nil {
...
} else if len(cachedPod.Spec.NodeName) != 0 {
// Pod 已经完成绑定,不再放回调度队列。
...
} else {
// 同名 Pod 的 UID 已变化,说明原 Pod 已删除并重新创建。
if cachedPod.UID != podInfo.Pod.UID {
return
}
podInfo.PodInfo, _ = framework.NewPodInfo(cachedPod.DeepCopy())
pod = podInfo.Pod
// 根据失败插件、调度期间的事件和退避状态选择重新入队位置。
if err := sched.SchedulingQueue.AddUnschedulableIfNotPresent(
logger,
podInfo,
sched.SchedulingQueue.SchedulingCycle(),
); err != nil {
...
}
calledDone = true
}
// 先更新 Scheduler 内部的 nominated Pod,再更新 API 中的状态。
if sched.SchedulingQueue != nil {
sched.SchedulingQueue.AddNominatedPod(
logger,
podInfo.PodInfo,
nominatingInfo,
)
}
// 记录失败事件,并更新 PodScheduled Condition 与 nominatedNodeName。
msg := truncateMessage(errMsg)
podFwk.EventRecorder().WithLogger(logger).Eventf(
pod, nil, v1.EventTypeWarning,
"FailedScheduling", "Scheduling", msg,
)
if err := updatePod(ctx, sched.client, podFwk.APICacher(), pod, &v1.PodCondition{
Type: v1.PodScheduled,
ObservedGeneration: podutil.CalculatePodConditionObservedGeneration(
&pod.Status, pod.Generation, v1.PodScheduled,
),
Status: v1.ConditionFalse,
Reason: reason,
Message: errMsg,
}, nominatingInfo); err != nil {
...
}
}
源码中的处理顺序如下:
- 分类失败结果:插件返回
Unschedulable或Pending时,Pod Condition 使用Unschedulable作为 Reason;调度器内部错误使用SchedulerError。FitError的诊断结果写入UnschedulablePlugins和PendingPlugins,SchedulingQueue 以后只需运行相关插件注册的 QueueingHint。 - 取得最新 Pod:Pod Lister 从 Informer Cache 读取同名对象。后续判断和重新入队使用这个对象,而不是调度开始时取出的旧副本。
- 检查对象身份与绑定状态:
spec.nodeName已设置表示 Pod 已经完成绑定;UID 变化表示原 Pod 已删除,同名的新 Pod 会通过自己的事件进入队列。这两种情况都不会重新入队旧对象。 - 交回 SchedulingQueue:
AddUnschedulableIfNotPresent根据失败插件、调度期间发生的事件和剩余退避时间,将 Pod 放入unschedulablePods、backoffQ或activeQ。 - 记录诊断状态:Scheduler 先更新队列中的 nominated Pod 信息,避免下一轮调度早于 API 状态更新。Event Recorder 随后写入
FailedSchedulingEvent;updatePod将PodScheduledCondition 更新为False,并按 PostFilter 的结果更新status.nominatedNodeName。
AddUnschedulableIfNotPresent 会结束该 Pod 在 SchedulingQueue 中的 in-flight 记录。Lister 无法取得 Pod、Pod 已绑定或同名重建时,函数通过延迟执行的 SchedulingQueue.Done 完成清理,不再创建重新入队记录。
自定义调度插件¶
前面介绍了 Scheduling Framework 的扩展点和插件调用流程。接下来编写一个自定义调度插件 GPUModelFilter,根据 Pod 指定的 GPU 型号,筛选出符合要求的 Node。
实现 GPUModelFilter¶
配置字段¶
插件通过 pluginConfig.args 接收两个必填字段,由 Args 保存:
nodeLabelKey:Node 声明 GPU 型号时使用的 Label Key。示例配置为nvidia.com/gpu.product,与 GPU Feature Discovery 生成的标签一致。podAnnotationKey:Pod 声明所需 GPU 型号时使用的 Annotation Key。示例配置为scheduler.example.com/gpu-product。
type Args struct {
// NodeLabelKey 指定 Node 声明 GPU 型号的 label key。
NodeLabelKey string `json:"nodeLabelKey,omitempty"`
// PodAnnotationKey 指定 Pod 声明所需 GPU 型号的 annotation key。
PodAnnotationKey string `json:"podAnnotationKey,omitempty"`
}
创建插件实例¶
Framework 调用 New() 读取并校验配置,确认两个 Key 均已填写且格式合法后,将它们保存到 GPUModelFilter 实例中。后续 Filter() 按这两个 Key 读取 Pod 注解和 Node 标签。
type GPUModelFilter struct {
nodeLabelKey string
podAnnotationKey string
}
func New(_ context.Context, configuration runtime.Object, _ fwk.Handle) (fwk.Plugin, error) {
args := Args{}
if err := frameworkruntime.DecodeInto(configuration, &args); err != nil {
return nil, fmt.Errorf("decode %s configuration: %w", Name, err)
}
if args.NodeLabelKey == "" {
return nil, fmt.Errorf("nodeLabelKey is required")
}
if messages := validation.IsQualifiedName(args.NodeLabelKey); len(messages) != 0 {
return nil, fmt.Errorf("nodeLabelKey %q is invalid: %s", args.NodeLabelKey, strings.Join(messages, "; "))
}
if args.PodAnnotationKey == "" {
return nil, fmt.Errorf("podAnnotationKey is required")
}
if messages := validation.IsQualifiedName(args.PodAnnotationKey); len(messages) != 0 {
return nil, fmt.Errorf("podAnnotationKey %q is invalid: %s", args.PodAnnotationKey, strings.Join(messages, "; "))
}
return &GPUModelFilter{
nodeLabelKey: args.NodeLabelKey,
podAnnotationKey: args.PodAnnotationKey,
}, nil
}
插件名称¶
Name() 返回固定名称 GPUModelFilter,与注册插件及 Profile 配置中使用的名称保持一致。
const Name = "GPUModelFilter"
func (p *GPUModelFilter) Name() string {
return Name
}
筛选节点¶
Filter() 接收待调度的 pod 和当前候选节点的 nodeInfo,依次完成以下判断:
- 读取 Pod 要求的型号:按
podAnnotationKey读取注解,并去掉首尾空白。注解缺失或值为空时,返回UnschedulableAndUnresolvable,说明缺少哪个注解。 - 读取 Node 的型号:通过
nodeInfo.Node()取得 Node,再按nodeLabelKey读取标签。Node 对象不可用时返回Error;标签缺失或值为空时,返回UnschedulableAndUnresolvable。 - 比较型号:型号不同则返回
UnschedulableAndUnresolvable,并列出节点型号与 Pod 要求的型号;相同则返回nil,表示当前节点通过该插件的检查。
func (p *GPUModelFilter) Filter(_ context.Context, _ fwk.CycleState, pod *v1.Pod, nodeInfo fwk.NodeInfo) *fwk.Status {
requiredProduct := strings.TrimSpace(pod.Annotations[p.podAnnotationKey])
if requiredProduct == "" {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf("pod annotation %q is required", p.podAnnotationKey),
)
}
node := nodeInfo.Node()
if node == nil {
return fwk.NewStatus(fwk.Error, "node information is unavailable")
}
actualProduct, exists := node.Labels[p.nodeLabelKey]
if !exists || strings.TrimSpace(actualProduct) == "" {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf("node %q does not define label %q", node.Name, p.nodeLabelKey),
)
}
if actualProduct != requiredProduct {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf(
"node %q GPU product %q does not match required product %q",
node.Name,
actualProduct,
requiredProduct,
),
)
}
return nil
}
GPUModelFilter 完整源码
package gpumodelfilter
import (
"context"
"fmt"
"strings"
v1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/util/validation"
fwk "k8s.io/kube-scheduler/framework"
frameworkruntime "k8s.io/kubernetes/pkg/scheduler/framework/runtime"
)
const Name = "GPUModelFilter"
type Args struct {
// NodeLabelKey 指定 Node 声明 GPU 型号的 label key。
NodeLabelKey string `json:"nodeLabelKey,omitempty"`
// PodAnnotationKey 指定 Pod 声明所需 GPU 型号的 annotation key。
PodAnnotationKey string `json:"podAnnotationKey,omitempty"`
}
// GPUModelFilter 根据 Pod 注解选择具有相同 GPU 产品标签的节点。
type GPUModelFilter struct {
nodeLabelKey string
podAnnotationKey string
}
var _ fwk.FilterPlugin = &GPUModelFilter{}
// New 解码并校验 pluginConfig.args,再创建插件实例。
func New(_ context.Context, configuration runtime.Object, _ fwk.Handle) (fwk.Plugin, error) {
args := Args{}
if err := frameworkruntime.DecodeInto(configuration, &args); err != nil {
return nil, fmt.Errorf("decode %s configuration: %w", Name, err)
}
if args.NodeLabelKey == "" {
return nil, fmt.Errorf("nodeLabelKey is required")
}
if messages := validation.IsQualifiedName(args.NodeLabelKey); len(messages) != 0 {
return nil, fmt.Errorf("nodeLabelKey %q is invalid: %s", args.NodeLabelKey, strings.Join(messages, "; "))
}
if args.PodAnnotationKey == "" {
return nil, fmt.Errorf("podAnnotationKey is required")
}
if messages := validation.IsQualifiedName(args.PodAnnotationKey); len(messages) != 0 {
return nil, fmt.Errorf("podAnnotationKey %q is invalid: %s", args.PodAnnotationKey, strings.Join(messages, "; "))
}
return &GPUModelFilter{
nodeLabelKey: args.NodeLabelKey,
podAnnotationKey: args.PodAnnotationKey,
}, nil
}
func (p *GPUModelFilter) Name() string {
return Name
}
// Filter 在每个候选节点上运行。该示例只匹配产品型号,不核算 GPU 数量。
func (p *GPUModelFilter) Filter(_ context.Context, _ fwk.CycleState, pod *v1.Pod, nodeInfo fwk.NodeInfo) *fwk.Status {
requiredProduct := strings.TrimSpace(pod.Annotations[p.podAnnotationKey])
if requiredProduct == "" {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf("pod annotation %q is required", p.podAnnotationKey),
)
}
node := nodeInfo.Node()
if node == nil {
return fwk.NewStatus(fwk.Error, "node information is unavailable")
}
actualProduct, exists := node.Labels[p.nodeLabelKey]
if !exists || strings.TrimSpace(actualProduct) == "" {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf("node %q does not define label %q", node.Name, p.nodeLabelKey),
)
}
if actualProduct != requiredProduct {
return fwk.NewStatus(
fwk.UnschedulableAndUnresolvable,
fmt.Sprintf(
"node %q GPU product %q does not match required product %q",
node.Name,
actualProduct,
requiredProduct,
),
)
}
return nil
}
注册插件¶
主程序调用 app.NewSchedulerCommand() 创建调度器启动命令,并通过 app.WithPlugin() 注册插件名称 GPUModelFilter 和工厂函数 New()。
package main
import (
"os"
"k8s.io/component-base/cli"
_ "k8s.io/component-base/logs/json/register"
_ "k8s.io/component-base/metrics/prometheus/clientgo"
_ "k8s.io/component-base/metrics/prometheus/version"
"k8s.io/kubernetes/cmd/kube-scheduler/app"
"example.com/kubernetes-ai-infra/gpu-model-scheduler/pkg/gpumodelfilter"
)
func main() {
// WithPlugin 将工厂注册到调度器 Registry;配置文件再按名称启用插件。
command := app.NewSchedulerCommand(
app.WithPlugin(gpumodelfilter.Name, gpumodelfilter.New),
)
os.Exit(cli.Run(command))
}
配置调度器¶
下面是自定义调度器的配置示例,用于启用 GPUModelFilter 并设置插件参数。
profiles:定义调度器使用的 Profile,每个 Profile 包含调度器名称、启用的插件和插件参数。这里配置一个 Profile。schedulerName:设为gpu-model-scheduler。需要使用这组调度配置的 Pod,其spec.schedulerName必须与这个值一致。plugins.filter.enabled:在 Filter 扩展点启用GPUModelFilter。这里的name必须与前面通过app.WithPlugin()注册的插件名称一致。pluginConfig:配置插件及其参数。它是一个列表,每个配置项包含两个同级字段:name:插件名称,这里为GPUModelFilter,与注册时使用的名称一致。args:传给插件工厂New()的参数,对应前面的Args结构。其中,nodeLabelKey和podAnnotationKey分别指定读取 GPU 型号的 Node 标签和 Pod 注解。
apiVersion: v1
kind: ConfigMap
metadata:
name: gpu-model-scheduler-config
namespace: kube-system
data:
scheduler-config.yaml: |
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
leaderElection:
leaderElect: false
profiles:
- schedulerName: gpu-model-scheduler
plugins:
filter:
enabled:
- name: GPUModelFilter
pluginConfig:
- name: GPUModelFilter
args:
nodeLabelKey: nvidia.com/gpu.product
podAnnotationKey: scheduler.example.com/gpu-product
上面的配置只在 enabled 中添加了 GPUModelFilter,没有显式禁用默认 Filter 插件,因此 NodeResourcesFit、NodeAffinity、TaintToleration 等检查仍然生效。
如果希望 Filter 阶段只运行 GPUModelFilter,可以在当前 Profile 的 plugins.filter.disabled 中使用 name: "*",禁用该扩展点的全部默认插件:
如果只需禁用某个默认 Filter 插件,将 "*" 替换为具体名称,例如 TaintToleration。这里的禁用范围仅限 Filter 扩展点;使用 "*" 会一并移除该阶段的资源、亲和性和污点等默认检查。
部署调度器¶
下面通过 Deployment 运行自定义调度器,并挂载前面定义的调度配置。
-
serviceAccountName:指定调度器访问 API Server 时使用的身份。这里为gpu-model-scheduler,由rbac.yaml创建 ServiceAccount,并绑定以下角色:system:kube-scheduler:ClusterRole,允许读取和监听 Pod、Node,提交 Pod 绑定请求、更新 Pod 状态及记录 Event 等。system:volume-scheduler:ClusterRole,允许读取、更新 PV 和 PVC,并读取 StorageClass,供卷调度与绑定使用。extension-apiserver-authentication-reader:kube-system中的 Role,允许读取该命名空间下的extension-apiserver-authenticationConfigMap,获取认证配置。
-
image:运行包含GPUModelFilter的自定义调度器镜像,名称和标签应与部署脚本构建、加载的镜像一致。 volumes:定义名为scheduler-config的配置卷。configMap.name指向前面的 ConfigMapgpu-model-scheduler-config。volumeMounts:通过name: scheduler-config引用该配置卷,并挂载到容器的/etc/kubernetes目录。ConfigMap 中的scheduler-config.yaml会成为该目录下的同名文件。args:传给调度器的启动参数。--config=/etc/kubernetes/scheduler-config.yaml指向挂载后的配置文件,调度器据此启用插件并读取参数。
apiVersion: apps/v1
kind: Deployment
metadata:
name: gpu-model-scheduler
namespace: kube-system
labels:
app: gpu-model-scheduler
spec:
replicas: 1
selector:
matchLabels:
app: gpu-model-scheduler
template:
metadata:
labels:
app: gpu-model-scheduler
spec:
serviceAccountName: gpu-model-scheduler
containers:
- name: scheduler
image: k8s-ai-infra/gpu-model-scheduler:v1.36.4
imagePullPolicy: Never
args:
- --config=/etc/kubernetes/scheduler-config.yaml
- --secure-port=0
- --v=4
resources:
requests:
cpu: 100m
memory: 128Mi
limits:
memory: 512Mi
securityContext:
allowPrivilegeEscalation: false
readOnlyRootFilesystem: true
runAsNonRoot: true
capabilities:
drop:
- ALL
volumeMounts:
- name: scheduler-config
mountPath: /etc/kubernetes
readOnly: true
volumes:
- name: scheduler-config
configMap:
name: gpu-model-scheduler-config
本地需要安装 Docker、Kind 和 kubectl,并确保 Docker 已启动。运行 setup-kind.sh 准备 Kind 集群、构建并加载镜像,然后应用调度配置、RBAC 和 Deployment:
脚本默认使用 scheduler-demo 作为集群名称,集群包含一个 control-plane 和两个 worker。两个 worker 分别添加以下 GPU 型号标签,供 GPUModelFilter 筛选:
scheduler-demo-worker: nvidia.com/gpu.product=NVIDIA-A100-SXM4-80GB
scheduler-demo-worker2: nvidia.com/gpu.product=NVIDIA-H100-80GB-HBM3
脚本会等待 kube-system/gpu-model-scheduler Deployment 就绪后返回。随后,spec.schedulerName 为 gpu-model-scheduler 的 Pod 就会使用这套调度配置。
验证调度结果¶
接下来,通过 test-pods.yaml 创建三个测试 Pod,交给 gpu-model-scheduler 调度,分别验证请求 A100、请求 H100 和缺少 GPU 型号注解时的调度结果:
apiVersion: v1
kind: Pod
metadata:
name: gpu-model-a100
namespace: default
labels:
app: gpu-model-filter-test
annotations:
scheduler.example.com/gpu-product: NVIDIA-A100-SXM4-80GB
spec:
schedulerName: gpu-model-scheduler
restartPolicy: Never
containers:
- name: pause
image: registry.k8s.io/pause:3.10
---
apiVersion: v1
kind: Pod
metadata:
name: gpu-model-h100
namespace: default
labels:
app: gpu-model-filter-test
annotations:
scheduler.example.com/gpu-product: NVIDIA-H100-80GB-HBM3
spec:
schedulerName: gpu-model-scheduler
restartPolicy: Never
containers:
- name: pause
image: registry.k8s.io/pause:3.10
---
apiVersion: v1
kind: Pod
metadata:
name: gpu-model-missing-annotation
namespace: default
labels:
app: gpu-model-filter-test
spec:
schedulerName: gpu-model-scheduler
restartPolicy: Never
containers:
- name: pause
image: registry.k8s.io/pause:3.10
运行验证脚本:
脚本重新创建三个 Pod,并检查以下结果:
gpu-model-a100绑定到带 A100 标签的 worker。gpu-model-h100绑定到带 H100 标签的 worker。gpu-model-missing-annotation保持 Pending,FailedSchedulingEvent 包含插件返回的缺少注解的原因。
脚本执行成功后,输出类似如下:
GPUModelFilter verification
Pod gpu-model-a100: node=scheduler-demo-worker
Pod gpu-model-h100: node=scheduler-demo-worker2
Pod gpu-model-missing-annotation: Pending
FailedScheduling event: FailedScheduling 0/3 nodes are available: 1 node(s) had untolerated taint(s), 2 pod annotation "scheduler.example.com/gpu-product" is required. no new claims to deallocate, preemption: 0/3 nodes are available: 3 Preemption is not helpful for scheduling.
GPU resource request: absent
Result: PASS
Event 中的 pod annotation "scheduler.example.com/gpu-product" is required 是插件返回的拒绝原因。缺少注解的 Pod 保持 Pending 是预期结果,PASS 表示三个测试 Pod 的检查均已通过。
验证完成后,执行清理脚本,删除本例的 scheduler-demo Kind 集群及其中的所有资源:
FAQ¶
PreFilter 和 Filter 有什么区别?
PreFilter 在当前 Pod 的每轮调度中运行一次,用于计算所有候选 Node 都会使用的数据。Filter 对每台候选 Node 运行一次,判断该 Node 是否满足 Pod 的硬约束。
| 阶段 | 执行次数 | 作用 |
|---|---|---|
PreFilter |
每个插件一次 | 计算所有 Node 都会用到的数据 |
Filter |
每台候选 Node 一次 | 判断当前 Node 是否满足硬约束 |
例如集群中有 1,000 台候选 Node,NodeResourcesFit 在 PreFilter 中只计算一次 Pod 的 CPU、内存和 GPU 请求,并将结果写入 CycleState。Filter 随后使用这份结果逐台检查 Node,无需重复解析 1,000 次 Pod 的资源配置。
插件没有可复用的筛选计算时,可以只实现 Filter。
PreScore 和 Score 有什么区别?
PreScore 在当前 Pod 的每轮调度中运行一次,用于汇总所有可行 Node 的评分背景数据。Score 对每台可行 Node 运行一次,计算该 Node 的分数。
| 阶段 | 执行次数 | 作用 |
|---|---|---|
PreScore |
每个插件一次 | 汇总可行 Node 的评分背景数据 |
Score |
每台可行 Node 一次 | 计算当前 Node 的分数 |
例如 PodTopologySpread 在 PreScore 中统计各个拓扑域中已有的 Pod 数量。Score 为某台 Node 打分时,直接查询该 Node 所在拓扑域的计数,无需重新扫描整个集群。
插件没有可复用的评分计算时,可以只实现 Score。
集群有几千个节点,Scheduler 会全部检查一遍吗?如果提前停止,会不会漏掉更合适的节点?
不一定会全部检查;提前停止时,也可能漏掉得分更高的节点。
进入常规 Filter 和 Score 路径时,Scheduler 会先确定本轮需要找到多少个可行节点;找到足够数量后,便取消后续扫描,再对返回的可行节点评分。
本轮需要找到的可行节点数量,按以下步骤确定:
- 读取配置:优先使用当前 Profile 的
percentageOfNodesToScore,未配置时使用全局值。 -
确定比例:配置值非零时,直接使用该比例;为
0时,根据候选节点数量自动计算,自动计算的比例最低为5%。下式中numAllNodes是候选节点数量,除法取整数商: -
换算节点数:候选节点总数不少于 100 时,按比例计算目标可行节点数,计算结果不足 100 则将目标提高到 100;候选节点总数不足 100 时,检查全部候选节点,只把通过 Filter 的节点加入可行节点列表。
例如,参与筛选的节点有 5,000 个时,两种 percentageOfNodesToScore 配置会得到不同的目标数量:
- 配置为
20:直接使用20%,目标是找到5000 × 20% = 1000个可行节点。 - 配置为
0:自动计算得到50 - 5000 / 125 = 10,即10%,目标是找到5000 × 10% = 500个可行节点。
目标数量指通过 Filter 的节点数量,不是固定扫描多少个节点。如果很多节点不满足条件,就需要检查更多节点;始终找不够时,会遍历全部候选节点。
未进入可行节点列表的节点不会参与本轮评分。提高这个比例可以扩大比较范围,但也会增加筛选和评分开销。为避免每轮总是优先检查同一批节点,Scheduler 会从上一轮停止的位置继续扫描,节点列表也会按 Zone 交错排列。这些措施让不同节点都有机会被检查,但不保证每个 Pod 都选到整个集群中得分最高的节点。
Pod 还在调度时收到资源释放事件,失败后还能据此重试吗?
可以,不会因为事件先于调度失败到达而错过。
SchedulingQueue 通过 activeQ 内的 inFlightEvents 记录 Pod 调度期间收到的相关事件。Pod 失败返回后,队列再通过 QueueingHint 判断这些变化是否有助于重试:值得重试的 Pod 进入 activeQ 或 backoffQ,否则留在 unschedulablePods 等待后续事件。重试仍受退避限制,不一定立即开始。
前一个 Pod 还没绑定完成,后一个 Pod 就开始调度,会不会把同一份资源分配两次?
在同一个 Scheduler 的普通 Pod 调度路径中,不会因前一个 Pod 尚未绑定完成而重复分配同一份资源。
Scheduler 会先为前一个 Pod 选好节点并记录本地占用,再异步执行绑定。下一个 Pod 开始调度时,已经能看到这部分占用,无需等待 API Server 确认绑定。
假设某个 Node 还可以容纳 4 CPU 的请求,Pod A 和 Pod B 各请求 3 CPU。A 选中该 Node 后,Cache.AssumePod() 会先把 A 计入目标 NodeInfo。B 的调度周期刷新 Snapshot 后,看到的剩余可分配量只有 1 CPU,因此不会因为 A 尚未绑定完成,就把同一份资源再分配给 B。
绑定成功后,Pod Informer 通过 Cache.AddPod() 用观察到的 Pod 更新临时记录,不会把 A 的资源请求再叠加一次;后续阶段失败需要回滚时,则通过 Cache.ForgetPod() 撤销这部分本地占用。这套协调依赖同一个 Scheduler 的 Cache,不能当作多个独立调度器之间的资源锁。
正在筛选节点时,集群中的资源占用发生了变化,这一轮调度应该使用旧数据还是新数据?
使用本轮开始时刷新的 Snapshot,不会在筛选过程中切换到实时更新的 Cache 数据。
Scheduler 在运行调度算法前更新 Snapshot,将 Cache 中已经观察到的变化复制过来;之后 Informer 对 Cache 的更新,不会直接改写本轮 Snapshot。
这样,Filter 和 Score 可以使用同一份节点视图,而不用一边筛选、一边跟随集群事件修改判断依据。对于本轮节点读取,Scheduler 只在刷新 Snapshot 时持有 Cache 锁,不用覆盖整个节点筛选和评分过程。接口注释明确说明,这份 Snapshot 在本轮 Permit 阶段结束前保持稳定。
例如,当前 Pod 已经根据 Snapshot 判断某个 Node 资源不足,随后另一个 Pod 被删除并释放资源,本轮不会立即重新检查这个 Node。若本次调度失败,队列可以根据删除事件安排重试,下一轮刷新 Snapshot 后再使用更新后的资源状态。这里的“最新”是 Scheduler 已经观察到的状态,不保证与 API Server 在同一时刻完全一致。
相关资料¶
- Kubernetes Scheduling Framework
- KubeSchedulerConfiguration v1
- Scheduling Policies
- Kubernetes v1.36.4 Scheduler source
- 深入理解 Kubernetes Scheduler Framework 调度框架(Part 1)
- 深入理解 Kubernetes Scheduler Framework 调度框架(Part 2)
- 深入理解 Kubernetes Scheduler Framework 调度框架(Part 3)
- 深入理解 Kubernetes Scheduler Framework 调度框架(Part 4)

















