Kubernetes Scheduling Framework 深度实战:从 Plugin 机制到 AI 训练任务多级抢占式调度

引言

Kubernetes 调度器是集群的"决策大脑",负责将 Pending 状态的 Pod 分配到合适的 Node 上运行。然而,默认调度器的通用逻辑无法满足 AI 训练任务、高性能计算、资源碎片治理等复杂场景的需求。

Kubernetes 1.15+ 引入的 Scheduling Framework 提供了一套可插拔的 Plugin 扩展机制,让开发者能够在不修改调度器源码的前提下,深度介入调度决策的每一个阶段。

本文将从调度框架的架构原理出发,深入剖析 Scheduling Cycle 和 Binding Cycle 的完整生命周期,然后通过一个 AI 训练任务多级抢占式调度 的完整实战案例,展示如何从零构建一个生产级调度插件。


一、调度框架架构全景

1.1 核心数据结构

调度框架的核心结构体是 framework.Framework,它内部维护了一组按阶段组织的 Plugin 链:

type Framework interface {
    QueueSortPlugin
    PreFilterPlugin
    FilterPlugin
    PostFilterPlugin
    PreScorePlugin
    ScorePlugin
    ReservePlugin
    PermitPlugin
    PreBindPlugin
    BindPlugin
    PostBindPlugin
}

每个 Plugin 接口对应调度流程的一个特定扩展点。调度器按顺序遍历这些扩展点上的插件,逐步缩小候选节点范围,最终确定最优分配。

1.2 两大调度周期

Scheduling Cycle(调度周期) 是单节点决策流程,它决定 Pod 应该被调度到哪个 Node 上:

QueueSort → PreEnqueue → PreFilter → Filter → PostFilter → PreScore → Score → NormalizeScore → Reserve → Permit

Binding Cycle(绑定周期) 是实际的节点绑定流程,与 Scheduling Cycle 异步执行:

PreBind → Bind → PostBind

两张周期的关键解耦点在于 Permit Plugin:Permit 可以设置 Pod 的等待状态(Wait/Allow/Deny),当设置为 Wait 时,Binding Cycle 被延迟执行,从而实现跨 Pod 的协调调度。


二、Plugin 扩展点深度解析

2.1 PreFilter:前置校验与状态初始化

PreFilter 在过滤阶段之前执行,用于检查 Pod 是否满足调度的基本前提条件,同时初始化后续阶段需要的共享状态:

func (pl *GPUFragmentationPlugin) PreFilter(ctx context.Context, state *framework.CycleState, pod *v1.Pod) *framework.Status {
    // 计算 Pod 所需的 GPU 资源
    resourceList := computePodGPURequest(pod)

    // 将计算结果存入 CycleState,供 Filter 阶段使用
    state.Write(cycleStateKey, &StateData{
        requiredGPU: resourceList,
        gpuModel:    getGPUModel(pod),
    })

    // 如果 Pod 请求的 GPU 数量超过单节点上限,直接拒绝
    if resourceList.NvidiaGPU.Value() > maxGPUPerNode {
        return framework.NewStatus(framework.Unschedulable, "request exceeds max GPU per node")
    }
    return nil
}

PreFilter 返回 Unschedulable 时,调度器不会尝试其他节点;返回 Error 会导致整个调度周期终止。

2.2 Filter:节点过滤(一票否决)

Filter 是核心的节点淘汰阶段,插件返回 Success 表示节点通过过滤,Unschedulable 表示节点不满足条件但不终止遍历。

func (pl *GPUFragmentationPlugin) Filter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeInfo *framework.NodeInfo) *framework.Status {
    s, err := state.Read(cycleStateKey)
    if err != nil {
        return framework.AsStatus(err)
    }
    cycleState := s.(*StateData)

    // 检查节点是否有足够的可用 GPU
    freeGPU := nodeInfo.Allocatable.ScalarResources[v1.ResourceNvidiaGPU]
    if freeGPU < cycleState.requiredGPU.NvidiaGPU.Value() {
        return framework.NewStatus(framework.Unschedulable, "insufficient GPU")
    }

    // 检查 GPU 型号是否匹配
    nodeGPUModel := nodeInfo.Node().Labels["gpu-model"]
    if cycleState.gpuModel != "" && nodeGPUModel != cycleState.gpuModel {
        return framework.NewStatus(framework.Unschedulable, "GPU model mismatch")
    }

    return nil
}

2.3 PreScore & Score:节点评分

Score 阶段为所有通过 Filter 的节点评分(0-100),最终得分最高者胜出:

