跳转至

Scheduler:Pod 的节点选择

Kubernetes 调度器是一个控制面进程,负责将 Pod 指派到节点上。调度器根据约束和可用资源,为调度队列中的每个 Pod 找出所有可以运行它的节点,再对这些节点排序,将 Pod 绑定到一个合适的节点。同一集群可以运行多个调度器;kube-scheduler 是 Kubernetes 的参考实现。

Scheduling Framework

Scheduling Framework 是 Kubernetes 调度器的可插拔架构。它由一组直接编译进调度器的插件 API 组成,使大多数调度功能可以通过插件实现,同时让调度核心保持轻量、易于维护。更多设计细节见 Scheduling Framework 设计提案。

整体架构

下图展示了 Pod 的调度上下文以及 Scheduling Framework 暴露的接口。一个插件可以实现多个接口,用于处理更复杂或需要保存状态的调度逻辑。部分接口对应可以通过调度器配置设置的扩展点。

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 调度过程

下面分别说明这三个阶段的处理过程。

进入调度队列

  1. 观察 Pod:Scheduler 观察 API Server 中的 Pod 变化。已经设置 spec.nodeName 的 Pod 表示它已被分配到某个 Node,Scheduler 将它记录到本地 Scheduler Cache,用于计算 Node 的资源占用;尚未指定 Node 且由当前 Scheduler 负责的 Pod 进入 SchedulingQueue,等待调度。
  2. 检查入队:PreEnqueue 检查 Pod 能否进入可调度队列。首次通过检查的 Pod 进入 activeQ;检查未通过时进入 unschedulablePods。
  3. 等待重试:调度未能继续的 Pod 进入 unschedulablePods,等待可能改变结果的集群事件。Unschedulable 表示当前条件下无法完成调度;Pending 表示已经选出目标 Node,但插件需要等待外部组件根据调度结果完成操作,Pod 暂时不能进入绑定阶段。相关事件触发重新入队后,Pending Pod 直接进入 activeQ;Unschedulable Pod 在退避尚未结束时进入 podBackoffQ,退避结束后进入 activeQ。

进入调度队列

调度阶段(Scheduling cycle)

  1. 取出 Pod:Scheduler 从 SchedulingQueue 取出 Pod,读取当前的 Node 和资源状态。
  2. 筛选 Node:PreFilter 准备筛选所需的状态,Filter 排除不满足硬约束的 Node。没有可行 Node 时,PostFilter 尝试抢占或其他补救措施;仍无法选出 Node 时,由 FailureHandler 记录失败原因,并将 Pod 交回 SchedulingQueue 等待重试。
  3. 选择 Node:存在多个可行 Node 时,PreScore 准备评分所需的状态,Score 选择得分最高的 Node;只有一个可行 Node 时直接使用该 Node。
  4. 记录选择:Scheduler 在本地 Cache 中记录目标 Node 和资源占用。Reserve 保存可回滚的临时状态,例如卷的预绑定信息;Permit 决定继续绑定、等待其他条件满足,或者拒绝本次选择。Reserve 或 Permit 失败时,Scheduler 先运行 Unreserve 撤销插件状态,再调用 ForgetPod 删除 Cache 中的 Assumed Pod。

调度阶段

绑定阶段(Binding cycle)

  1. 检查 PreBind 插件:Scheduler 逐个运行 PreBindPreFlight,记录需要跳过的插件,以及哪些相邻插件可以在后续 PreBind 中并行执行。预检本身按顺序执行。
  2. 等待 Permit 放行:Scheduler 调用 WaitOnPermit。只有 Permit 插件在调度阶段返回 Wait 时,这一步才会阻塞;否则立即进入 PreBind。等待期间保留目标 Node 的本地资源占用,插件拒绝或等待超时会撤销本次选择。
  3. 准备绑定:PreBind 完成卷绑定等准备工作。
  4. 提交绑定:Bind 负责提交绑定结果。默认 DefaultBinder 请求 Pod 的 binding 子资源,API Server 将目标 Node 写入 spec.nodeName。
  5. 确认结果: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 的使用计数

这两个方法主要用于以下两类调度计算:

这些操作只更新 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 和失败处理流程。

Scheduler 结构

  • 保存本地数据:SchedulingQueue 保存待处理 Pod,Cache 汇总 Node 与 Pod 状态,Profiles 保存各调度策略对应的 Framework,nodeInfoSnapshot 提供本轮调度使用的节点视图。
  • 取出 Pod:NextPod 默认指向 SchedulingQueue.Pop,队列为空时阻塞等待。
  • 选择 Node:SchedulePod 默认指向 schedulePod,使用选定的 Framework 和 NodeInfo Snapshot 完成筛选与评分。
  • 处理失败:FailureHandler 默认指向 handleSchedulingFailure,负责记录诊断信息并把 Pod 交回 SchedulingQueue。

