乐于分享
好东西不私藏

kube-scheduler源码阅读

kube-scheduler源码阅读

前言

kube-scheduler 是 Kubernetes 控制平面的核心组件之一,负责为待调度 Pod 选择合适的 Node。集群中新建或重新调度的 Pod 在创建时通常没有指定运行节点;调度器监听这些未绑定 Pod,结合节点资源、亲和性、污点容忍等约束做出决策,并将结果写回 Pod 的spec.nodeName,后续由 kubelet 在目标节点上拉起容器。

本文基于Kubernetes1.36.1,按数据流梳理 kube-scheduler 的核心路径:从 Informer 创建与事件回调,到调度队列入队/出队,再到一次完整的 Pod 调度。

一、从 Informer 到 调度循环

Informer → Pod变更回调 → 调度队列 → ScheduleOne拿到Pod → 调度循环

1.1 启动阶段 - 构建Informer

调用链路(server.go → options.go → scheduler.go)。

scheduler.goNewInformerFactory 构造 SharedInformerFactory,注册 Pod Informer。

funcNewInformerFactory(cs clientset.Interface, resyncPeriod time.Duration)informers.SharedInformerFactory { informerFactory := informers.NewSharedInformerFactory(cs, resyncPeriod) informerFactory.InformerFor(&v1.Pod{}, newPodInformer)return informerFactory}funcnewPodInformer(cs clientset.Interface, resyncPeriod time.Duration)cache.SharedIndexInformer { informer := coreinformers.NewFilteredPodInformer(cs, metav1.NamespaceAll, resyncPeriod, cache.Indexers{}, tweakListOptions)return informer}

1.2 启动阶段 - 注册事件监听

eventhandlers.goaddAllEventHandlers 注册pod、node等事件监听到informer。

// pod变更informerFactory.Core().V1().Pods().Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{ AddFunc:    sched.addPod, UpdateFunc: sched.updatePod, DeleteFunc: sched.deletePod,});// node变更informerFactory.Core().V1().Nodes().Informer().AddEventHandler( cache.ResourceEventHandlerFuncs{  AddFunc:    sched.addNodeToCache,  UpdateFunc: sched.updateNodeInCache,  DeleteFunc: sched.deleteNodeFromCache, },);// pvc变更informerFactory.Core().V1().PersistentVolumeClaims().Informer().AddEventHandler(    buildEvtResHandler(at, fwk.PersistentVolumeClaim),);// 其他变更...

1.3 启动阶段 - 开启调度循环

scheduler.go:开启调度循环,单协程执行ScheduleOne。