func (pl *GPUFragmentationPlugin) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
    nodeInfo, err := pl.nodeInfoLister.Get(nodeName)
    if err != nil {
        return 0, framework.AsStatus(err)
    }

    // 评分策略:优先选择能更高效利用 GPU 的节点(减少碎片)
    freeGPU := nodeInfo.Available.ScalarResources[v1.ResourceNvidiaGPU]
    totalGPU := nodeInfo.Allocatable.ScalarResources[v1.ResourceNvidiaGPU]
    requestedGPU := computePodGPURequest(pod).NvidiaGPU.Value()

    // 如果调度后剩余 GPU 为 0 或接近 0,给予最高分(无碎片)
    remainingAfterSchedule := freeGPU - requestedGPU
    if remainingAfterSchedule == 0 {
        return framework.MaxNodeScore, nil
    }

    // 剩余越少分数越高(Best-Fit 策略)
    fragmentationRatio := float64(remainingAfterSchedule) / float64(totalGPU)
    score := int64((1.0 - fragmentationRatio) * float64(framework.MaxNodeScore))

    return score, nil
}

func (pl *GPUFragmentationPlugin) ScoreExtensions() framework.ScoreExtensions {
    return pl
}

func (pl *GPUFragmentationPlugin) NormalizeScore(ctx context.Context, state *framework.CycleState, pod *v1.Pod, scores framework.NodeScoreList) *framework.Status {
    // 归一化分数到 0-MaxNodeScore
    var maxScore int64
    for _, score := range scores {
        if score.Score > maxScore {
            maxScore = score.Score
        }
    }
    if maxScore == 0 {
        return nil
    }
    for i := range scores {
        scores[i].Score = scores[i].Score * framework.MaxNodeScore / maxScore
    }
    return nil
}

2.4 Reserve:资源预留

Reserve 阶段在 Permit 返回 Allow 后执行,用于在 NodeInfo 中"预占"资源,防止并发调度时资源超分配:

func (pl *GPUFragmentationPlugin) Reserve(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) *framework.Status {
    nodeInfo, err := pl.nodeInfoLister.Get(nodeName)
    if err != nil {
        return framework.AsStatus(err)
    }

    // 在 Available 资源中扣减请求量
    nodeInfo.AddPod(pod)
    return nil
}

func (pl *GPUFragmentationPlugin) Unreserve(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) {
    // 释放预留资源(当 Permit Deny 或调度失败时回调)
    nodeInfo, err := pl.nodeInfoLister.Get(nodeName)
    if err != nil {
        klog.Errorf("Unreserve failed for node %s: %v", nodeName, err)
        return
    }
    nodeInfo.RemovePod(pod)
}

2.5 Permit:批准等待(跨 Pod 协调)

Permit 是最强大的扩展点之一,它允许插件延迟 Pod 的绑定,实现等待条件满足后继续:

func (pl *AIPriorityPlugin) Permit(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (*framework.Status, time.Duration) {
    priority := getPodAIPriority(pod)

    // 高优先级任务:立即调度
    if priority >= PriorityHigh {
        return framework.NewStatus(framework.Success), 0
    }

    // 低优先级任务:检查高优先级任务是否正在等待资源消耗
    if pl.highPriorityQueue.Len() > 0 && priority < PriorityMedium {
        // 等待 30 秒,让高优先级任务先调度
        return framework.NewStatus(framework.Wait), 30 * time.Second
    }

    return framework.NewStatus(framework.Success), 0
}

三、实战:AI 训练任务多级抢占式调度

3.1 需求分析

在 AI 训练场景中,我们需要实现以下调度策略:

  1. 多级优先级:P0(生产训练)> P1(实验训练)> P2(开发调试)
  2. 资源碎片治理:GPU 调度采用 Best-Fit 最小化碎片
  3. 抢占式调度:高优先级任务可以抢占低优先级任务的资源
  4. 公平性保障:同一用户不能独占所有 GPU 资源

3.2 完整 Plugin 实现

package scheduling

import (
    "context"
    "time"

    v1 "k8s.io/api/core/v1"
    "k8s.io/apimachinery/pkg/runtime"
    framework "k8s.io/kubernetes/pkg/scheduler/framework"
)

const (
    // 优先级分层
    PriorityCritical = iota // P0: 生产级训练(可跨队列抢占)
    PriorityHigh            // P1: 实验训练(同队列内抢占)
    PriorityNormal          // P2: 普通任务(不抢占)

    cycleStateKey = "AISchedulingPluginData"
)

type AISchedulingPlugin struct {
    handle    framework.FrameworkHandle
    nodeCache map[string]*NodeInfo
    mu        sync.RWMutex
}

type StateData struct {
    podPriority int
    gpuRequest  int64
    queueID     string
}

func NewAISchedulingPlugin(obj runtime.Object, handle framework.FrameworkHandle) (framework.Plugin, error) {
    return &AISchedulingPlugin{
        handle:    handle,
        nodeCache: make(map[string]*NodeInfo),
    }, nil
}

// Name 返回插件名称
func (pl *AISchedulingPlugin) Name() string {
    return "AISchedulingPlugin"
}

// QueueSort:按 AI 优先级对调度队列排序
func (pl *AISchedulingPlugin) Less(podInfo1, podInfo2 *framework.QueuedPodInfo) bool {
    p1 := getPodPriority(podInfo1.Pod)
    p2 := getPodPriority(podInfo2.Pod)
    if p1 != p2 {
        return p1 < p2 // 数字越小优先级越高
    }
    // 同优先级按提交时间排序(FIFO)
    return podInfo1.Timestamp.Before(podInfo2.Timestamp)
}

// Filter:节点过滤 + 抢占决策
func (pl *AISchedulingPlugin) Filter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeInfo *framework.NodeInfo) *framework.Status {
    podPriority := getPodPriority(pod)
    requestedGPU := computeGPURequest(pod)
    availableGPU := getAvailableGPU(nodeInfo)

    // 资源充足,直接通过
    if availableGPU >= requestedGPU {
        return nil
    }

    // 资源不足,只有高优先级任务可以考虑抢占
    if podPriority > PriorityNormal {
        return framework.NewStatus(framework.Unschedulable, "insufficient GPU, low priority")
    }

    // 尝试抢占:找出可被抢占的低优先级 Pod
    victimPods := findPreemptablePods(nodeInfo, pod, requestedGPU-availableGPU)
    if len(victimPods) == 0 {
        return framework.NewStatus(framework.Unschedulable, "no preemptable pods found")
    }

    // 将抢占目标存入 CycleState
    state.Write(cycleStateKey, &StateData{
        podPriority: podPriority,
        gpuRequest:  requestedGPU,
        queueID:     getQueueID(pod),
        victims:     victimPods,
    })

    return nil
}

