前言
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.go:NewInformerFactory 构造 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.go:addAllEventHandlers 注册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.go:addPod 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.go:updatePod 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.go(PriorityQueue)。 配套三个队列:active_queue.go、backoff_queue.go、unschedulable_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:
activeQ:待调度的pod,按 QueueSort 排序
backoffQ:调度失败后退避中,到期进 activeQ
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}核心流转:
新pod:PreEnqueue通过,进入activeQ;否则进入unschedulablePods;
ScheduleOne调度循环:从activeQ弹出pod,如果失败,根据情况进入activeQ/backoffQ/unschedulablePods;
集群事件:EnqueueExtensions订阅事件,根据QueueingHint返回决定是否入队;
退避结束 / 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, 0, len(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)}| 单协程串行 | 异步并发go runBindingCycle) | |
schedulingCycle | bindingCycle | |
FailureHandler | unreserveAndForgetFailureHandler |
整体流程如下:
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:调度周期主流程。
UpdateSnapshot:将cache中的node列表保存到Scheduler自己的快照里,本次调度周期中以这个node列表快照为准;
schedulingAlgorithm:Filter+Score计算pod应该调度到哪个Node上,如果失败,则返回上层进入backoffQ/unschedulablePods;
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
findNodesThatFitPod:执行PreFilter+Filter插件,从node列表中过滤出部分node,如果node数量为1,直接选中,不走后续流程;
RunPostFilterPlugins:如果过滤结果为空,执行PostFilter,可能得到抢占节点nominatingInfo;
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.go:findNodesThatFitPod先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, 0, len(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范围,根据插件返回状态
Skip:后面跳过该插件Filter;
UnschedulableAndUnresolvable:返回上层,PostFilter默认抢占逻辑不生效;
Unschedulable:返回上层,PostFilter默认可抢占;
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 {returnnil, nil// 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)→ 无法缩集,全量节点returnnil, nil } 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 }returnnil, nil}Filter
scheduler_one.go:findNodesThatPassFilters并行对每个节点执行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.go:RunFilterPluginsWithNominatedPods对单个节点调用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 {returnnil, 0, fwk.NewStatus(fwk.UnschedulableAndUnresolvable) }// 2. 如果把所有潜在受害者都移走后,新 Pod 还是放不进来,说明这个节点不适合通过抢占解决。if status := pl.fh.RunFilterPluginsWithNominatedPods(ctx, state, pod, nodeInfo); !status.IsSuccess() {returnnil, 0, 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 {returnnil, 0, 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:
Assume:内存更新,pod绑定nodeName,pod挂到node上,更新node的资源使用情况;
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 }returnnil, nil}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:绑定周期并行执行,如果处理失败
unreserveAndForget:调用ReservePlugin的Unreserve方法,清理内存node中的pod;
发送EventAssignedPodDelete事件,当做当前pod删除。因为node资源释放,可以重新评估其他unschedulablePods,可能可以重新回队列处理;
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, nil, func(pod *v1.Pod)bool {return assumedPod.UID != pod.UID }) } else { sched.SchedulingQueue.MoveAllToActiveOrBackoffQueue(logger, framework.EventAssignedPodDelete, assumedPod, nil, nil) } }// 3. 当前pod进入backoffQ/unschedulablePods sched.FailureHandler(ctx, fwk, podInfo, status, clearNominatedNode, start)}scheduler_one.go:绑定周期
RunPreBindPreFlights:控制哪些插件可以 并行PreBind or 跳过PreBind;
WaitOnPermit:调度周期中PermitPlugin可以设置等待时长,如果超时未批准,这里返回FitError;
RunPreBindPlugins:执行PreBind,比如Pod的PVC绑定节点;
bind:执行BindPlugin,默认只有DefaultBinder,调用apiserver,设置pod.spec.NodeName;
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绑定方案
静态绑定StaticBindings:与PVC匹配的PV已经存在,只需要建立绑定关系;
动态绑定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
静态绑定,建立PV和PVC的关系,pv.spec.claimRef = pvc;
动态绑定,建立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 {returntrue, nil }// 静态绑定 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
返回PreBindPreFlightResult,AllowParallel=true允许PreBind阶段与其他插件并行执行;
Skip场景:Reserve阶段计算,如果所有Volume对应的PVC的注解pv.kubernetes.io/bind-completed存在,代表PV已经与PVC绑定,allBound=true,跳过PreBind处理;
其他场景,需要走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, false, func(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
入口:Informer 监听 Pod/Node 等变更;未绑定且由本 scheduler 负责的 Pod 进入调度队列。
队列:
activeQ出队调度,backoffQ做失败退避,unschedulablePods等待事件唤醒;PreEnqueue / QueueingHint / QueueSort 分别控制入队门槛、重入时机和出队顺序。调度周期(串行):快照 Node 列表 → Filter(硬过滤)→ Score(打分选优)→ 失败可走 PostFilter 抢占 → Assume / Reserve / Permit 在内存中预占。
绑定周期(异步):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
夜雨聆风