下面的源码保留普通 Pod 调度主路径使用的字段。

kubernetes/pkg/scheduler/scheduler.go
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:主循环调用的调度入口

下图展示了从命令入口到调度主循环的启动顺序。

Scheduler 初始化与启动

命令与对象创建

  • 创建命令: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 初始化: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。
  • 启动调度 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 配置。

kubernetes/cmd/kube-scheduler/app/server.go
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:

kubernetes/cmd/kube-scheduler/app/server.go
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 注册位置:

kubernetes/pkg/scheduler/scheduler.go
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 主循环。

kubernetes/cmd/kube-scheduler/app/server.go
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:

kubernetes/pkg/scheduler/scheduler.go
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。

Scheduling Profile

  • 读取 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 会生成默认参数,供创建插件实例时使用。

kubernetes/pkg/scheduler/apis/config/v1/defaults.go
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 编写自定义插件;后面的自定义调度插件通过完整示例介绍插件接口、注册、配置和验证过程。

examples/chapter-01/05-scheduler/manifests/multi-profile-config.yaml
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:

kubernetes/pkg/scheduler/profile/profile.go
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:

kubernetes/pkg/scheduler/schedule_one.go
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 与扩展点插件接口

Framework 插件管理和调用

插件来源与配置

  • 读取插件工厂: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 本身再声明调度和绑定阶段需要的其他入口:

kubernetes/pkg/scheduler/framework/interface.go
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 调用这些方法:

kubernetes/staging/src/k8s.io/kube-scheduler/framework/interface.go
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 插件:

kubernetes/pkg/scheduler/framework/plugins/defaultpreemption/default_preemption.go
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 的定义如下:

kubernetes/pkg/scheduler/framework/runtime/registry.go
// 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 方法:

kubernetes/pkg/scheduler/framework/runtime/framework.go
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 按以下顺序创建插件并建立扩展点列表:

  1. 选择插件:读取当前 Profile 的 plugins 和 pluginConfig,再从 Registry 找到对应的插件工厂。

  2. 创建插件实例:插件工厂收到配置参数和 Handle。Handle 是 Framework 提供给插件的共享能力接口,可以访问 NodeInfo Snapshot、Informer 和 Kubernetes Client 等能力。创建完成的实例按名称写入 pluginsMap。

  3. 建立扩展点列表:updatePluginList 按 Profile 中的配置顺序,将 pluginsMap 中的实例加入 PreFilter、Filter、Score、Bind 等执行列表。同一个插件实现多个扩展点时,这些列表引用同一个实例。

  4. 展开 MultiPoint:expandMultiPointPlugins 将 MultiPoint 中启用的插件加入它所实现的全部扩展点。某个扩展点存在显式的 enabled 或 disabled 配置时,以显式配置为准。

  5. 保存 Score 权重:getValidScoreWeights 校验已启用 Score 插件的权重,并按插件名称写入 scorePluginWeight。权重没有设置或配置为 0 时使用 1。该权重只作用于 Score 扩展点。

kubernetes/pkg/scheduler/framework/runtime/framework.go
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 的插件名。

kubernetes/pkg/scheduler/framework/runtime/framework.go
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 还会统计插件的执行时间和返回状态:

kubernetes/pkg/scheduler/framework/runtime/framework.go
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:

kubernetes/pkg/scheduler/framework/plugins/nodeunschedulable/node_unschedulable.go
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 之间的主要流转路径。

SchedulingQueue

  • 接收待调度 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 的重试路径。
    • 绑定阶段成功: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。

PriorityQueue

Scheduler 创建过程中会调用 NewSchedulingQueue,得到默认实现 PriorityQueue。PriorityQueue 使用三个字段保存不同状态的 Pod:

  • activeQ:现在可以尝试调度的 Pod。
  • backoffQ:保存已经满足重新入队条件、但仍处于退避期的 Pod。podBackoffQ 保存插件拒绝后等待普通退避结束的 Pod;podErrorBackoffQ 保存发生 Scheduler Error 后等待错误退避结束的 Pod。
  • unschedulablePods:保存尚未满足重新入队条件的 Pod,等待可能改变上次调度结果的集群事件。