// PostFilter:处理无合适节点时的抢占逻辑
func (pl *AISchedulingPlugin) PostFilter(ctx context.Context, state *framework.CycleState, pod *v1.Pod, filteredNodeStatusMap framework.NodeToStatusMap) (*framework.PostFilterResult, *framework.Status) {
    s, err := state.Read(cycleStateKey)
    if err != nil {
        return nil, framework.NewStatus(framework.Error, err.Error())
    }
    cycleState := s.(*StateData)

    // 如果 Filter 阶段已识别抢占目标,执行抢占
    if len(cycleState.victims) > 0 {
        pl.executePreemption(cycleState.vitims)
        // 返回 Unschedule 让调度框架立即重试
        return nil, framework.NewStatus(framework.Unschedulable, "preemption in progress, will retry")
    }

    return nil, framework.NewStatus(framework.Unschedulable, "no suitable node")
}

// Score:Best-Fit GPU 碎片治理
func (pl *AISchedulingPlugin) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
    nodeInfo, _ := pl.handle.SnapshotSharedLister().NodeInfos().Get(nodeName)

    freeGPU := getAvailableGPU(nodeInfo)
    totalGPU := getTotalGPU(nodeInfo)
    requestedGPU := computeGPURequest(pod)

    if freeGPU < requestedGPU {
        return 0, nil
    }

    // Best-Fit:调度后剩余 GPU 越少分数越高
    remaining := freeGPU - requestedGPU

    // 完美契合:剩余为 0,给满分
    if remaining == 0 {
        return framework.MaxNodeScore, nil
    }

    // 剩余 1 张卡的情况:如果请求是奇数卡,可能产生碎片
    // 对 A100 (40/80GB) 等场景,1 张碎片影响很大
    remainingRatio := float64(remaining) / float64(totalGPU)

    // 碎片惩罚:剩余越少分数越高
    score := int64(float64(framework.MaxNodeScore) * (1.0 - remainingRatio))

    // 额外加分:NVLink 拓扑匹配
    if isNVLinkTopoMatch(pod, nodeInfo) {
        score += 20
    }

    return min(score, framework.MaxNodeScore), nil
}

// Reserve:预占资源
func (pl *AISchedulingPlugin) Reserve(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) *framework.Status {
    pl.mu.Lock()
    defer pl.mu.Unlock()

    nodeInfo, err := pl.handle.SnapshotSharedLister().NodeInfos().Get(nodeName)
    if err != nil {
        return framework.AsStatus(err)
    }

    // 添加到节点已分配 Pod 列表
    nodeInfo.AddPod(pod)
    return nil
}

