一、开篇:一个真实的生产事故
去年双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)。
| 阶段 | 负责内容 | 扩展点 |
|---|---|---|
| QueueSort | Pod入队排序(默认按优先级) | 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.go 的 ScheduleOne 方法。我们直接看最核心的流程。
三、源码深度拆解:调度周期
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.2s | 0.32s | 0.11s (per Pod) |
| P99调度延迟 | 4.0s | 0.8s | 0.25s |
| 调度吞吐量 (pods/s) | 8 | 300 | 800 |
| 调度器CPU使用 | 3.2 cores | 1.8 cores | 2.5 cores |
| 调度器内存使用 | 4.5GB | 1.2GB | 2.0GB |
数据证明:自定义调度器将单Pod调度延迟降低了73%,吞吐量提升了37倍(批量模式)。当然,代价是牺牲了部分调度精度(没有考虑节点亲和性、污点容忍等)。如果业务对这些特性依赖很强,建议在Filter中按需添加,不要全关。
七、避坑指南
这4个坑是我们踩过之后才知道的:
- 抢占与优先级死循环:自定义调度器没有实现PostFilter(抢占),如果Pod请求资源太大且没有节点满足,Pod会一直待在unschedulableQ中反复重试,最后backoff到10秒间隔。我们后来加了一个超时机制:超过5次调度失败直接标记为Unschedulable,并发送告警。
- 节点资源统计遗漏:在FastMatch中,我们只统计了容器的Requests,忽略了InitContainers、EphemeralContainers,导致调度后InitContainer启动时OOM。后来改为统计所有容器类型。
- 批量调度导致Pod抢占排序错乱:批量Pop10个Pod,然后顺序绑定,如果对应节点恰好被其他Pod占用,会绑定失败。解决办法是绑定前二次校验资源,失败则重新入队列。
- 调度器版本兼容:不同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— 框架接口定义
有了这些底子,你也能写出适配自己业务的自定义调度器。
```