K8s调度器源码拆解:从队列到绑定
发布日期: 2026/07/30 阅读总量: 0
```

一、开篇:一个真实的生产事故

去年双11大促前一周,我们线上K8s集群从200节点扩容到5000节点,开始压测时Pod一直处于Pending状态,最长的等了30分钟才调度上。查看kube-scheduler日志,发现scheduling queue中积压了8000+个Pod,单个Pod调度周期平均1.2秒。如果放任不管,大促当天集群必然雪崩。

问题出在哪?——调度器成了瓶颈。于是我翻开kube-scheduler源码(v1.28.2),从根上搞懂调度器的工作机制,并手写一个性能更好的自定义调度器。本文全程实战,无论你用哪个版本,核心思想不变。

二、调度器核心架构:三个角色、四个阶段

kube-scheduler本质是一个Event Loop,从队列中拿Pod,经过筛选、打分、绑定后返回成功/失败。官方将这个过程抽象成Scheduling Framework,提供扩展点(Extender/Plugin)。

阶段负责内容扩展点
QueueSortPod入队排序(默认按优先级)QueueSortPlugin
PreFilter预筛选检查(如Pod是否满足节点选择器)PreFilterPlugin
Filter节点筛选(资源、污点、亲和性等)FilterPlugin
PostFilter筛选后处理(如抢占)PostFilterPlugin
PreScore为打分做预计算PreScorePlugin
Score节点打分(0-100)ScorePlugin
NormalizeScore分数归一化NormalizeScorePlugin
Reserve预留资源ReservePlugin
Permit授权(可阻塞等待)PermitPlugin
PreBind绑定前处理(如卷挂载)PreBindPlugin
Bind将Pod绑定到节点BindPlugin
PostBind绑定后清理PostBindPlugin

调度器就是按顺序执行这些插件。每次调度一个Pod,都称为一个调度周期(ScheduleOne)。

源码入口在 pkg/scheduler/scheduler.goScheduleOne 方法。我们直接看最核心的流程。

三、源码深度拆解:调度周期

3.1 调度队列(Priority Queue)

Pod创建后并不会立刻被调度,而是先进入一个优先级队列。队列默认按Pod的priority排序,同优先级则按QueuedTime排序。源码在 pkg/scheduler/internal/queue/priority_queue.go

核心结构体:

// 精简版 PriorityQueue
type PriorityQueue struct {
    // 使用双向链表+分桶(activeQ, backoffQ, unschedulableQ)
    activeQ     *heap.Heap
    backoffQ    *heap.Heap
    unschedulableQ map[string]*framework.QueuedPodInfo
    // 每个Pod的backoff时间
    backoffDuration time.Duration
    // 最大backoff(默认10秒)
    maxBackoffDuration time.Duration
}

Pod被放入activeQ,然后通过Pop取出。注意:Pop一次只取一个,且会锁定Queue。

// Pop 函数逻辑
func (p *PriorityQueue) Pop() (*framework.QueuedPodInfo, error) {
    p.lock.Lock()
    defer p.lock.Unlock()
    for p.activeQ.Len() == 0 {
        // 如果active为空,等待信号
        p.cond.Wait()
    }
    obj, err := p.activeQ.Pop()
    if err != nil {
        return nil, err
    }
    pInfo := obj.(*framework.QueuedPodInfo)
    pInfo.Attempts++  // 调度尝试次数+1
    return pInfo, nil
}

问题暴露:当activeQ中Pod数量巨大时,Pop操作本身开销极小(O(log n)),但后续的Filter/Score才是大头。而如果你的Pod优先级都一样,队列退化为先进先出,没有任何区分度。

3.2 筛选(Filter)与打分(Score)

从队列拿到Pod后,调度器会调用 findNodesThatFitPod 做筛选。源码在 pkg/scheduler/core/generic_scheduler.go

// 简化版调度周期
func (sched *Scheduler) scheduleOne(ctx context.Context) {
    podInfo := sched.NextPod() // 从队列Pop
    pod := podInfo.Pod
    
    // 执行调度插件
    scheduleResult, err := sched.Algorithm.Schedule(ctx, sched.Profiles[pod.Spec.SchedulerName], sched.Cache, pod)
    if err != nil {
        // 调度失败,重新入队列(带backoff)
        sched.FailureHandler(ctx, sched.Profiles[pod.Spec.SchedulerName], podInfo, err)
        return
    }
    // 绑定
    err = sched.bind(ctx, pod, scheduleResult.SuggestedHost)
    if err != nil {
        // 绑定失败处理
        ...
    }
}

筛选逻辑:遍历所有节点,对每个节点执行FilterPlugin,只有通过所有Filter的节点才进入Score阶段。

默认Filter包括:NodeResourcesFit(CPU/内存)、NodePorts、NodeAffinity、TaintToleration等。

节点筛选在多核下可以并行,但默认每次调度都要重新计算所有节点。官方提供了一个 NodeCache 缓存部分结果,但仍然不够快。

打分阶段对通过筛选的节点依次执行ScorePlugin,最后乘以权重求和,取分数最高的节点。所有打分插件的总分范围0-100,但不同插件权重可以不同。

3.3 抢占(Preemption)

当没有节点满足Pod资源请求时,调度器尝试抢占——驱逐低优先级Pod来腾出空间。源码在 pkg/scheduler/core/preemption.go。核心函数 preempt

func (sched *Scheduler) preempt(ctx context.Context, pod *v1.Pod, nodeName string) (string, error) {
    // 1. 选出所有可能被抢占的节点
    node := sched.Cache.Node(nodeName)
    // 2. 计算出需要释放的资源量
    needed := calculateResourceDelta(pod, node)
    // 3. 从该节点上的低优先级Pod中选出要移除的Pod
    victims := selectVictims(node, needed, pod)
    // 4. 模拟删除victims后的节点状态
    if fitsAfterRemoval(node, victims, pod) {
        // 5. 实际删除(标记为抢占,异步清理)
        for _, v := range victims {
            sched.Cache.RemovePod(v)
        }
        return nodeName, nil
    }
    return "", nil // 抢占失败
}

抢占带来的风险:低优先级Pod被驱逐可能造成级联故障,而且抢占本身需要重新计算节点状态,增加延迟。

四、方案对比:默认调度器 vs 自定义调度器

我的需求:在5000节点+秒级Pod创建量下,将单Pod调度延迟从1.2秒降到0.5秒以内。默认调度器无法满足。

方案一:调优默认调度器
- 增加kube-scheduler的并发数:--kube-api-qps=100 --kube-api-burst=150
- 调整Filter并发:--concurrent-scheduling-workers=50(默认16)
- 开启NodeCache
效果:调度延迟从1.2s降至0.9s,但内存占用飙到8GB,且Filter的并行度受节点数限制,瓶颈在APIServer。

方案二:手写自定义调度器
- 只监听我们自己关心的Pod(通过label选择)
- 跳过复杂的Filter,使用简化版的资源计算(只算CPU/Mem)
- 打分不再遍历所有节点,而是维护一个按资源排序的节点索引
- 支持批量调度(一次Pop多个Pod,按批打分)

我们选择了方案二。以下是自定义调度器的核心代码。

五、自定义调度器完整实现

使用Go语言,基于 sigs.k8s.io/scheduler-plugins 框架(v0.28.0)。注:这不是玩具,我们实际在生成环境跑了一年。

5.1 调度器主入口

// cmd/scheduler/main.go
package main

import (
    "flag"
    "k8s.io/component-base/logs"
    "k8s.io/kubernetes/cmd/kube-scheduler/app"
    // 引入自定义插件
    "example.com/scheduler-plugins/pkg/fastmatch"
    "example.com/scheduler-plugins/pkg/fastscore"
)

func main() {
    logs.InitLogs()
    defer logs.FlushLogs()
    
    command := app.NewSchedulerCommand(
        app.WithPlugin(fastmatch.Name, fastmatch.New),
        app.WithPlugin(fastscore.Name, fastscore.New),
        // 禁用默认的Filter和Score插件(可选)
        // app.WithPlugin(..., app.WithDefaultPluginWeight(...)),
    )
    if err := command.Execute(); err != nil {
        os.Exit(1)
    }
}

5.2 快速匹配插件(FastMatch)

// pkg/fastmatch/fastmatch.go
package fastmatch

import (
    "context"
    "k8s.io/kubernetes/pkg/scheduler/framework"
    v1 "k8s.io/api/core/v1"
)

const Name = "FastMatch"

type FastMatch struct {
    handle framework.Handle
}

func New(obj runtime.Object, handle framework.Handle) (framework.Plugin, error) {
    return &FastMatch{handle: handle}, nil
}

func (f *FastMatch) Name() string { return Name }

// Filter 只检查Pod需要的CPU和内存是否满足
func (f *FastMatch) Filter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeInfo *framework.NodeInfo) *framework.Status {
    // 获得Pod请求资源
    podsReqs := computePodResourceRequest(pod)
    // 节点可分配资源
    nodeAlloc := nodeInfo.AllocatableResource()
    // 节点已分配资源(包括正在调度的Pod)
    nodeUsed := nodeInfo.RequestedResource()
    // 是否够
    for k, v := range podsReqs {
        if nodeAlloc.ScalarResources[k]-nodeUsed.ScalarResources[k] < v {
            return framework.NewStatus(framework.Unschedulable, "insufficient resources")
        }
    }
    return framework.NewStatus(framework.Success, "")
}

func computePodResourceRequest(pod *v1.Pod) map[v1.ResourceName]int64 {
    req := map[v1.ResourceName]int64{
        v1.ResourceCPU:    0,
        v1.ResourceMemory: 0,
    }
    for _, c := range pod.Spec.Containers {
        req[v1.ResourceCPU] += c.Resources.Requests.Cpu().MilliValue()
        req[v1.ResourceMemory] += c.Resources.Requests.Memory().Value()
    }
    return req
}

5.3 快速打分插件(FastScore)

// pkg/fastscore/fastscore.go
package fastscore

import (
    "context"
    "math"
    "k8s.io/kubernetes/pkg/scheduler/framework"
    v1 "k8s.io/api/core/v1"
)

const Name = "FastScore"

type FastScore struct {
    handle framework.Handle
}

func New(obj runtime.Object, handle framework.Handle) (framework.Plugin, error) {
    return &FastScore{handle: handle}, nil
}

func (f *FastScore) Name() string { return Name }

// Score 根据节点剩余资源比例打分,剩余越多分越高
func (f *FastScore) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
    nodeInfo, err := f.handle.SnapshotSharedLister().NodeInfos().Get(nodeName)
    if err != nil {
        return 0, framework.NewStatus(framework.Error, "node not found")
    }
    alloc := nodeInfo.AllocatableResource()
    used := nodeInfo.RequestedResource()
    // CPU剩余比例
    cpuRatio := float64(alloc.ScalarResources[v1.ResourceCPU]-used.ScalarResources[v1.ResourceCPU]) / float64(alloc.ScalarResources[v1.ResourceCPU])
    memRatio := float64(alloc.ScalarResources[v1.ResourceMemory]-used.ScalarResources[v1.ResourceMemory]) / float64(alloc.ScalarResources[v1.ResourceMemory])
    // 取最小值作为分数(0-100)
    score := int64(math.Min(cpuRatio, memRatio) * 100)
    return score, framework.NewStatus(framework.Success, "")
}

5.4 调度器配置文件

# scheduler-config.yaml
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
clientConnection:
  kubeconfig: /etc/kubernetes/scheduler.conf
profiles:
  - schedulerName: fast-scheduler
    plugins:
      filter:
        enabled:
          - name: FastMatch    # 只保留自己
        disabled:
          - name: "*"          # 禁用所有默认
      score:
        enabled:
          - name: FastScore
        disabled:
          - name: "*"
      preFilter:
        disabled:
          - name: "*"
      postFilter:
        disabled:
          - name: "*"
      reserve:
        disabled:
          - name: "*"
      preBind:
        disabled:
          - name: "*"
      bind:
        enabled:
          - name: DefaultBinder   # 必须保留默认绑定器
---

5.5 Pod使用自定义调度器

apiVersion: v1
kind: Pod
metadata:
  name: my-pod
spec:
  schedulerName: fast-scheduler   # 指定自定义调度器
  containers:
    - name: nginx
      image: nginx
      resources:
        requests:
          cpu: "1"
          memory: "1Gi"

5.6 批量调度(可选优化)

为了减少锁竞争和网络延迟,我们实现了批量调度:每次Pop出M个Pod,然后一起筛选、打分。在调度器内部维护一个缓冲队列:

// 伪代码:批量调度逻辑
func (s *FastScheduler) ScheduleBulk(ctx context.Context, batchSize int) {
    pods := make([]*v1.Pod, 0, batchSize)
    for i := 0; i < batchSize; i++ {
        pod, err := s.podQueue.Pop()
        if err != nil {
            break
        }
        pods = append(pods, pod)
    }
    if len(pods) == 0 {
        return
    }
    // 对所有Pod进行筛选(多节点并行)
    feasibleNodes := s.batchFilter(pods)
    // 对每个Pod+节点组合打分
    assignments := s.batchScore(pods, feasibleNodes)
    // 并发绑定
    s.batchBind(assignments)
}

注意:批量调度需要开启Permit阶段与节点预占,防止多个Pod竞争同一个节点。详细实现略。

六、效果数据

测试环境:K8s v1.28.2,集群节点5000个(8C16G),Pod规格average: 1C1Gi,调度器部署在独立的8C16G机器上。

指标默认调度器自定义调度器(单Pod)自定义调度器(批量10个)
平均调度延迟1.2s0.32s0.11s (per Pod)
P99调度延迟4.0s0.8s0.25s
调度吞吐量 (pods/s)8300800
调度器CPU使用3.2 cores1.8 cores2.5 cores
调度器内存使用4.5GB1.2GB2.0GB

数据证明:自定义调度器将单Pod调度延迟降低了73%,吞吐量提升了37倍(批量模式)。当然,代价是牺牲了部分调度精度(没有考虑节点亲和性、污点容忍等)。如果业务对这些特性依赖很强,建议在Filter中按需添加,不要全关。

七、避坑指南

这4个坑是我们踩过之后才知道的:

  1. 抢占与优先级死循环:自定义调度器没有实现PostFilter(抢占),如果Pod请求资源太大且没有节点满足,Pod会一直待在unschedulableQ中反复重试,最后backoff到10秒间隔。我们后来加了一个超时机制:超过5次调度失败直接标记为Unschedulable,并发送告警。
  2. 节点资源统计遗漏:在FastMatch中,我们只统计了容器的Requests,忽略了InitContainers、EphemeralContainers,导致调度后InitContainer启动时OOM。后来改为统计所有容器类型。
  3. 批量调度导致Pod抢占排序错乱:批量Pop10个Pod,然后顺序绑定,如果对应节点恰好被其他Pod占用,会绑定失败。解决办法是绑定前二次校验资源,失败则重新入队列。
  4. 调度器版本兼容:不同K8s版本的Framework API不兼容。我们一开始基于v1.27开发的自定义插件,升级到v1.28时接口变了,编译不过。强烈建议锁定K8s版本,并用vendor或go模块精确控制。

八、总结与建议

调度器源码不难,难在理解整个调度框架与扩展机制。如果你不需要万级节点吞吐,完全没必要手写——调整默认调度器的并发参数足够了。但如果你想实现抢占式调度、批量调度、甚至基于机器学习的调度,那就必须深入源码,拆解每个阶段。

最后推荐几个必读源码文件:

  • pkg/scheduler/scheduler.go — 主循环
  • pkg/scheduler/core/generic_scheduler.go — 调度算法核心
  • pkg/scheduler/internal/queue/priority_queue.go — 优先级队列
  • pkg/scheduler/framework/interface.go — 框架接口定义

有了这些底子,你也能写出适配自己业务的自定义调度器。

```