// 执行驱逐
func (pl *AISchedulingPlugin) executePreemption(victims []*v1.Pod) {
    for _, victim := range victims {
        eviction := &policyv1.Eviction{
            ObjectMeta: metav1.ObjectMeta{
        Name:      victim.Name,
        Namespace: victim.Namespace,
            },
            DeleteOptions: &metav1.DeleteOptions{
                GracePeriodSeconds: ptr.To[int64](30),
            },
        }

        err := pl.handle.ClientSet().PolicyV1().Evictions(victim.Namespace).Evict(context.Background(), eviction)
        if err != nil {
            klog.Errorf("Failed to evict pod %s/%s: %v", victim.Namespace, victim.Name, err)
        } else {
            klog.Infof("Successfully evicted pod %s/%s for high-priority task", victim.Namespace, victim.Name)
        }
    }
}

3.3 注册与配置

# scheduler-config.yaml
apiVersion: kubescheduler.config.k8s.io/v1
kind: KubeSchedulerConfiguration
profiles:
  - schedulerName: ai-scheduler
    plugins:
      queueSort:
        enabled:
          - name: AISchedulingPlugin
        disabled:
          - name: "*"
      preFilter:
        enabled:
          - name: AISchedulingPlugin
      filter:
        enabled:
          - name: AISchedulingPlugin
      score:
        enabled:
          - name: AISchedulingPlugin
          - name: NodeResourcesFit
            weight: 50
      reserve:
        enabled:
          - name: AISchedulingPlugin
      permit:
        enabled:
          - name: AISchedulingPlugin
      postFilter:
        enabled:
          - name: AISchedulingPlugin
    pluginConfig:
      - name: AISchedulingPlugin
        args:
          priorityPreemptionThreshold: 1
          maxPreemptPodsPerCycle: 3
          fairnessEnabled: true

四、生产部署关键经验

4.1 并发安全

调度器是高度并发的,多个 Pod 的 Scheduling Cycle 可能同时进行。CycleState 本身是线程安全的(每个调度周期独立),但 Plugin 内部维护的状态(如 nodeCache)需要自行加锁:

// 正确:使用 CycleState 传递插件私有数据
state.Write(cycleStateKey, &StateData{...})

// 错误:直接用 plugin 字段存储调度周期数据
pl.currentPod = pod // 并发不安全!

4.2 抢占延迟控制

抢占式调度的核心挑战在于 Graceful Termination:被抢占的 Pod 需要时间保存 Checkpoint、通知训练框架退出。生产环境的典型配置:

# 被抢占 Pod 的 PreStop Hook 配置
lifecycle:
  preStop:
    exec:
      command:
        - /bin/sh
        - -c
        - python /opt/checkpoint/save.py && sleep 5
terminationGracePeriodSeconds: 120  # 给足够时间保存状态

4.3 调度性能调优

在高规模集群(1000+ Node)中,Score 插件的性能至关重要:

// 预计算节点 GPU 拓扑,避免重复遍历
func (pl *AISchedulingPlugin) Score(ctx context.Context, state *framework.CycleState, pod *v1.Pod, nodeName string) (int64, *framework.Status) {
    // 使用缓存的 nodeInfo 而不是实时查询
    cachedInfo := pl.getCachedNodeInfo(nodeName)
    if cachedInfo == nil {
        return 0, nil
    }

    // 简单算术运算,避免 I/O
    score := computeScore(cachedInfo.freeGPU, cachedInfo.totalGPU, requestedGPU)
    return score, nil
}

4.4 多租户公平性

在多租户场景中,为避免单一用户占用全部 GPU,可以在 Score 阶段加入公平性因子:

// 如果某个用户的 Pod 已使用超过其配额的 70%,降低其新任务的调度分数
userUtilization := calculateUserGPUUtilization(nodeInfo, pod.Namespace)
if userUtilization > 0.7 {
    score = int64(float64(score) * 0.5) // 惩罚系数
}

五、总结

Kubernetes Scheduling Framework 通过精心设计的 10+ 扩展点,实现了调度逻辑的完全可插拔化。掌握这套框架后,你可以:

  1. 资源碎片治理:通过 Best-Fit Score 最大化 GPU 利用率
  2. 多级抢占:按优先级实现任务抢占,保障关键训练任务 SLA
  3. 公平调度:跨用户、跨队列的资源配额管理
  4. 拓扑感知:NVLink、InfiniBand 拓扑感知调度

这些能力对于构建高效、稳定的 AI 训练平台至关重要。调度框架的设计哲学 —— "不修改源码,只扩展逻辑" —— 本身就是 Kubernetes 设计原则的精髓。

源码参考:kubernetes/pkg/scheduler/framework/interface.go 定义了所有扩展点接口,kubernetes/pkg/scheduler/ 目录下的 generic_scheduler.go 展示了 Scheduler 如何驱动这些 Plugin。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部