func(sched *Scheduler)Run(ctx context.Context) {// 开启调度队列后台任务 sched.SchedulingQueue.Run(logger)// 单协程ScheduleOnego wait.UntilWithContext(ctx, sched.ScheduleOne, 0)// 阻塞等待进程退出 <-ctx.Done()}

1.4 运行阶段 - pod变更入调度队列

eventhandlers.goaddPod Pod新增 addPodToSchedulingQueue 入调度队列。

func(sched *Scheduler)addPod(obj interface{}) {// ...if assignedPod(pod) {  sched.addAssignedPodToCache(pod) // 已绑定:进 Cache } elseif responsibleForPod(pod, sched.Profiles) { // 当前scheduler负责这个pod  sched.addPodToSchedulingQueue(pod) // 未绑定:进 SchedulingQueue }}// pod.spec.schedulerName(默认default-scheduler) // 在 当前scheduler中有配置(默认包含default-scheduler)funcresponsibleForPod(pod *v1.Pod, profiles profile.Map)bool {return profiles.HandlesSchedulerName(pod.Spec.SchedulerName)}func(sched *Scheduler)addPodToSchedulingQueue(pod *v1.Pod) { sched.SchedulingQueue.Add(..., pod)}

eventhandlers.goupdatePod Pod更新,更新缓存。

func(sched *Scheduler)updatePod(oldObj, newObj interface{}) { logger := sched.logger oldPod, ok := oldObj.(*v1.Pod) newPod, ok := newObj.(*v1.Pod)if assignedPod(oldPod) {// pod早就完成调度,更新内存  sched.updateAssignedPodInCache(oldPod, newPod) } elseif assignedPod(newPod) {// pod刚分配到node,更新缓存  sched.addAssignedPodToCache(newPod)if responsibleForPod(oldPod, sched.Profiles) {// 当前scheduler负责这个pod的调度,从调度队列清除   sched.deletePodFromSchedulingQueue(oldPod, true)  } } elseif responsibleForPod(oldPod, sched.Profiles) {// 当前scheduler负责这个pod的调度,更新队列里的对象  sched.updatePodInSchedulingQueue(oldPod, newPod) }}// pod.spec.nodeName 非空funcassignedPod(pod *v1.Pod)bool {returnlen(pod.Spec.NodeName) != 0}

1.5 运行阶段 - 调度循环

schedule_one.go:ScheduleOne调度循环。

func(sched *Scheduler)ScheduleOne(ctx context.Context) {// 阻塞,直到队列里有 Pod podInfo, err := sched.NextPod(logger)// 调度Pod sched.scheduleOnePod(ctx, podInfo)}

NextPod 即 从 SchedulingQueue 中获取Pod。

funcNew()(*Scheduler, error) { podQueue := internalqueue.NewSchedulingQueue( ) sched := &Scheduler{  SchedulingQueue: podQueue, } sched.NextPod = podQueue.Popreturn sched, nil}

二、调度队列

Add 入队 → PreEnqueue 闸门 → activeQ / backoffQ / unschedulablePods → Pop 出队; 事件用 QueueingHint 决定是否重试

主调度队列:scheduling_queue.goPriorityQueue)。 配套三个队列:active_queue.gobackoff_queue.gounschedulable_pods.go

scheduler.go:创建SchedulingQueue(PriorityQueue)调度队列,注入 QueueSort、PreEnqueue、EnqueueExtensions扩展(QueueHint)插件。

funcNew()(*Scheduler, error) { podQueue := internalqueue.NewSchedulingQueue(// QueueSort插件:activeQ堆比较函数  profiles[options.profiles[0].SchedulerName].QueueSortFunc(),  informerFactory,// PreEnqueue插件  internalqueue.WithPreEnqueuePluginMap(preEnqueuePluginMap),// EnqueueExtensions扩展,对于一个事件返回QueueingHint,判断是否能重新入队  internalqueue.WithQueueingHintMapPerProfile(queueingHintsPerProfile),// ... ) sched.NextPod = podQueue.Pop}

2.1 三个 Queue

scheduling_queue.go:调度队列实现为PriorityQueue:

  1. activeQ:待调度的pod,按 QueueSort 排序

  2. backoffQ:调度失败后退避中,到期进 activeQ

  3. unschedulablePods:不可调度 / PreEnqueue 挡住,等事件或超时再唤醒

type PriorityQueue struct {// pod抢占节点信息 *nominator// unschedulablePods中的pod停留最大5分钟,要进入activeQ/backoffQ podMaxInUnschedulablePodsDuration time.Duration// 待调度的pod,按 QueueSort 排序 activeQ  activeQueuer// 调度失败后退避中,到期进 activeQ backoffQ backoffQueuer// 不可调度 / PreEnqueue 挡住,等事件或超时再唤醒 unschedulablePods *unschedulablePods// profile-pluginName-plugin preEnqueuePluginMap map[string]map[string]fwk.PreEnqueuePlugin// profile-focus的事件-EnqueueExtensions扩展 queueingHintMap QueueingHintMapPerProfile// plugin-focus的事件 pluginToEventsMap map[string][]fwk.ClusterEvent}

核心流转:

  1. 新pod:PreEnqueue通过,进入activeQ;否则进入unschedulablePods;

  2. ScheduleOne调度循环:从activeQ弹出pod,如果失败,根据情况进入activeQ/backoffQ/unschedulablePods;

  3. 集群事件:EnqueueExtensions订阅事件,根据QueueingHint返回决定是否入队;

  4. 退避结束 / unschedulable 超时:从backoffQ/unschedulablePods入队activeQ;

// 入队(informer pod新增)func(p *PriorityQueue)Add(ctx context.Context, pod *v1.Pod) { pInfo := p.newQueuedPodInfo(ctx, pod)if added := p.moveToActiveQ(logger, pInfo, ...); added {  p.activeQ.broadcast() // 唤醒可能阻塞在 Pop 的 ScheduleOne }}// 出队(调度循环 ScheduleOne 的 NextPod)func(p *PriorityQueue)Pop(logger klog.Logger)(*framework.QueuedPodInfo, error) {return p.activeQ.pop(logger) // activeQ 空则阻塞}// 后台:退避到期、unschedulable 滞留过久强制再试func(p *PriorityQueue)Run(logger klog.Logger) {go p.backoffQ.waitUntilAlignedWithOrderingWindow(func() {  p.flushBackoffQCompleted(logger) // backoffQ → activeQ }, p.stop)go wait.Until(func() {  p.flushUnschedulablePodsLeftover(logger) // 默认超过 5min 强制 flush }, 30*time.Second, p.stop)}

2.2 PreEnqueue

scheduling_queue.go:PreEnqueue是进入activeQ前的闸门,全部 Success 才进 activeQ,否则进 unschedulablePods。

func(p *PriorityQueue)moveToActiveQ(...)bool { p.runPreEnqueuePlugins(context.Background(), pInfo)if pInfo.Gated() {// PreEnqueue 失败  p.unschedulablePods.addOrUpdate(...)returnfalse }// 放入activeQ,可以进入调度循环 unlockedActiveQ.add(logger, pInfo, event)returntrue}

scheduling_gates.go:如默认开启的SchedulingGates插件,如果pod.spec.schedulingGates非空,则拦截。

func(pl *SchedulingGates)PreEnqueue(ctx context.Context, p *v1.Pod) *fwk.Status {iflen(p.Spec.SchedulingGates) == 0 {returnnil// Success }return fwk.NewStatus(fwk.UnschedulableAndUnresolvable,  fmt.Sprintf("waiting for scheduling gates: %v", gates))}

2.3 QueueingHint

对于unschedulable的pod,可以通过集群事件重新投递到activeQ/backoffQ。

不同调度插件通过实现EnqueueExtensions.EventsToRegister登记「关心哪些事件 + Hint 函数」。

type EnqueueExtensions interface { Plugin EventsToRegister(context.Context) ([]ClusterEventWithHint, error)}type ClusterEventWithHint struct {// 关心事件,如node新增 Event ClusterEvent// 判断是否回队Hint函数 QueueingHintFn QueueingHintFn}// unschedulable的pod,老对象,新对象(如老node、新node)type QueueingHintFn func(logger klog.Logger, pod *v1.Pod, oldObj, newObj interface{})(QueueingHint, error)// 是否重新回activeQ/backoffQtype QueueingHint intconst ( QueueSkip QueueingHint = iota Queue)

node_affinity.go:比如节点亲和性,关心节点新增等事件,isSchedulableAfterNodeChange判断是否允许pod重新回队。

func(pl *NodeAffinity)EventsToRegister(_ context.Context)([]fwk.ClusterEventWithHint, error) { nodeActionType := fwk.Add | fwk.UpdateNodeLabel | fwk.UpdateNodeTaintif pl.enableSchedulingQueueHint {  nodeActionType = fwk.Add | fwk.UpdateNodeLabel }return []fwk.ClusterEventWithHint{// 关心Node变更  {Event: fwk.ClusterEvent{Resource: fwk.Node, ActionType: nodeActionType}, // QueueingHintFn判断是否需要重新入队  QueueingHintFn: pl.isSchedulableAfterNodeChange}, }, nil}func(pl *NodeAffinity)isSchedulableAfterNodeChange(logger klog.Logger, pod *v1.Pod, oldObj, newObj interface{})(fwk.QueueingHint, error) { originalNode, modifiedNode, err := util.As[*v1.Node](oldObj, newObj)if err != nil {return fwk.Queue, err }if pl.addedNodeSelector != nil && !pl.addedNodeSelector.Match(modifiedNode) {return fwk.QueueSkip, nil }// ...其他判断}

eventhandlers.go:启动时informer注册node变更回调。

// pod变更...// node变更informerFactory.Core().V1().Nodes().Informer().AddEventHandler( cache.ResourceEventHandlerFuncs{  AddFunc:    sched.addNodeToCache,  UpdateFunc: sched.updateNodeInCache,  DeleteFunc: sched.deleteNodeFromCache, },);

eventhandlers.go:把node变更发送给调度队列。

func(sched *Scheduler)addNodeToCache(obj interface{}) { evt := fwk.ClusterEvent{Resource: fwk.Node, ActionType: fwk.Add} node, ok := obj.(*v1.Node) nodeInfo := sched.Cache.AddNode(logger, node) sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(logger, evt, nil, node, preCheckForNode(logger, nodeInfo))}

scheduling_queue.go:对每个unschedulable pod执行相关QueueingHintFn,某个QueueingHintFn返回Queue,则允许重新入队,否则停留在unschedulable。 调度循环中,调度失败时会把否决插件名记在 podInfo.UnschedulablePlugins / PendingPlugins,被Pending拒绝的进入activeQ,被Unschedulable拒绝的进入backoffQ。

func(p *PriorityQueue)moveAllToActiveOrBackoffQueue(logger klog.Logger, event fwk.ClusterEvent, oldObj, newObj interface{}, preCheck PreEnqueueCheck) {if !p.isEventOfInterest(logger, event) {// 无插件关心事件,不处理return }// 所有unschedulablePods unschedulablePods := make([]*framework.QueuedPodInfo, 0len(p.unschedulablePods.podInfoMap))for _, pInfo := range p.unschedulablePods.podInfoMap {if preCheck == nil || preCheck(pInfo.Pod) {   unschedulablePods = append(unschedulablePods, pInfo)  } } p.movePodsToActiveOrBackoffQueue(logger, unschedulablePods, event, oldObj, newObj)}func(p *PriorityQueue)movePodsToActiveOrBackoffQueue(logger klog.Logger, podInfoList []*framework.QueuedPodInfo, event fwk.ClusterEvent, oldObj, newObj interface{}) {if !p.isEventOfInterest(logger, event) {return } activated := false// 循环所有unschedulablePodsfor _, pInfo := range podInfoList {// 执行所有插件的QueueingHintFn// 只要有插件返回Queue,则需要入队;所有插件Skip,则留在unschedulablePods  schedulingHint := p.isPodWorthRequeuing(logger, pInfo, event, oldObj, newObj)if schedulingHint == queueSkip {continue  }  p.unschedulablePods.delete(pInfo.Pod, pInfo.Gated())// 根据情况进入activeQ/backoffQ  queue := p.requeuePodWithQueueingStrategy(logger, pInfo, schedulingHint, event.Label()) }}

2.4 QueueSort

activeQ入队,需要通过QueueSortPlugin排序,决定后续Pop给调度循环的优先级。

type activeQueue struct {// 大顶堆 queue *heap.Heap[*framework.QueuedPodInfo]}

同一时刻只能启用一个QueueSortPlugin,lessFunc返回true,代表pod1>pod2。

type QueueSortPlugin struct { lessFunc func(info1, info2 fwk.QueuedPodInfo)bool}

priority_sort.go:默认实现,pod优先级(pod.Spec.Priority)高先出,同优先级按进入调度队列的时间戳FIFO。

func(pl *PrioritySort)Less(pInfo1, pInfo2 fwk.QueuedPodInfo)bool { p1 := corev1helpers.PodPriority(pInfo1.GetPodInfo().GetPod()) p2 := corev1helpers.PodPriority(pInfo2.GetPodInfo().GetPod())return (p1 > p2) || (p1 == p2 && pInfo1.GetTimestamp().Before(pInfo2.GetTimestamp()))}

三、Pod 调度

3.1 调度周期与绑定周期

一次Pod调度尝试 = 调度周期(Scheduling Cycle) + 绑定周期(Binding Cycle)

schedule_one.go:ScheduleOne调度循环。

func(sched *Scheduler)ScheduleOne(ctx context.Context) {// 阻塞,直到activeQ里有Pod podInfo, err := sched.NextPod(logger)// 调度Pod sched.scheduleOnePod(ctx, podInfo)}func(sched *Scheduler)scheduleOnePod(ctx context.Context, podInfo *framework.QueuedPodInfo) {// 调度周期 scheduleResult, assumedPodInfo, status := sched.schedulingCycle(...)if !status.IsSuccess() {// 回backoffQ/unschedulablePods  sched.FailureHandler(...)return }// 绑定周期go sched.runBindingCycle(ctx, state, fwk, scheduleResult, assumedPodInfo, start, podsToActivate)}
调度周期
绑定周期
目的
选节点并在内存中预占
将决策落到 API Server
执行
单协程串行异步并发
go runBindingCycle
主函数
schedulingCyclebindingCycle
失败
FailureHandler
 → 回队列
unreserveAndForget
 + FailureHandler

整体流程如下:

Pod调度中各阶段涉及异常码如下:

type Code intconst (// 执行正常 Success Code = iota// 未知异常,一般进入backOffQ Error// 一般在PreFilter或Filter插件返回,不可调度但可以PostFilter(DefaultPreemption)尝试抢占节点 Unschedulable// 不可调度且无法PostFilter(DefaultPreemption)抢占节点 UnschedulableAndUnresolvable// 调度周期中Permit插件决定pod需要等待一段时间// 绑定周期中会等待这个时间,超时则本周期失败 Wait// 1. Bind插件跳过自身// 2. PreFilter插件 决定跳过 自己的Filter// 3. PreScore插件 决定跳过 自己的Score Skip)

schedule_one.go:Pod调度失败:

1)更新pod信息,根据QueueingHintFn进入activeQ/backoffQ/unschedulablePods;

2)记录调度失败事件,reason=FailedScheduling;

3)更新Pod的Condition,PodScheduled=false;

func(sched *Scheduler)handleSchedulingFailure(...) {// 异常SchedulerError reason := v1.PodReasonSchedulerErrorif status.IsRejected() {// 不可调度Unschedulable  reason = v1.PodReasonUnschedulable } pod := podInfo.Pod// 异常信息,常见是FitError.Error() errMsg := status.Message() podLister := podFwk.SharedInformerFactory().Core().V1().Pods().Lister() cachedPod, e := podLister.Pods(pod.Namespace).Get(pod.Name)// 更新pod信息    podInfo.PodInfo, _ = framework.NewPodInfo(cachedPod.DeepCopy())    pod = podInfo.Pod// 根据QueueingHintFn进入activeQ/backoffQ/unschedulablePods    sched.SchedulingQueue.AddUnschedulableIfNotPresent(logger, podInfo, sched.SchedulingQueue.SchedulingCycle());// 记录 抢占节点-pods+当前pod 当前pod-抢占节点    sched.SchedulingQueue.AddNominatedPod(logger, podInfo.PodInfo, nominatingInfo) msg := truncateMessage(errMsg)// 记录调度失败事件// kubectl get events --field-selector reason=FailedScheduling podFwk.EventRecorder().WithLogger(logger).Eventf(pod, nil, v1.EventTypeWarning, "FailedScheduling""Scheduling", msg)// 更新Pod的Condition PodScheduled=falseif err := updatePod(ctx, sched.client, podFwk.APICacher(), pod, &v1.PodCondition{  Type:               v1.PodScheduled,  Status:             v1.ConditionFalse,  Reason:             reason,  Message:            errMsg, }, nominatingInfo);}

3.2 调度周期-主流程

schedule_one.go:调度周期主流程。

  1. UpdateSnapshot:将cache中的node列表保存到Scheduler自己的快照里,本次调度周期中以这个node列表快照为准;

  2. schedulingAlgorithm:Filter+Score计算pod应该调度到哪个Node上,如果失败,则返回上层进入backoffQ/unschedulablePods;

  3. prepareForBindingCycle:assume内存构造新的podInfo(pod.Spec.NodeName=选中节点),执行ReservePlugin和PermitPlugin;

// schedulingCycle tries to schedule a single Pod.func(sched *Scheduler)schedulingCycle(...)(ScheduleResult, *framework.QueuedPodInfo, *fwk.Status) {// 1. 更新内存 node列表 快照nodeInfoSnapshot,后面用这个快照来算调度if err := sched.Cache.UpdateSnapshot(klog.FromContext(ctx), sched.nodeInfoSnapshot); err != nil {return ScheduleResult{nominatingInfo: clearNominatedNode}, podInfo, fwk.AsStatus(err) }// 2. 计算pod调度到哪个node上(Filter -> Score) scheduleResult, status := sched.schedulingAlgorithm(ctx, state, schedFramework, podInfo, start)if !status.IsSuccess() {return scheduleResult, podInfo, status }// 3. assume->reserve->permit// assume-更新cache里的pod.Spec.NodeName// reserve - ReservePlugin --- 更新插件自己的状态// permit - PermitPlugin --- 可 批准 / 拒绝 / 等待(带超时),全部批准才进入绑定 assumedPodInfo, status := sched.prepareForBindingCycle(ctx, state, schedFramework, podInfo, podsToActivate, scheduleResult)if !status.IsSuccess() {return ScheduleResult{nominatingInfo: clearNominatedNode}, assumedPodInfo, status }return scheduleResult, assumedPodInfo, nil}

schedule_one.go:schedulingAlgorithm

  1. findNodesThatFitPod:执行PreFilter+Filter插件,从node列表中过滤出部分node,如果node数量为1,直接选中,不走后续流程;

  2. RunPostFilterPlugins:如果过滤结果为空,执行PostFilter,可能得到抢占节点nominatingInfo;

  3. prioritizeNodes:如果1中返回node数量>1,执行PreScore+Score插件,给各个node评分,选出最高分node返回;

// 调度结果type ScheduleResult struct {// Filter+Score 选中节点 SuggestedHost string// Filter失败,PostFilter抢占节点 nominatingInfo *fwk.NominatingInfo}// 抢占节点type NominatingInfo struct { NominatedNodeName string NominatingMode    NominatingMode}func(sched *Scheduler)schedulingAlgorithm(...)(ScheduleResult, *fwk.Status) { pod := podInfo.Pod logger := klog.FromContext(ctx)// Filter+Score scheduleResult, err := sched.SchedulePod(ctx, schedFramework, state, podInfo)if err != nil {if err == ErrNoNodesAvailable {// 无可用节点,直接返回   status := fwk.NewStatus(fwk.UnschedulableAndUnresolvable).WithError(err)return ScheduleResult{nominatingInfo: clearNominatedNode}, status  }  fitError, ok := err.(*framework.FitError)if !ok {// 如果不是Filter返回空node列表,不走PostFilter   logger.Error(err, "Error selecting node for pod""pod", klog.KObj(pod))return ScheduleResult{nominatingInfo: clearNominatedNode}, fwk.AsStatus(err)  }// 只有Filter返回空node列表,才走PostFilter执行抢占逻辑  result, status := schedFramework.RunPostFilterPlugins(ctx, state, pod, fitError.Diagnosis.NodeToStatus)  msg := status.Message()  fitError.Diagnosis.PostFilterMsg = msgvar nominatingInfo *fwk.NominatingInfoif result != nil {// 需要抢占节点   nominatingInfo = result.NominatingInfo  }return ScheduleResult{nominatingInfo: nominatingInfo}, fwk.NewStatus(fwk.Unschedulable).WithError(err) }return scheduleResult, nil}// Filter+Scorefunc(sched *Scheduler)schedulePod(...)(result ScheduleResult, err error) { pod := podInfo.Pod// 无可用节点if sched.nodeInfoSnapshot.NumNodesInPlacement() == 0 {return result, ErrNoNodesAvailable }// 执行Filter,返回可用节点feasibleNodes feasibleNodes, diagnosis, nodeHint, err := sched.findNodesThatFitPod(ctx, fwk, state, podInfo)if err != nil {return result, err }// 如果Filter后没有可用节点,返回FitError,执行PostFilter抢占iflen(feasibleNodes) == 0 {return result, &framework.FitError{   Pod:         pod,   NumAllNodes: sched.nodeInfoSnapshot.NumNodesInPlacement(),   Diagnosis:   diagnosis,  } }// 如果Filter后只有一个Node,不走Score,直接返回iflen(feasibleNodes) == 1 {  node := feasibleNodes[0].Node().Namereturn ScheduleResult{   SuggestedHost:  node,  }, nil }// 对feasibleNodes中的节点,执行PreScore + Score priorityList, err := prioritizeNodes(ctx, sched.Extenders, fwk, state, pod, feasibleNodes)if err != nil {return result, err }// 排序后选择评分最高的节点返回 sortedPrioritizedNodes := newSortedNodeScores(priorityList) node := sortedPrioritizedNodes.Pop()return ScheduleResult{  SuggestedHost:  node, }, err}

schedule_one.go:prepareForBindingCycle

func(sched *Scheduler)prepareForBindingCycle( ctx context.Context, state fwk.CycleState, schedFramework framework.Framework, podInfo *framework.QueuedPodInfo, podsToActivate *framework.PodsToActivate, scheduleResult ScheduleResult,)(*framework.QueuedPodInfo, *fwk.Status) {// 1.// assume: 更新内存cache里pod.Spec.NodeName = scheduleResult.SuggestedHost// reserve: ReservePlugin assumedPodInfo, status := sched.assumeAndReserve(ctx, state, schedFramework, podInfo, scheduleResult)if !status.IsSuccess() {return assumedPodInfo, status } assumedPod := assumedPodInfo.Pod// 2. PermitPlugin 返回是否需要等待 等待多长时间 pluginsWaitTime, runPermitStatus := schedFramework.RunPermitPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)if runPermitStatus.IsWait() {  schedFramework.AddWaitingPod(assumedPod, pluginsWaitTime) } elseif !runPermitStatus.IsSuccess() {// 回滚第一步的assumeAndReserve  err := sched.unreserveAndForget(ctx, state, schedFramework, assumedPodInfo, scheduleResult.SuggestedHost)if runPermitStatus.IsRejected() {   fitErr := &framework.FitError{    NumAllNodes: 1,    Pod:         podInfo.Pod,    Diagnosis: framework.Diagnosis{     NodeToStatus: framework.NewDefaultNodeToStatus(),    },   }   fitErr.Diagnosis.NodeToStatus.Set(scheduleResult.SuggestedHost, runPermitStatus)   fitErr.Diagnosis.AddPluginStatus(runPermitStatus)return assumedPodInfo, fwk.NewStatus(runPermitStatus.Code()).WithError(fitErr)  }return assumedPodInfo, runPermitStatus }return assumedPodInfo, nil}

3.3 调度周期-Filter阶段(预选)

schedule_one.gofindNodesThatFitPod先PreFilter,再Filter,返回feasibleNodes可行节点列表;可行节点为空时由上层 schedulingAlgorithm跑PostFilter(抢占)。

func(sched *Scheduler)findNodesThatFitPod(...)([]fwk.NodeInfo, framework.Diagnosis, string, error) {// 1. PreFilter:预处理 / 缩候选集 preRes, s, unscheduledPlugins := schedFramework.RunPreFilterPlugins(ctx, state, pod)if !s.IsSuccess() {returnnil, diagnosis, "", ... }// 2. 优先试 NominatedNode(对于该pod,上轮PostFilter抢占写下的抢占节点名)iflen(pod.Status.NominatedNodeName) > 0 {// 跑Filter和Extender,如果成功了,直接返回抢占节点  feasibleNodes, _ := sched.evaluateNominatedNode(...)iflen(feasibleNodes) != 0 {return feasibleNodes, diagnosis, ...  } } nodes := allNodesif !preRes.AllNodes() {// 对于PreFilter结果,校验node在本轮snapshot  nodes = make([]fwk.NodeInfo, 0len(preRes.NodeNames))for nodeName := range preRes.NodeNames {if nodeInfo, err := sched.nodeInfoSnapshot.GetNodeInPlacement(nodeName); err == nil {    nodes = append(nodes, nodeInfo)   }  } }// 3. 按 PreFilterResult 缩小 nodes 后,并行跑 Filter feasibleNodes, err := sched.findNodesThatPassFilters(ctx, ..., nodes)return feasibleNodes, diagnosis, ...}// 上层 schedulingAlgorithm:Filter 结果为空时才跑 PostFilterscheduleResult, err := sched.SchedulePod(...) if fitError, ok := err.(*framework.FitError); ok { result, status := schedFramework.RunPostFilterPlugins(...) // 抢占等}

PreFilter

framework.go:RunPreFilterPlugins按序跑每个 PreFilter 插件,缩小node范围,根据插件返回状态

  1. Skip:后面跳过该插件Filter;

  2. UnschedulableAndUnresolvable:返回上层,PostFilter默认抢占逻辑不生效;

  3. Unschedulable:返回上层,PostFilter默认可抢占;

  4. N个PreFilter后,node交集为空,同UnschedulableAndUnresolvable;

func(f *frameworkImpl)RunPreFilterPlugins(            ctx context.Context, state fwk.CycleState, pod *v1.Pod)          (_ *fwk.PreFilterResult, status *fwk.Status, _ sets.Set[string]) { skipPlugins := sets.New[string]() nodes, err := f.SnapshotSharedLister().NodeInfos().List()var result *fwk.PreFilterResult// 有返回缩小node范围的插件 pluginsWithNodes := sets.New[string]()var returnStatus *fwk.Statusfor _, pl := range f.preFilterPlugins {  r, s := f.runPreFilterPlugin(ctx, pl, state, pod, nodes)if s.IsSkip() {// 不执行这个插件的Filter阶段   skipPlugins.Insert(pl.Name())continue  }if !s.IsSuccess() {   s.SetPlugin(pl.Name())if s.Code() == fwk.UnschedulableAndUnresolvable {// UnschedulableAndUnresolvablereturnnil, s, nil   }if s.Code() == fwk.Unschedulable {// 一个插件Unschedulable,允许跑完所有PreFilter再返回    returnStatus = scontinue   }// 异常直接返回returnnil, fwk.AsStatus(fmt.Errorf().WithPlugin(pl.Name()), nil  }if !r.AllNodes() {   pluginsWithNodes.Insert(pl.Name())  }// node交集  result = result.Merge(r)if !result.AllNodes() && len(result.NodeNames) == 0 {// 交集为空,返回UnschedulableAndUnresolvablereturn result, fwk.NewStatus(fwk.UnschedulableAndUnresolvable, msg), pluginsWithNodes  } }return result, returnStatus, pluginsWithNodes}

内置插件示例node_affinity.go

无约束 Skip;有约束则写 CycleState;能钉到具体节点名则返回 NodeNames 缩候选。

CycleState:单次调度周期内的共享状态(scheduleOnePod 里 NewCycleState() 创建)。每个插件用 Write/Read 按 key 存自己的中间结果,调度周期结束即丢弃,不跨 Pod、不持久化。

func(pl *NodeAffinity)PreFilter(...)(*fwk.PreFilterResult, *fwk.Status) { affinity := pod.Spec.Affinity noNodeAffinity := (affinity == nil ||  affinity.NodeAffinity == nil ||  affinity.NodeAffinity.RequiredDuringSchedulingIgnoredDuringExecution == nil)// 1) 无 NodeSelector、无 required NodeAffinity、插件也未注入 selector → Skip//    框架会跳过本插件后续 Filter / PreFilterExtensionsif noNodeAffinity && pl.addedNodeSelector == nil && pod.Spec.NodeSelector == nil {returnnil, fwk.NewStatus(fwk.Skip) }// 2) 解析 required 亲和,写入 CycleState,供 Filter 阶段 Match 复用 state := &preFilterState{requiredNodeSelectorAndAffinity: nodeaffinity.GetRequiredNodeAffinity(pod)} cycleState.Write(preFilterStateKey, state)// 只有 NodeSelector、没有 required NodeAffinity terms → Success 但不缩节点集if noNodeAffinity || len(affinity.NodeAffinity.RequiredDuringSchedulingIgnoredDuringExecution.NodeSelectorTerms) == 0 {returnnilnil// Success,PreFilterResult 为空 = 全量节点 }// 3) 从 MatchFields 里抠 metadata.name In [...],terms 之间 OR、同一 term 内 AND terms := affinity.NodeAffinity.RequiredDuringSchedulingIgnoredDuringExecution.NodeSelectorTermsvar nodeNames sets.Set[string]for _, t := range terms {var termNodeNames sets.Set[string]for _, r := range t.MatchFields {if r.Key == metav1.ObjectNameField && r.Operator == v1.NodeSelectorOpIn {    s := sets.New(r.Values...)if termNodeNames == nil {     termNodeNames = s    } else {     termNodeNames = termNodeNames.Intersection(s) // 同 term 内求交    }   }  }if termNodeNames == nil {// 本 term 未钉死节点名(例如只按 label)→ 无法缩集,全量节点returnnilnil  }  nodeNames = nodeNames.Union(termNodeNames) // terms 之间求并 }if nodeNames != nil && len(nodeNames) == 0 {// 各 term 对 node.Name 约束互相矛盾 → 永远匹配不到节点returnnil, fwk.NewStatus(fwk.UnschedulableAndUnresolvable, errReasonConflict) } elseiflen(nodeNames) > 0 {// 钉死候选节点名 → 框架后续 Filter 只扫这些 Node// 典型配置:required NodeAffinity 用 matchFields 指定 metadata.name In [node-a]return &fwk.PreFilterResult{NodeNames: nodeNames}, nil }returnnilnil}

Filter

scheduler_one.gofindNodesThatPassFilters并行对每个节点执行Filter。

func(sched *Scheduler)findNodesThatPassFilters(...)([]fwk.NodeInfo, error) { numAllNodes := len(nodes)// 并不要求扫描所有node// 1. 集群小于100节点需要全部扫描;2. 超过100有特定公式计算 numNodesToFind := sched.numFeasibleNodesToFind(schedFramework.PercentageOfNodesToScore(), int32(numAllNodes))// node结果集 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)// Filter某个node checkNode := func(i int) {  nodeInfo := nodes[(sched.nextStartNodeIndex+i)%numAllNodes]// Filter执行  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 {    cancel(errors.New("findNodesThatPassFilters has found enough nodes"))    atomic.AddInt32(&feasibleNodesLen, -1)   } else {    feasibleNodes[length-1] = nodeInfo   }  } else {// 如果是Unschedulable,后面PostFilter会尝试抢占这个node   result[i] = &nodeStatus{node: nodeInfo.Node().Name, status: status}  } }// 并行处理节点Filter schedFramework.Parallelizer().Until(ctx, numAllNodes, checkNode, metrics.Filter) feasibleNodes = feasibleNodes[:feasibleNodesLen]for _, item := range result {if item == nil {continue  }  diagnosis.NodeToStatus.Set(item.node, item.status)  diagnosis.AddPluginStatus(item.status) }if err := errCh.Receive(); err != nil {return feasibleNodes, err }return feasibleNodes, nil}

scheduler_one.goRunFilterPluginsWithNominatedPods单个节点调用N个Filter插件,任一插件非 Success 则该节点失败并短路。最多跑两次Filter,先加上节点上更高优先级的抢占 Pod(可能是之前调度周期的PostFilter)做一次,再不加跑一次。

Filter阶段需要找出运行这个Pod的Node集合,只要有一个Filter插件返回Node不满足,Node就不能被选中。如果该Node结果只是Unschedulable,代表还能走PostFilter,如果Pod优先级高,能抢占这个Node,驱逐这个Node上的其他Pod。

func(f *frameworkImpl)RunFilterPluginsWithNominatedPods(ctx context.Context, state fwk.CycleState, pod *v1.Pod, info fwk.NodeInfo) *fwk.Status {var status *fwk.Status podsAdded := false// 第一次Filter,加上之前抢占这个node的pods// 第二次Filter,不加上抢占podsfor i := 0; i < 2; i++ {  stateToUse := state  nodeInfoToUse := infoif i == 0 {var err error   podsAdded, stateToUse, nodeInfoToUse, err = addGENominatedPods(ctx, f, pod, state, info)if err != nil {return fwk.AsStatus(err)   }  } elseif !podsAdded || !status.IsSuccess() {// 如果第一次没有抢占pod加入,或抢占pod加入该节点无法满足Filter,则不进行第二次尝试break  }  status = f.RunFilterPlugins(ctx, stateToUse, pod, nodeInfoToUse)if !status.IsSuccess() && !status.IsRejected() {return status  } }return status}func(f *frameworkImpl)RunFilterPlugins(...) *fwk.Status {for _, pl := range f.filterPlugins {if state.GetSkipFilterPlugins().Has(pl.Name()) {continue// PreFilter Skip 过来的  }if status := f.runFilterPlugin(...); !status.IsSuccess() {   status.SetPlugin(pl.Name())return status // 该节点淘汰,不再跑后续 Filter 插件  } }returnnil}

PostFilter

framework.go当Filter后节点列表为空,则进入PostFilter,默认内置插件只有一个DefaultPreemption抢占实现。

func(f *frameworkImpl)RunPostFilterPlugins(...)(*fwk.PostFilterResult, *fwk.Status) {if state.ShouldSkipAllPostFilterPlugins() {returnnil, Unschedulable // 全部跳过,不跑插件 }var result *fwk.PostFilterResult // 记下「最后一个有意义的」结果(非 ModeNoop)var reasons []stringvar rejectorPlugin stringfor _, pl := range f.postFilterPlugins {  r, s := f.runPostFilterPlugin(...)if s.IsSuccess() {return r, s // 立刻停:如抢占成功,r 常带 NominatingInfo  } elseif s.Code() == fwk.UnschedulableAndUnresolvable {return r, s.WithPlugin(pl.Name()) // 立刻停:抢占也解决不了  } elseif !s.IsRejected() {// 非 Success / Unschedulable* / Pending → 当 Error,立刻停;丢弃 resultreturnnil, fwk.AsStatus(s.AsError()).WithPlugin(pl.Name())  }// 走到这里:IsRejected 且不是 UnschedulableAndUnresolvable// 即 Unschedulable 或 Pending → 不立刻返回,继续下一个插件if r != nil && r.Mode() != fwk.ModeNoop {   result = r // 覆盖:只保留最后一个非 Noop 的 PostFilterResult(如提名节点)  }  reasons = append(reasons, s.Reasons()...)if rejectorPlugin == "" {   rejectorPlugin = pl.Name() // 记录第一个拒绝者  } }// 全部插件都是 Unschedulable/Pending:合并 reasons,带上「最后一个」非 Noop 的 resultreturn result, fwk.NewStatus(fwk.Unschedulable, reasons...).WithPlugin(rejectorPlugin)}

上层 schedulingAlgorithm:无论 PostFilter返回什么,本周期最终都Unschedulable,但是会更新pod.status.nominatedNodeName,持久化抢占节点名。

本轮pod调度失败,需要等感知到被抢占节点的pod被删除后,该pod重新进入调度队列,在Filter阶段会优先尝试被抢占节点(findNodesThatFitPod→evaluateNominatedNode)。

return ScheduleResult{nominatingInfo: result.NominatingInfo}, fwk.NewStatus(fwk.Unschedulable).WithError(fitError)

preemption.go:默认PostFilter实现是抢占。

1)PodEligibleToPreemptOthers:如果pod.spec.preemptionPolicy=Never,不执行抢占;如果pod已经有抢占节点(之前的调度周期),且抢占在执行中,则不执行抢占;

2)findCandidates:循环node列表,对每个节点选出驱逐pods;(如果先前是PreFilter返回Unschedulable,这里node列表是所有节点,但是不是全部扫描,集群小于1000节点,最多扫描100个,大于1000节点,用0.1*节点数;如果先前Filter返回Unschedulable,则针对Filter跑的node列表)

3)SelectCandidate:从2里选择最合适的一个node;

4)actuatePodPreemption:执行3删除node上需要被驱逐的pods;

func(ev *Evaluator)Preempt(...)(*fwk.PostFilterResult, *fwk.Status) { nominatedNodeStatus := m.Get(pod.Status.NominatedNodeName)// 1)pod抢占的节点上有 比自己优先级低 被scheduler驱逐 的 terminating状态 pod// 2)preemptionPolicy=Never// 返回Unschedulableif ok, msg := ev.PodEligibleToPreemptOthers(ctx, pod, nominatedNodeStatus); !ok {returnnil, fwk.NewStatus(fwk.Unschedulable, msg) } allNodes, err := ev.Handler.SnapshotSharedLister().NodeInfos().List()if err != nil {returnnil, fwk.AsStatus(err) }// 一个 Candidate = 节点 + 节点上被驱逐pods + PDB违反数// 通过删掉这个节点上的这些受害者,让 pod 可以调度到这个节点 candidates, nodeToStatusMap, err := ev.findCandidates(ctx, state, allNodes, pod, m)if err != nil && len(candidates) == 0 {returnnil, fwk.AsStatus(err) }iflen(candidates) == 0 {return framework.NewPostFilterResultWithNominatedNode(""), fwk.NewStatus(fwk.Unschedulable, fitError.Error()) }// 选一个Candidate 驱逐上面的pods bestCandidate := ev.SelectCandidate(ctx, candidates)if bestCandidate == nil || len(bestCandidate.Name()) == 0 {returnnil, fwk.NewStatus(fwk.Unschedulable, "no candidate node for preemption") }// 执行抢占 删除podif status := ev.executor.actuatePodPreemption(ctx, bestCandidate.Name(), bestCandidate.Victims(), pod, ev.PluginName); !status.IsSuccess() {returnnil, status }return framework.NewPostFilterResultWithNominatedNode(bestCandidate.Name()), fwk.NewStatus(fwk.Success)}

default_preemption.go:findCandidates→SelectVictimsOnNode,挑选node上被驱逐的pods

1)先把比当前pod优先级低的pods都从node上删了,跑一次Filter,如果行不通,则这个node不能被抢占;

2)候选pods按照重要性排序,优先级越高越优先,启动时间越早越优先;

3)候选pods一个一个加回去(先加重要的),执行Filter,如果失败则进入victims需要被驱逐,如果成功则不需要被驱逐;

func(pl *DefaultPreemption)SelectVictimsOnNode(....)([]*v1.Pod, int, *fwk.Status) {var potentialVictims []fwk.PodInfo removePod := func(rpi fwk.PodInfo)error {  nodeInfo.RemovePod(logger, rpi.GetPod()); } addPod := func(api fwk.PodInfo)error {  nodeInfo.AddPodInfo(api); }// 1. 先扫描当前节点上的所有 Pod,只挑出“优先级低于抢占者”的 Pod 作为潜在受害者。for _, pi := range nodeInfo.GetPods() {if pl.isPreemptionAllowed(nodeInfo, pi, pod) {   potentialVictims = append(potentialVictims, pi)  } }for _, pi := range potentialVictims {  removePod(pi); }// 一般如果没有配置优先级,都是0,不会发生抢占iflen(potentialVictims) == 0 {returnnil0, fwk.NewStatus(fwk.UnschedulableAndUnresolvable) }// 2. 如果把所有潜在受害者都移走后,新 Pod 还是放不进来,说明这个节点不适合通过抢占解决。if status := pl.fh.RunFilterPluginsWithNominatedPods(ctx, state, pod, nodeInfo); !status.IsSuccess() {returnnil0, status }var victims []fwk.PodInfo numViolatingVictim := 0// 按优先级降序排列,优先级高->优先级低,先启动->后启动// 后面优先把优先级高的加回去 sort.Slice(potentialVictims, func(i, j int)bool {return pl.MoreImportantPod(potentialVictims[i].GetPod(), potentialVictims[j].GetPod()) })// 3. 按 PDB 是否会被违反分成两组 violatingVictims, nonViolatingVictims :=             filterPodsWithPDBViolation(potentialVictims, pdbs) reprievePod := func(pi fwk.PodInfo)(bool, error) {// pod加回去if err := addPod(pi);// 跑一次Filter  status := pl.fh.RunFilterPluginsWithNominatedPods(ctx, state, pod, nodeInfo)  fits := status.IsSuccess()if !fits {// 如果不行,pod还是要进入victims被驱逐if err := removePod(pi);    victims = append(victims, pi)  }return fits, nil }// 4. 逐个把 Pod 加回节点做试探;如果加回后新 Pod 就放不下了,这个 Pod 就必须保留在 victims 里。for _, p := range violatingVictims {if fits, err := reprievePod(p); err != nil {returnnil0, fwk.AsStatus(err)  } elseif !fits {   numViolatingVictim++  } }for _, p := range nonViolatingVictims {if _, err := reprievePod(p); }// 5. 最终返回“确实不能加回去”的那批 Pod,它们就是该节点上需要被抢占的 victims。var victimPods []*v1.Podfor _, pi := range victims {  victimPods = append(victimPods, pi.GetPod()) }return victimPods, numViolatingVictim, fwk.NewStatus(fwk.Success)}

preemption.go:SelectCandidate→pickOneNodeForPreemption,对于候选node的驱逐方案,按照下面的优先级选择

1)违反PDB数最少的节点;

2)驱逐Pod优先级最低的所在节点;

3)驱逐Pod优先级总和最少的节点;

4)驱逐Pod最少的节点;

5)驱逐Pod启动时间最晚的所在节点;

funcpickOneNodeForPreemption(...)string { minNumPDBViolatingScoreFunc := func(node string)int64 {return -nodesToVictims[node].NumPDBViolations } minHighestPriorityScoreFunc := func(node string)int64 {  highestPodPriority := corev1helpers.PodPriority(nodesToVictims[node].Pods[0])return -int64(highestPodPriority) } minSumPrioritiesScoreFunc := func(node string)int64 {var sumPriorities int64for _, pod := range nodesToVictims[node].Pods {   sumPriorities += int64(corev1helpers.PodPriority(pod)) + int64(math.MaxInt32+1)  }return -sumPriorities } minNumPodsScoreFunc := func(node string)int64 {return -int64(len(nodesToVictims[node].Pods)) } latestStartTimeScoreFunc := func(node string)int64 {  earliestStartTimeOnNode := util.GetEarliestPodStartTime(nodesToVictims[node])return earliestStartTimeOnNode.UnixNano() } scoreFuncs = []func(string)int64{  minNumPDBViolatingScoreFunc, // 违反PDB数最少的节点  minHighestPriorityScoreFunc, // 驱逐Pod优先级最低的所在节点  minSumPrioritiesScoreFunc, // 驱逐Pod优先级总和最少的节点  minNumPodsScoreFunc, // 驱逐Pod最少的节点  latestStartTimeScoreFunc, // 驱逐Pod启动时间最晚的所在节点 }for _, f := range scoreFuncs {  selectedNodes := []string{}  maxScore := int64(math.MinInt64)for _, node := range allCandidates {   score := f(node)if score > maxScore {    maxScore = score    selectedNodes = []string{}   }if score == maxScore {    selectedNodes = append(selectedNodes, node)   }  }iflen(selectedNodes) == 1 {return selectedNodes[0]  }  allCandidates = selectedNodes }return allCandidates[0]}

3.4 调度周期-Score阶段(优选)

prioritizeNodes:PreScore → Score → NormalizeScore → 加权求和 → 选最高分节点

仅当 Filter 后 可行节点 > 1 才打分;只有 1 个则 schedulePod 直接选用。

PreScore

PreScore:如果插件返回Skip,跳过对应插件的Score阶段;如果Error,本调度周期失败,回队列。

func(f *frameworkImpl)RunPreScorePlugins(...) *fwk.Status {for _, pl := range f.preScorePlugins {  status = f.runPreScorePlugin(...)if status.IsSkip() {   skipPlugins.Insert(pl.Name()) // 本周期跳过同名 Scorecontinue  }if !status.IsSuccess() {return Error // 非 Success/Skip → 中止打分,本调度周期失败  } } state.SetSkipScorePlugins(skipPlugins)returnnil}

Score

framework.go:循环节点Score打分 → 循环插件NormalizeScore归一化打分 → 加权求和,最终得到每个节点的打分情况。

type NodePluginScores struct {// 节点名 Name string// 每个插件的打分 Scores []PluginScore// 加权求和总分 TotalScore int64}func(f *frameworkImpl)RunScorePlugins(...)([]fwk.NodePluginScores, *fwk.Status) {// 过滤掉 PreScore 标记 Skip 的插件 plugins := scorePlugins - SkipScorePlugins pluginToNodeScores := map[pluginName][]NodeScore{} // 每个插件对所有节点的分// ── 1. Score:按节点并行,节点内串行跑各插件 ── Parallelizer.Until(len(nodes), func(i int) {for _, pl := range plugins {   s, status := pl.Score(pod, nodes[i])if !status.IsSuccess() {    cancel(); return Error // 任一失败 → 整轮打分失败   }   pluginToNodeScores[pl.Name()][i] = NodeScore{Name: nodeName, Score: s}  } })// ── 2. NormalizeScore:按插件并行,把该插件在各节点的原始分归一化到 [0,100] ── Parallelizer.Until(len(plugins), func(i int) {  pl := plugins[i]if pl.ScoreExtensions() == nil {return// 无可选扩展则跳过  }  status := pl.ScoreExtensions().NormalizeScore(pod, pluginToNodeScores[pl.Name()])if !status.IsSuccess() {   cancel(); return Error // 非 Success → 整轮失败  } })// ── 3. 加权求和:按节点并行,校验区间后 score * weight,累加 TotalScore ── Parallelizer.Until(len(nodes), func(i int) {var total int64for _, pl := range plugins {   score := pluginToNodeScores[pl.Name()][i].Score// 要求每个插件对每个节点的打分,必须在[0, 100]之间if score < MinNodeScore || score > MaxNodeScore {    cancel(); return Error   }// 每个插件有自己的权重配置   weighted := score * int64(scorePluginWeight[pl.Name()])// TotalScore = Σ (normalizedScore × weight)   total += weighted  }  allNodePluginScores[i] = NodePluginScores{Name: nodeName, TotalScore: total, ...} })return allNodePluginScores, nil}

least_allocated.go:内置插件NodeResourcesFit,默认使用LeastAllocated策略,打散负载,尽量将Pod调度到资源相对充足的节点上。对每种资源算完分后做加权平均:nodeScore = Σ(resourceScore_i × weight_i) / Σ(weight_i)。默认资源是 CPU、Memory,权重都是 1,所以得分=(cpu((capacity-requested) × 100/capacity) + memory((capacity-requested) × 100/capacity))/2。

比如request 1c1g,现在有两个node,规格和评分如下:

1)4c8g:((4-1)*100/4 + (8-1)*100/8) / 2 = 81.25

2)2c4g:((2-1)*100/2 + (4-1)*100/4) / 2 = 62.5

funcleastResourceScorer(resources []config.ResourceSpec)func([]int64, []int64, []int64)int64 {// 参数 1. 请求资源 2. 已分配资源 3. 可分配资源returnfunc(requested, _, allocable []int64)int64 {var nodeScore, weightSum int64for i := range requested {if allocable[i] == 0 {continue              }              weight := resources[i].Weight// (可用资源 - 请求请求) / 可用资源 * 100              resourceScore := leastRequestedScore(requested[i], allocable[i])// cpu和内存默认权重都是1              nodeScore += resourceScore * weight              weightSum += weight          }if weightSum == 0 {return0          }return nodeScore / weightSum      }}funcleastRequestedScore(requested, capacity int64)int64 {if capacity == 0 {return0 }if requested > capacity {return0 }// MaxNodeScore=100return ((capacity - requested) * fwk.MaxNodeScore) / capacity}

3.5 调度周期-Assume→Reserve→Permit

选中节点后:先改内存(Assume),再让插件预留状态(Reserve),最后 Permit 批准/等待/拒绝

Assume&Reserve

scheduler_one.go

  1. Assume:内存更新,pod绑定nodeName,pod挂到node上,更新node的资源使用情况;

  2. Reserve:执行ReservePlugin预留资源,某个插件预留失败,执行插件的Unreserve方法,并回滚Assume内存状态;

完成了Assume和Reserve,后面绑定周期才能并发异步执行。

func(sched *Scheduler)assumeAndReserve(...)(*framework.QueuedPodInfo, *fwk.Status) { logger := klog.FromContext(ctx)// 复制新pod信息,更新到cache缓存 assumedPodInfo := podInfo.DeepCopy() assumedPod := assumedPodInfo.Pod err := sched.assume(logger, state, assumedPodInfo, scheduleResult.SuggestedHost)if err != nil {return assumedPodInfo, fwk.AsStatus(err) }// 执行reserve插件,预留资源if sts := schedFramework.RunReservePluginsReserve(ctx, state, assumedPod, scheduleResult.SuggestedHost); !sts.IsSuccess() {// 如果发生异常,执行reserve插件的unreserve方法,回滚状态  err := sched.unreserveAndForget(ctx, state, schedFramework, assumedPodInfo, scheduleResult.SuggestedHost)if sts.IsRejected() {return assumedPodInfo, fwk.NewStatus(sts.Code()).WithError(fitErr)  }return assumedPodInfo, sts }// 正常情况,返回新pod信息,spec.nodeName已经赋值return assumedPodInfo, nil}func(sched *Scheduler)assume(logger klog.Logger, state fwk.CycleState, assumedPodInfo *framework.QueuedPodInfo, host string)error {// 设置nodeName assumedPodInfo.Pod.Spec.NodeName = host// 更新缓存pod信息,将pod挂到某个node上,// 并内存更新这个node的资源情况,比如cpu、内存、使用HostPortif err := sched.Cache.AssumePod(logger, assumedPodInfo.Pod); err != nil {return err }returnnil}func(f *frameworkImpl)RunReservePluginsReserve(...)(status *fwk.Status) {for _, pl := range f.reservePlugins {  status = f.runReservePluginReserve(ctx, pl, state, pod, nodeName)if !status.IsSuccess() {// 单插件不成功,返回errorif status.IsRejected() {    status.SetPlugin(pl.Name())return status   }   err := status.AsError()return fwk.AsStatus(fmt.Errorf("running Reserve plugin %q: %w", pl.Name(), err))  } }returnnil}

Permit

framework.go:执行Permit插件,如果插件返回Wait,则收集插件需要等待的时间。

func(f *frameworkImpl)RunPermitPlugins(...)(pluginsWaitTime map[string]time.Duration, status *fwk.Status) {var waitStatus *fwk.Status pluginsWaitTime = make(map[string]time.Duration)for _, pl := range f.permitPlugins {  status, timeout := f.runPermitPlugin(ctx, pl, state, pod, nodeName)if !status.IsSuccess() {if status.IsRejected() {// Unschedulable直接返回returnnil, status.WithPlugin(pl.Name())   }if status.IsWait() {// Wait,记录每个插件的等待时间    pluginsWaitTime[pl.Name()] = timeout    waitStatus = status   } else {// Error直接返回    err := status.AsError()returnnil, fwk.AsStatus(fmt.Errorf("running Permit plugin %q: %w", pl.Name(), err)).WithPlugin(pl.Name())   }  } }if waitStatus.IsWait() {return pluginsWaitTime, waitStatus }returnnilnil}

scheduler_one.go:如果Wait则加入waitingPods并开启定时器,如果异常,则回滚assume和reserve。

func(sched *Scheduler)prepareForBindingCycle(...)(...) { pluginsWaitTime, runPermitStatus := schedFramework.RunPermitPlugins(...)if runPermitStatus.IsWait() {  schedFramework.AddWaitingPod(assumedPod, pluginsWaitTime) } elseif !runPermitStatus.IsSuccess() {  sched.unreserveAndForget(...) // 回滚 Assume + Reservereturn ..., runPermitStatus }}func(f *frameworkImpl)AddWaitingPod(pod *v1.Pod, pluginsWaitTime map[string]time.Duration) { waitingPod := newWaitingPod(pod, pluginsWaitTime) f.waitingPods.add(waitingPod)}funcnewWaitingPod(pod *v1.Pod, pluginsMaxWaitTime map[string]time.Duration) *waitingPod { wp := &waitingPod{  pod: pod,// chan可等待permit结束  s: make(chan *fwk.Status, 1), } wp.pendingPlugins = make(map[string]*time.Timer, len(pluginsMaxWaitTime))for k, v := range pluginsMaxWaitTime {  plugin, waitTime := k, v// 对于每个插件,开启定时器,超时Reject发送到s这个chan  wp.pendingPlugins[plugin] = time.AfterFunc(waitTime, func() {   wp.Reject(plugin, msg)  }) }return wp}

至此调度周期结束。

3.6 绑定周期

scheduler_one.go:绑定周期并行执行,如果处理失败

  1. unreserveAndForget:调用ReservePlugin的Unreserve方法,清理内存node中的pod;

  2. 发送EventAssignedPodDelete事件,当做当前pod删除。因为node资源释放,可以重新评估其他unschedulablePods,可能可以重新回队列处理;

  3. FailureHandler,当前pod进入backoffQ/unschedulablePods;

func(sched *Scheduler)scheduleOnePod(...) {// 串行 调度周期 scheduleResult, assumedPodInfo, status := sched.schedulingCycle(...)if !status.IsSuccess() {  sched.FailureHandler(schedulingCycleCtx, fwk, assumedPodInfo, status, scheduleResult.nominatingInfo, start)return }// 并行 绑定周期go sched.runBindingCycle(ctx, state, fwk, scheduleResult, assumedPodInfo, start, podsToActivate)}func(sched *Scheduler)runBindingCycle(...) {// 执行 绑定周期 status := sched.bindingCycle(bindingCycleCtx, state, schedFramework, scheduleResult, assumedPodInfo, start, podsToActivate)if !status.IsSuccess() {// 处理失败  sched.handleBindingCycleError(bindingCycleCtx, state, schedFramework, assumedPodInfo, start, scheduleResult, status)return }}func(sched *Scheduler)handleBindingCycleError(...) { logger := klog.FromContext(ctx) assumedPod := podInfo.Pod// 1. unreserve & 清理内存node中的podif forgetErr := sched.unreserveAndForget(ctx, state, fwk, podInfo, scheduleResult.SuggestedHost); forgetErr != nil {  utilruntime.HandleErrorWithContext(ctx, forgetErr, "ForgetPod failed") } else {// 2. 发布EventAssignedPodDelete事件// 其他pod可能因为当前pod删除而重新有资格调度,重新进入activeQ或backoffQif status.IsRejected() {defer sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(logger, framework.EventAssignedPodDelete, assumedPod, nilfunc(pod *v1.Pod)bool {return assumedPod.UID != pod.UID   })  } else {   sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(logger, framework.EventAssignedPodDelete, assumedPod, nilnil)  } }// 3. 当前pod进入backoffQ/unschedulablePods sched.FailureHandler(ctx, fwk, podInfo, status, clearNominatedNode, start)}

scheduler_one.go:绑定周期

  1. RunPreBindPreFlights:控制哪些插件可以 并行PreBind or 跳过PreBind;

  2. WaitOnPermit:调度周期中PermitPlugin可以设置等待时长,如果超时未批准,这里返回FitError;

  3. RunPreBindPlugins:执行PreBind,比如Pod的PVC绑定节点;

  4. bind:执行BindPlugin,默认只有DefaultBinder,调用apiserver,设置pod.spec.NodeName;

  5. RunPostBindPlugins:PostBind无内置实现;

func(sched *Scheduler)bindingCycle(...) *fwk.Status { logger := klog.FromContext(ctx) assumedPod := assumedPodInfo.Pod// 1. PreBindPreFlightvar preFlightStatus *fwk.Statusif sched.nominatedNodeNameForExpectationEnabled {  preFlightStatus = schedFramework.RunPreBindPreFlights(ctx, state, assumedPod, scheduleResult.SuggestedHost)if preFlightStatus.Code() == fwk.Error ||   preFlightStatus.IsRejected() {return preFlightStatus  }if preFlightStatus.IsSuccess() || schedFramework.WillWaitOnPermit(ctx, assumedPod) {// 如果preFlight成功,但是需要permit,更新NominatedNodeName=nodeNameif err := updatePod(ctx, sched.client, schedFramework.APICacher(), assumedPod, nil, &fwk.NominatingInfo{    NominatedNodeName: scheduleResult.SuggestedHost,    NominatingMode:    fwk.ModeOverride,   }); err != nil {    logger.Error()   }  } }// 2. 等Permit插件放行if status := schedFramework.WaitOnPermit(ctx, assumedPod); !status.IsSuccess() {if status.IsRejected() {   fitErr := &framework.FitError{}   fitErr.Diagnosis.NodeToStatus.Set(scheduleResult.SuggestedHost, status)return fwk.NewStatus(status.Code()).WithError(fitErr)  }return status } sched.SchedulingQueue.Done(assumedPod.UID)// 3. PreBind --- pvc pvif status := schedFramework.RunPreBindPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost); !status.IsSuccess() {return status }// 4. DefaultBinder调用apiserver pod/binding 设置pod.spec.NodeName=if status := sched.bind(ctx, schedFramework, assumedPod, scheduleResult.SuggestedHost, state); !status.IsSuccess() {return status }// 5. PostBind 无内置实现 schedFramework.RunPostBindPlugins(ctx, state, assumedPod, scheduleResult.SuggestedHost)returnnil}

PreBind

VolumeBinding插件为例,如果Pod的Volume引用PVC,需要等PVC和PV绑定完成,再让Pod与Node绑定。

PreBind需要把Reserve阶段内存assume的PV/PVC变更写回apiserver,并阻塞等待卷真正绑定完成,再让 Pod 继续绑节点。

volume_binding.go:调度周期Filter阶段,VolumeBinding根据node筛选出合适的pvc绑定方案

  1. 静态绑定StaticBindings:与PVC匹配的PV已经存在,只需要建立绑定关系;

  2. 动态绑定DynamicProvisions:PV不存在,需要更新PVC,让provisioner(存储插件)感知到后创建PV并绑定;

func(pl *VolumeBinding)Filter(ctx context.Context, cs fwk.CycleState, pod *v1.Pod, nodeInfo fwk.NodeInfo) *fwk.Status { logger := klog.FromContext(ctx) node := nodeInfo.Node() state, err := getStateData(cs)// 根据node找出合适的绑定关系 podVolumes, reasons, err := pl.Binder.FindPodVolumes(logger, pod, state.podVolumeClaims, node)if err != nil {return fwk.AsStatus(err) }// 如果找不到,返回UnschedulableAndUnresolvableiflen(reasons) > 0 {  status := fwk.NewStatus(fwk.UnschedulableAndUnresolvable)return status } state.Lock()// 把绑定关系存储到状态 state.podVolumesByNode[node.Name] = podVolumes state.Unlock()returnnil}type stateData struct {// node -> pvc绑定关系 podVolumesByNode map[string]*PodVolumes}type PodVolumes struct {// 静态绑定,与PVC匹配的PV已经存在,只需要建立绑定关系 StaticBindings []*BindingInfo// 动态绑定,PV不存在,更新PVC让provisioner创建PV与之绑定 DynamicProvisions []*DynamicProvision}type DynamicProvision struct { PVC          *v1.PersistentVolumeClaim}type BindingInfo struct { pvc *v1.PersistentVolumeClaim pv *v1.PersistentVolume}

volume_binding.go:调度周期Reserve,取出绑定方案,更新内存中PV和PVC

  1. 静态绑定,建立PV和PVC的关系,pv.spec.claimRef = pvc;

  2. 动态绑定,建立PVC和Node的关系,pvc注解volume.kubernetes.io/selected-node = nodeName;

func(pl *VolumeBinding)Reserve(ctx context.Context, cs fwk.CycleState, pod *v1.Pod, nodeName string) *fwk.Status { state, err := getStateData(cs)// 取出绑定方案 podVolumes, ok := state.podVolumesByNode[nodeName]if ok {  allBound, err := pl.Binder.AssumePodVolumes(klog.FromContext(ctx), pod, nodeName, podVolumes)if err != nil {return fwk.AsStatus(err)  }  state.allBound = allBound } else {  state.allBound = true }returnnil}func(b *volumeBinder)AssumePodVolumes(logger klog.Logger, assumedPod *v1.Pod, nodeName string, podVolumes *PodVolumes)(allFullyBound bool, err error) {// PVC 注解 pv.kubernetes.io/bind-completed = yesif allBound := b.arePodVolumesBound(logger, assumedPod); allBound {returntruenil }// 静态绑定 pv.spec.claimRef = pvc newBindings := []*BindingInfo{}for _, binding := range podVolumes.StaticBindings {  newPV, dirty, err := volume.GetBindVolumeToClaim(binding.pv, binding.pvc)// ...        err = b.pvCache.Assume(newPV)  newBindings = append(newBindings, &BindingInfo{pv: newPV, pvc: binding.pvc}) }// 动态绑定 更新PVC注解 volume.kubernetes.io/selected-node=nodeName,等provisioner创建PV绑定 newProvisionedPVCs := []*DynamicProvision{}for _, dynamicProvision := range podVolumes.DynamicProvisions {  claimClone := dynamicProvision.PVC.DeepCopy()  metav1.SetMetaDataAnnotation(&claimClone.ObjectMeta, volume.AnnSelectedNode, nodeName)  err = b.pvcCache.Assume(claimClone)  newProvisionedPVCs = append(newProvisionedPVCs, &DynamicProvision{PVC: claimClone}) } podVolumes.StaticBindings = newBindings podVolumes.DynamicProvisions = newProvisionedPVCsreturn}

volume_binding.go:绑定周期-PreBindPreFlight

  1. 返回PreBindPreFlightResult,AllowParallel=true允许PreBind阶段与其他插件并行执行;

  2. Skip场景:Reserve阶段计算,如果所有Volume对应的PVC的注解pv.kubernetes.io/bind-completed存在,代表PV已经与PVC绑定,allBound=true,跳过PreBind处理;

  3. 其他场景,需要走PreBind,绑定PV和PVC;

func(pl *VolumeBinding)PreBindPreFlight(ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string)(*fwk.PreBindPreFlightResult, *fwk.Status) {// PreBind允许和其他插件并行 result := &fwk.PreBindPreFlightResult{AllowParallel: true}// state中存储本轮调度当前插件的状态 s, err := getStateData(state)if err != nil {return result, fwk.AsStatus(err) }// pod对应的pvc已经绑定好pv,跳过PreBindif s.allBound {return result, fwk.NewStatus(fwk.Skip) }return result, nil}type stateData struct {// Reserve阶段计算,是否所有PVC都绑定成功 allBound bool}

volume_binding.go:绑定周期-PreBind,取出Filter阶段的绑定方案,持久化到apiserver,等待PVC完成绑定(注解pv.kubernetes.io/bind-completed=yes)。

func(pl *VolumeBinding)PreBind(ctx context.Context, cs fwk.CycleState, pod *v1.Pod, nodeName string) *fwk.Status { s, err := getStateData(cs)if err != nil {return fwk.AsStatus(err) }if s.allBound {returnnil }// pod在这个node上的绑定方案 podVolumes, ok := s.podVolumesByNode[nodeName]// 调用apiserver执行绑定 err = pl.Binder.BindPodVolumes(ctx, pod, podVolumes)if err != nil {return fwk.AsStatus(err) }returnnil}func(b *volumeBinder)BindPodVolumes(ctx context.Context, assumedPod *v1.Pod, podVolumes *PodVolumes)(err error) { bindings := podVolumes.StaticBindings claimsToProvision := convertDynamicProvisionsToPVCs(podVolumes.DynamicProvisions)// 调用apiserver,持久化绑定关系 err = b.bindAPIUpdate(ctx, assumedPod, bindings, claimsToProvision)if err != nil {return err }// 等待PVC注解pv.kubernetes.io/bind-completed=yes err = wait.PollUntilContextTimeout(ctx, time.Second, b.bindTimeout, falsefunc(ctx context.Context)(bool, error) {  b, err := b.checkBindings(logger, assumedPod, bindings, claimsToProvision)return b, err })if err != nil {return fmt.Errorf("binding volumes: %w", err) }returnnil}func(b *volumeBinder)bindAPIUpdate(...)error { podName := getPodName(pod)// 静态绑定 更新pv引用pvcfor _, binding = range bindings {  newPV, err := b.kubeClient.CoreV1().PersistentVolumes().Update(ctx, binding.pv, metav1.UpdateOptions{}) }// 动态绑定 pvc关联nodefor i, claim = range claimsToProvision {  newClaim, err := b.kubeClient.CoreV1().PersistentVolumeClaims(claim.Namespace).Update(ctx, claim, metav1.UpdateOptions{}) }returnnil}

Bind

framework.go:RunBindPlugins,只要有一个BindPlugin返回非Skip,结束流程。

func(f *frameworkImpl)RunBindPlugins(ctx context.Context, state fwk.CycleState, pod *v1.Pod, nodeName string)(status *fwk.Status) {for _, pl := range f.bindPlugins {  status = f.runBindPlugin(ctx, pl, state, pod, nodeName)if status.IsSkip() {continue  }if !status.IsSuccess() {if status.IsRejected() {    status.SetPlugin(pl.Name())return status   }   err := status.AsError()return fwk.AsStatus(fmt.Errorf("running Bind plugin %q: %w", pl.Name(), err))  }return status }return status}

default_binder.go:DefaultBinder是默认BindPlugin,调用apiserver /api/v1/namespaces/{namespace}/pods/{podName}/binding,执行绑定。

funcNew(_ context.Context, _ runtime.Object, handle fwk.Handle)(fwk.Plugin, error) {// 插件构造,可以拿到Handle,Handle有一些工具和数据可以给插件运行时用return &DefaultBinder{handle: handle}, nil}func(b DefaultBinder)Bind(ctx context.Context, state fwk.CycleState, p *v1.Pod, nodeName string) *fwk.Status { binding := &v1.Binding{// pod  ObjectMeta: metav1.ObjectMeta{Namespace: p.Namespace, Name: p.Name, UID: p.UID},// node  Target:     v1.ObjectReference{Kind: "Node", Name: nodeName}, } err := b.handle.ClientSet().CoreV1().Pods(binding.Namespace).Bind(ctx, binding, metav1.CreateOptions{})if err != nil {return fwk.AsStatus(err) }returnnil}

pkg/registry/core/pod/storage/storage.go:apiserver的binding接口,修改pod的nodeName,至此全流程结束。

func(r *BindingREST)Create(...)(out runtime.Object, err error) { binding, ok := obj.(*api.Binding) err = r.assignPod(...)return}func(r *BindingREST)assignPod(...)(err error) {if _, err = r.setPodNodeAndMetadata(...); err != nil { }return}func(r *BindingREST)setPodNodeAndMetadata(...)(finalPod *api.Pod, err error) {// 乐观更新 err = r.store.Storage.GuaranteedUpdate(...,            storage.SimpleUpdate(func(obj runtime.Object)(runtime.Object, error) {  pod, ok := obj.(*api.Pod)// machine=scheduler选的节点名  pod.Spec.NodeName = machine  finalPod = podreturn pod, nil }), dryRun, nil)return finalPod, err}

总结

kube-scheduler的核心路径:

Informer 事件 → 调度队列 → ScheduleOne(调度周期 → 绑定周期)→ API Server 写回 nodeName

  1. 入口:Informer 监听 Pod/Node 等变更;未绑定且由本 scheduler 负责的 Pod 进入调度队列。

  2. 队列activeQ 出队调度,backoffQ 做失败退避,unschedulablePods 等待事件唤醒;PreEnqueue / QueueingHint / QueueSort 分别控制入队门槛、重入时机和出队顺序。

  3. 调度周期(串行):快照 Node 列表 → Filter(硬过滤)→ Score(打分选优)→ 失败可走 PostFilter 抢占 → Assume / Reserve / Permit 在内存中预占。

  4. 绑定周期(异步):WaitOnPermit → PreBind → Bind(DefaultBinder 调 Binding API),最终由 API Server 把 pod.spec.nodeName 持久化。

参考资料:

  • 调度框架:https://kubernetes.io/docs/concepts/scheduling-eviction/scheduling-framework/

  • kube-scheduler配置:https://kubernetes.io/docs/reference/scheduling/config/

  • out-of-tree扩展kube-scheduler的案例:https://github.com/kubernetes-sigs/scheduler-plugins