kubernetes/pkg/scheduler/backend/queue/scheduling_queue.go
type PriorityQueue struct {
    // ...

    // activeQ 保存当前可以调度的 Pod。
    activeQ  activeQueuer
    // backoffQ 保存仍处于退避时间内的 Pod。
    backoffQ backoffQueuer

    // unschedulablePods 等待可能改变调度结果的集群事件。
    unschedulablePods *unschedulablePods
}

activeQ

activeQueue 使用 Heap 保存可立即调度的 Pod,并记录正在调度的 Pod、调度期间发生的事件和调度周期序号:

kubernetes/pkg/scheduler/backend/queue/active_queue.go
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 顺序:

kubernetes/pkg/scheduler/backend/queue/scheduling_queue.go
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 根据失败记录选择内部队列:

kubernetes/pkg/scheduler/backend/queue/backoff_queue.go
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 不会沿这条路径提前取出:

kubernetes/pkg/scheduler/backend/queue/active_queue.go
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:

kubernetes/pkg/scheduler/backend/queue/unschedulable_pods.go
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 重新具备调度条件:

kubernetes/pkg/scheduler/backend/queue/scheduling_queue.go
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 再根据判断结果和剩余退避时间选择目标队列:

kubernetes/pkg/scheduler/backend/queue/scheduling_queue.go
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 选择目标队列:

kubernetes/pkg/scheduler/backend/queue/scheduling_queue.go
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。

Scheduler Cache

  • 更新 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() 将临时记录更新为已确认状态。
  • 增量刷新 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 调度路径直接使用的字段:

kubernetes/pkg/scheduler/backend/cache/cache.go
// 每个链表项保存一个 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 缓存状态有两个更新来源:

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,它的资源只计算一次。

Pod 缓存状态

  • 进入 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:

kubernetes/pkg/scheduler/backend/cache/cache.go
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 会被两类数据变化更新:

下图展示 Scheduler Cache 如何维护 Node 的名称索引和更新时间链表,以及如何利用这些结构增量更新 Snapshot;nodeTree 则负责生成按 Zone 交错排列的节点顺序。

Node 状态与索引

  • 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() 增量刷新,而不是重新创建完整副本。

kubernetes/pkg/scheduler/backend/cache/snapshot.go
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:

Snapshot 增量更新与读取

  • 进入调度周期: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。

kubernetes/pkg/scheduler/schedule_one.go
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 中运行。

kubernetes/pkg/scheduler/schedule_one.go
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 插件。

kubernetes/pkg/scheduler/schedule_one.go
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 会在后续周期重新尝试。

kubernetes/pkg/scheduler/schedule_one.go
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 发起

Filter 的节点并行

Score:Node 并行由 RunScorePlugins 发起

Score 的节点并行

Filter 并行源码

findNodesThatPassFilters 先计算本轮需要找到多少台可行 Node,再把“检查一台 Node”封装成 checkNode。Parallelizer 并行执行这些工作项;每台检查通过的 Node 通过原子计数取得结果数组中的位置,找到足够数量后取消剩余工作:

kubernetes/pkg/scheduler/schedule_one.go
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 中的配置顺序运行。

kubernetes/pkg/scheduler/framework/runtime/framework.go
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 中的临时占用。

kubernetes/pkg/scheduler/schedule_one.go
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:

kubernetes/pkg/scheduler/schedule_one.go
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。

kubernetes/pkg/scheduler/schedule_one.go
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。

kubernetes/pkg/scheduler/framework/plugins/defaultbinder/default_binder.go
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 记录 FailedScheduling Event,并把 PodScheduled Condition 更新为 False。

失败场景

unreserveAndForget

unreserveAndForget 先按 Reserve 插件的逆序运行 Unreserve,再调用 Cache.ForgetPod 移除 Assumed Pod。逆序回滚可以让后执行的 Reserve 插件先释放依赖前一个插件建立的临时状态。

kubernetes/pkg/scheduler/schedule_one.go
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。

kubernetes/pkg/scheduler/schedule_one.go
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 的主要处理路径:

kubernetes/pkg/scheduler/schedule_one.go
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 {
        ...
    }
}

源码中的处理顺序如下:

  1. 分类失败结果:插件返回 Unschedulable 或 Pending 时,Pod Condition 使用 Unschedulable 作为 Reason;调度器内部错误使用 SchedulerError。FitError 的诊断结果写入 UnschedulablePlugins 和 PendingPlugins,SchedulingQueue 以后只需运行相关插件注册的 QueueingHint。
  2. 取得最新 Pod:Pod Lister 从 Informer Cache 读取同名对象。后续判断和重新入队使用这个对象,而不是调度开始时取出的旧副本。
  3. 检查对象身份与绑定状态:spec.nodeName 已设置表示 Pod 已经完成绑定;UID 变化表示原 Pod 已删除,同名的新 Pod 会通过自己的事件进入队列。这两种情况都不会重新入队旧对象。
  4. 交回 SchedulingQueue:AddUnschedulableIfNotPresent 根据失败插件、调度期间发生的事件和剩余退避时间,将 Pod 放入 unschedulablePods、backoffQ 或 activeQ。
  5. 记录诊断状态:Scheduler 先更新队列中的 nominated Pod 信息,避免下一轮调度早于 API 状态更新。Event Recorder 随后写入 FailedScheduling Event;updatePod 将 PodScheduled Condition 更新为 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。
examples/chapter-01/05-scheduler/pkg/gpumodelfilter/plugin.go
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 标签。

examples/chapter-01/05-scheduler/pkg/gpumodelfilter/plugin.go
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 配置中使用的名称保持一致。

examples/chapter-01/05-scheduler/pkg/gpumodelfilter/plugin.go
const Name = "GPUModelFilter"

func (p *GPUModelFilter) Name() string {
    return Name
}

筛选节点

Filter() 接收待调度的 pod 和当前候选节点的 nodeInfo,依次完成以下判断:

  1. 读取 Pod 要求的型号:按 podAnnotationKey 读取注解,并去掉首尾空白。注解缺失或值为空时,返回 UnschedulableAndUnresolvable,说明缺少哪个注解。
  2. 读取 Node 的型号:通过 nodeInfo.Node() 取得 Node,再按 nodeLabelKey 读取标签。Node 对象不可用时返回 Error;标签缺失或值为空时,返回 UnschedulableAndUnresolvable。
  3. 比较型号:型号不同则返回 UnschedulableAndUnresolvable,并列出节点型号与 Pod 要求的型号;相同则返回 nil,表示当前节点通过该插件的检查。
examples/chapter-01/05-scheduler/pkg/gpumodelfilter/plugin.go
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 完整源码
examples/chapter-01/05-scheduler/pkg/gpumodelfilter/plugin.go
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()。

examples/chapter-01/05-scheduler/cmd/gpu-model-scheduler/main.go
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 注解。
examples/chapter-01/05-scheduler/manifests/scheduler-config.yaml
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: "*",禁用该扩展点的全部默认插件:

plugins:
  filter:
    disabled:
      - name: "*"
    enabled:
      - name: GPUModelFilter

如果只需禁用某个默认 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-authentication ConfigMap,获取认证配置。
  • image:运行包含 GPUModelFilter 的自定义调度器镜像,名称和标签应与部署脚本构建、加载的镜像一致。

  • volumes:定义名为 scheduler-config 的配置卷。configMap.name 指向前面的 ConfigMap gpu-model-scheduler-config。
  • volumeMounts:通过 name: scheduler-config 引用该配置卷,并挂载到容器的 /etc/kubernetes 目录。ConfigMap 中的 scheduler-config.yaml 会成为该目录下的同名文件。
  • args:传给调度器的启动参数。--config=/etc/kubernetes/scheduler-config.yaml 指向挂载后的配置文件,调度器据此启用插件并读取参数。
examples/chapter-01/05-scheduler/manifests/deployment.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:

cd examples/chapter-01/05-scheduler
./scripts/setup-kind.sh

脚本默认使用 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 型号注解时的调度结果:

examples/chapter-01/05-scheduler/manifests/test-pods.yaml
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

运行验证脚本:

./scripts/verify.sh

脚本重新创建三个 Pod,并检查以下结果:

  • gpu-model-a100 绑定到带 A100 标签的 worker。
  • gpu-model-h100 绑定到带 H100 标签的 worker。
  • gpu-model-missing-annotation 保持 Pending,FailedScheduling Event 包含插件返回的缺少注解的原因。

脚本执行成功后,输出类似如下:

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 集群及其中的所有资源:

./scripts/cleanup.sh

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 会先确定本轮需要找到多少个可行节点;找到足够数量后,便取消后续扫描,再对返回的可行节点评分。

本轮需要找到的可行节点数量,按以下步骤确定:

  1. 读取配置:优先使用当前 Profile 的 percentageOfNodesToScore,未配置时使用全局值。
  2. 确定比例:配置值非零时,直接使用该比例;为 0 时,根据候选节点数量自动计算,自动计算的比例最低为 5%。下式中 numAllNodes 是候选节点数量,除法取整数商:

    percentage = max(5, 50 - numAllNodes / 125)
    
  3. 换算节点数:候选节点总数不少于 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 在同一时刻完全一致。

相关资料