Kubernetes 调度器深度实战:从抢占式调度到拓扑感知的 AI 训练任务编排

Kubernetes 调度器深度实战:从抢占式调度到拓扑感知的 AI 训练任务编排

在 AI/ML 工作负载云原生化的浪潮中,Kubernetes 已经从通用编排器进化为 AI 基础设施的核心调度引擎。本文将从 K8s 调度框架的源码级架构出发,深入剖析抢占式调度、拓扑感知调度、Gang Scheduling 三大核心机制的生产级实战细节,并结合 Volcano/Kueue 等 AI 调度扩展,揭示如何在大规模 GPU 集群中实现高效的任务编排。

一、调度框架架构演进:从 Monolithic 到插件化

Kubernetes 1.15 正式 GA 的 Scheduling Framework(KEP-624)彻底重构了调度器架构。在理解 AI 训练调度之前,必须掌握这套插件化扩展机制。

1.1 调度生命周期

一个 Pod 进入调度流程后,依次穿越以下阶段:

┌─────────────┐    ┌─────────────┐    ┌──────────────┐
│ QueueSort   │───▶│ PreFilter   │───▶│ Filter       │
│ (队列排序)   │    │ (预检查)     │    │ (节点过滤)    │
└─────────────┘    └─────────────┘    └──────────────┘
                                            │
                                            ▼
┌─────────────┐    ┌─────────────┐    ┌──────────────┐
│ Bind        │◀───│ Score       │◀───│ PostFilter   │
│ (绑定)       │    │ (评分)       │    │ (后置过滤)    │
└─────────────┘    └─────────────┘    └──────────────┘

每个阶段都是一个扩展点(Extension Point),允许注册多个插件。调度器按优先级队列排序 Pod,然后遍历过滤/评分流水线,最终选出最优节点。

1.2 源码级架构:SchedulingQueue 与 ScheduleAlgorithm

// pkg/scheduler/internal/cache.go — 节点快照缓存
type Snapshot interface {
    NodeInfos() NodeInfoLister
    HavePodsWithAffinityMap() sets.String
}

// pkg/scheduler/schedule_one.go — 单次调度入口
func (sched *ScheduleAlgorithm) scheduleOne(
    ctx context.Context,
    state *framework.CycleState,
    pod *v1.Pod,
) {
    // 1. 选择下一个高优 Pod
    podInfo := sched.NextPod()

    // 2. 调用 ScheduleAlgorithm.Schedule
    result, err := sched.Schedule(ctx, state, podInfo)

    // 3. 处理抢占或绑定
    if fitError != nil {
        // 尝试抢占
        nominatedNode := sched.preempt(ctx, ...).nominatedNode
    }
}

对于 AI 训练任务,这里有两个关键调度特征:

  • 资源需求量大:单个 Pod 可能需要 8xA100 GPU + 多节点 GPU 互联
  • Gang Scheduling 要求:所有 Worker 必须同时调度,否则整个 Job 无法启动

原生调度器无法直接满足 Gang Scheduling 需求,这是 AI 调度扩展的核心动机。

二、抢占式调度:源码级深度解析

抢占(Preemption)是批量任务队列的基础设施。当一个高优 Pod 找不到可调度节点时,调度器会主动驱逐低优 Pod,释放资源空间。

2.1 抢占流程源码分析

// pkg/scheduler/scheduling_queue.go — 触发抢占的入口
func (p *PriorityQueue) Pop() (*framework.QueuedPodInfo, error) {
    // ...
    pod := p.podInfo.pop().Pod

    // 检查是否前一个周期有 scheduling cycle  失败
    // 若有,触发Nominated节点抢占回调
    return podInfo, nil
}

// pkg/scheduler/scheduler.go — preempt 核心逻辑
func (sched *Scheduler) preempt(
    ctx context.Context,
    fwk framework.Framework,
    state *framework.CycleState,
    pod *v1.Pod,
    fitError *FitError,
) (string, error) {
    // 1. 找到被抢占者所在节点
    preemptor := pod
    node, victims, nominatedPods, err := sched.Algorithm.FindCandidates(
        ctx, fwk, preemptor, ...,
    )

    // 2. 执行抢占:先删除 victim Pod,再激活 Nominated
    for _, victim := range victims {
        if err := sched.Client.CoreV1().
            Pods(victim.Namespace).
            Delete(ctx, victim.Name, ...); err != nil {
            return "", err
        }
    }

    return node, nil
}

2.2 生产实践中抢占的三大陷阱

陷阱一:级联抢占引发的雪崩

当集群持续涌入高优 Pod 时,频繁的抢占-驱逐循环会造成整个集群不稳定。实测某电商大促期间曾因抢占策略配置不当,导致核心应用 P99 延迟飙涨 300%。

# 反例:未限制的 PriorityClass
apiVersion: scheduling.k8s.io/v1
kind: PriorityClass
metadata:
  name: critical-gpu-job
value: 1000000  # 设为最高值,但缺乏多Anti-camp机制
globalDefault: false
preemptionPolicy: PreemptLowerPriority  # 直接抢占

# 正例:可控的抢占策略 + 租户隔离
apiVersion: scheduling.k8s.io/v1
kind: PriorityClass
metadata:
  name: ml-batch-high
value: 10000
preemptionPolicy: PreemptLowerPriority
description: "ML训练高优任务 - 不应抢占在线服务"

# 配合 PodDisruptionBudget 保护关键 Pod
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata:
  name: training-node-pdb
spec:
  minAvailable: 0
  selector:
    matchLabels:
      app: training-worker
# minAvailable=0 允许被抢占,但限制并发抢占比例

陷阱二:抢占后的 Nominated Pods 泄漏

被抢占驱逐的 Pod 不会自动触发重调度。当资源碎片化严重时,大量 Nominated Pod 会堆积在调度队列中。

// 源码中 FindCandidates 的过滤逻辑
// pkg/scheduler/framework/plugins/helper/nominator.go
func NominatedNodeName(pod *v1.Pod) string {
    return pod.Status.NominatedNodeName
}

生产上需配合 nominatedNodeName 状态清理机制——通过 CronJob 扫描并驱逐超过 TTL 的 Nominated Pod。

陷阱三:跨 Namespace 抢占的安全隐患

默认抢占策略允许跨 Namespace 抢占。在多租户场景中,高优租户可能抢占低租户的资源。解决方案是配合 Kueue 的 Cohort/Queue 机制,在租户分账层面限制抢占范围。

三、拓扑感知调度:NUMA 与 GPU 互联拓扑

对于 AI 训练任务,GPU 之间通过 NVLink/NVSwitch 互联。如果 Worker Pod 内的 GPU 跨越 NUMA 节点,训练效率会下降 30%~50%。

3.1 Topology Manager 工作原理

TopologyManager 是 kubelet 而非 scheduler 的职责,但与调度结果密切相关:

┌──────────────────────────────────────────────────────┐
│                    TopologyManager                    │
├──────────────────────────────────────────────────────┤
│  Policy: single-numa-node / best-effort /          │
│          restricted / none                         │
├──────────────────────────────────────────────────────┤
│  1. Pod 进入 Admission                               │
│  2. Gather resource containers (CPU, GPU, NIC)      │
│  3. Compute Affinity Mask across NUMA nodes        │
│  4. Deny if cross-NUMA (policy=restricted)         │
└──────────────────────────────────────────────────────┘

kubelet 配置示例:

# /var/lib/kubelet/config.yaml
cpuManagerPolicy: static
topologyManagerPolicy: single-numa-node
reservedSystemCPUs: "0-7"

3.2 Scheduler 侧的拓扑感知 plugin

调度器通过 NodeResourceTopology Informer 获取 NUMA 拓扑信息,在 Filter/Score 阶段做出优化决策:

// pkg/scheduler/framework/plugins/noderesourcetopology/
// filter.go — NUMA 对齐过滤
type NUMANodeAffinity struct {
    nodes sets.Int  // 包含关键资源的 NUMA 节点集合
}

func (nm *NUMANodeAffinity) Preferred() bool {
    return nm.nodes.Len() <= 1  // 单 NUMA 才最优
}

// score.go — NUMA 跨度越少分数越高
func (f *NodeResourceTopologyMatch) Score(
    ctx context.Context,
    cycleState *framework.CycleState,
    pod *v1.Pod,
    nodeName string,
) (int64, *framework.Status) {
    nodeTopology := f.nodeTopologiesMap[nodeName]

    // 计算该 Pod 所需资源的 NUMA 跨度
    numaID := nodeTopology.GuessNUMANodeAffinity(podSpec)

    // 跨 NUMA 越多,分数越低
    if numaID.Span() > 1 {
        return 100 - int64(numaID.Span()*20), nil
    }
    return 100, nil
}

3.3 GPU 拓扑实战案例

某 AI 训练集群配置:8 节点,每节点 8x A100-80G(通过 NVSwitch 全互联)。关键配置:

apiVersion: v1
kind: Pod
metadata:
  annotations:
    # GPU 拓扑亲和性:偏好 NVSwitch 连接的 GPU 组
    scheduling.k8s.io/nvidia-gpu-topology: "nvswitch-full-mesh"
spec:
  containers:
  - name: trainer
    resources:
      limits:
        nvidia.com/gpu: 8
        memory: "640Gi"  # 匹配 A100-80G + NVLink 带宽
        hugepages-2Mi: "128Mi"
  topologySpreadConstraints:
  - maxSkew: 1
    topologyKey: topology.kubernetes.io/zone
    whenUnsatisfiable: DoNotSchedule
    labelSelector:
      matchLabels:
        app: ai-training

四、AI 训练调度:Volcano 与 Kueke 的工程实践

针对 AI 训练任务的 Gang Scheduling 需求,社区形成了以 Volcano 和 Kueue 为代表的两大调度扩展方案。

4.1 Volcano:AI 训练专属调度器

Volcano 在 K8s Job 之上抽象了 vcjob/v1alpha1,支持:

  • Gang Scheduling:minAvailable 保证最少 Worker 数
  • Binpack / Spread / FFD 多维评分策略
  • Task-level Priority:同一 Job 内不同 Stage 设置不同优先级
  • Reclaim:比抢占更优雅的优先级资源释放
apiVersion: batch.volcano.sh/v1alpha1
kind: Job
metadata:
  name: pytorch-distributed-llama
spec:
  minAvailable: 8  # Gang:必须 8 个 Worker 同时 Ready
  schedulerName: volcano
  queue: ml-high
  policies:
  - event: PodFailed
    action: RestartJob
  - event: TaskCompleted
    action: CompleteJob
  plugins:
    ssh: []
    env: []
    svc: []
  tasks:
  - replicas: 1
    name: master
    template:
      spec:
        containers:
        - name: trainer
          image: nvcr.io/nvidia/pytorch:24.01-py3
          command: ["torchrun", "--nproc_per_node=8",
                    "--nnodes=8", "--node_rank=0",
                    "train.py", "--model", "llama-70b"]
          resources:
            limits:
              nvidia.com/gpu: 8
              cpu: "128"
              memory: "640Gi"
    - replicas: 7
      name: worker
      template:
        spec:
          containers:
          - name: trainer
            image: nvcr.io/nvidia/pytorch:24.01-py3
            command: ["sleep", "infinity"]  # 等待 Master 启动
            resources:
              limits:
                nvidia.com/gpu: 8
                cpu: "128"
                memory: "640Gi"

Volcano 的 Gang Scheduling 实现基于两个组件:

// volcano.sh/volcano/pkg/scheduler/plugins/gang/gang.go
type gangPlugin struct {
    // Job 维度:追踪每个 vcjob 的已调度 Pod 数
    jobWaitingPods map[api.JobID]int32
    jobMinTaskNum  map[api.JobID]int32
}

func (gp *gangPlugin) OnSessionOpen(ssn *framework.Session) {
    // 过滤阶段:不满足 minAvailable 时整体过滤该 Job 的 Pod
    ssn.AddJobValidFn(gp.Name(), func(jobInfo *api.JobInfo) bool {
        return jobInfo.MinAvailable >= len(jobInfo.WaitingTaskList)
    })

    // 资源分配阶段:采用 All-or-Nothing 策略
    ssn.AddReclaimableFn(gp.Name(), func(task *api.TaskInfo, ...) bool {
        return task.Job != targetTask.Job
    })
}

4.2 Kueue:云原生 Batch 队列入门

Kueue 1.0(2024 GA)以 Criterion 和 Cohort 为核心,提供多租户公平调度:

apiVersion: kueue.x-k8s.io/v1beta1
kind: ResourceFlavor
metadata:
  name: a100-80g
spec:
  nodeLabels:
    nvidia.com/gpu.product: "NVIDIA-A100-SXM4-80GB"
  tolerations:
  - key: nvidia.com/gpu
    operator: Exists
    effect: NoSchedule

---
apiVersion: kueue.x-k8s.io/v1beta1
kind: ClusterQueue
metadata:
  name: ml-training-cq
spec:
  namespaceSelector: {}
  resourceGroups:
  - coveredResources: ["cpu", "memory", "nvidia.com/gpu"]
    flavors:
    - name: a100-80g
      resources:
      - name: cpu
        nominalQuota: "1024"
      - name: memory
        nominalQuota: "4Ti"
      - name: nvidia.com/gpu
        nominalQuota: "64"
  cohort: "gpu-cluster"  # 共享 Cohort 资源池

  # 多租户公平性:DRF(Dominant Resource Fairshare)
  fairSharing:
    enable: true
    weight: "1"

---
apiVersion: kueue.x-k8s.io/v1beta1
kind: LocalQueue
metadata:
  namespace: ml-team-a
  name: training-queue
spec:
  clusterQueue: ml-training-cq
  priorityClassName: ml-batch-normal

Kueue 的三大创新点:

  1. 两级配额管理:ClusterQueue 分配总额,LocalQueue 按 Namespace 二次分账
  2. 自动暂停/恢复:资源不足时 Suspend 入场 Job,释放后自动 Unsuspend
  3. 资源借用:Cohort 内闲置资源可跨 CQ 借用,用后归还

4.3 Volcano vs Kueue 选型

维度 Volcano Kueue
控制粒度的 基于任务的细粒度 Job 级 批处理队列级配额
多租户公平性 粗粒度(Queue) 细粒度(CQ + DRF)
抢占/回收 内置 Reclaim 依赖 PriorityClass
AI 框架集成 原生 TFJob/PyTorch 需配合 TorchX/RunAI
适用场景 大规模纯训练集群 多业务共享集群

实测数据:在某 512 节点 A100 集群中,使用 Volcano Gang Scheduling 的 Job 平均排队时间为原生 K8s 的 1/3,集群 GPU 利用率从 65% 提升至 89%。

五、生产级实战:多集群调度与 NVMe 亲和性

5.1 资源碎片与 Binpack 策略

AI 训练任务对 GPU 数量要求固定(如 4/8/16),容易产生碎片。Binpack 优先填满已有负载节点,但可能导致热点;Spread 则分散部署,但加剧网络开销。

// volcano/pkg/scheduler/plugins/binpack/binpack.go
func (bp *binpackPlugin) Score(
    state *framework.CycleState,
    task *api.TaskInfo,
    nodeName string,
) (int64, *framework.Status) {
    nodeInfo := bp.session.Nodes[nodeName]
    requested := task.Resreq

    // 剩余资源越少(越紧凑)分数越高
    idleCPU := nodeInfo.Allocatable.MilliCPU - nodeInfo.Requested.MilliCPU
    idleMEM := nodeInfo.Allocatable.Memory - nodeInfo.Requested.Memory

    // 自定义评分:考虑 GPU 带宽与网络拓扑
    score := (idleCPU + idleMEM) / 100  // 越小越好
    return score, framework.NewStatus(framework.Success)
}

最佳实践:对小型推理任务使用 Spread(避免热点),对大型训练任务使用 Binpack(减少碎片)。

5.2 NVMe 亲和性与本地 SSD 调度

AI 训练数据集通常缓存在节点本地 NVMe 盘上(可达 30GB/s 读速)。实践中发现多个 Worker Pod 争抢同一节点 NVMe 带宽会导致 IO 抖动。

解决方案:将 NVMe 作为扩展资源声明:

apiVersion: v1
kind: Node
metadata:
  labels:
    nvme-bandwidth: "high"
    nvme-iops: "2000000"
spec:
  taints:
  - key: nvme.local.data
    value: "cached"
    effect: NoSchedule
---
apiVersion: v1
kind: Pod
spec:
  containers:
  - resources:
      limits:
        nvidia.com/gpu: 8
        local-nvme-iops: "500000"  # 自定义资源
  tolerations:
  - key: nvme.local.data
    operator: Equal
    value: "cached"
    effect: NoSchedule
// 自定义 Scheduler Plugin:NVMe 带宽感知过滤
type NVMeAffinityPlugin struct{}

func (n *NVMeAffinityPlugin) Filter(
    ctx, cycleState, pod, nodeInfo,
) (*framework.Status) {
    requestedIOPS := pod.Annotations["nvme-iops-request"]
    availableIOPS := nodeInfo.Node().Capacity["local-nvme-iops"]

    if requestedIOPS > availableIOPS {
        return framework.NewStatus(
            framework.Unschedulable,
            "niibe bandwidth insufficient",
        )
    }
    return nil
}

5.3 故障域与 Gpu 代码的强关联调度

当某 GPU 节点出现 ECC 错误、XID 错误(常见于超频或散热故障)时,调度器应自动避开该 node。结合 Node Feature Discovery (NFD) 和自定义插件:

// package gpuhealth
func (h *GPUHealthPlugin) Filter(
    ctx context.Context,
    state *framework.CycleState,
    pod *v1.Pod,
    nodeInfo *framework.NodeInfo,
) *framework.Status {
    node := nodeInfo.Node()

    // 检查 XID 错误污点
    for _, taint := range node.Spec.Taints {
        if taint.Key == "nvidia.com/xid-error" {
            return framework.NewStatus(
                framework.Unschedulable,
                fmt.Sprintf("node %s has GPU XID unhealthy", node.Name),
            )
        }
    }

    // 检查 ECC 错误数(自定义 Node Label)
    eccCount := node.Labels["nvidia.com/ecc-uncorrected"]
    if strconv.Atoi(eccCount) > 10 {
        return framework.NewStatus(framework.Unschedulable,
            "GPU ECC error count exceeds threshold")
    }
    return nil
}

六、前沿趋势:DRA 与动态资源分配

Kubernetes 1.30 GA 的 Dynamic Resource Allocation (DRA) 为 AI 训练任务提供了更灵活的资源调度能力——GPU 切片(MIG)、vGPU 时分复用、跨节点 GPU 共享。

apiVersion: resource.k8s.io/v1alpha3
kind: ResourceClaimTemplate
metadata:
  name: gpu-claim-template
spec:
  spec:
    resourceClassName: nvidia-gpu-a100
    parameters:
      apiVersion: gpu.resource.nvidia.com/v1alpha1
      kind: GPUClaimParameters
      flags:
        mig:
          enabled: true
          profile: "1g.10gb"  # MIG 切片:10GB 显存切片
---
apiVersion: v1
kind: Pod
spec:
  resourceClaims:
  - name: training-gpu
    source:
      resourceClaimTemplateName: gpu-claim-template
  containers:
  - name: trainer
    resources:
      claims:
      - name: training-gpu

DRA 的正式发布使得单节点可承载更多小模型推理任务,将 GPU 利用率从传统分区的 60% 提升至 95%。

七、总结

Kubernetes 调度器为 AI 训练任务提供了强大的声明式编排能力,但原生调度器在设计之初并未考虑 Gang Scheduling、拓扑亲和、NUMA 对齐等 AI 场景需求。生产实践中:

  1. 小集群(<32节点):K8s 原生调度器 + 合理 PriorityClass + PDB 即可满足
  2. 中等集群(32-256节点):Volcano Gang Scheduling + 抢占策略是最佳选择
  3. 超大规模(>256节点):Kueue 多租户配额管理 + Volcano 队列 + Binpack/Spread 混合策略

GPU 拓扑感知、NVMe 带宽亲和、XID 健康检测等细节,决定了生产集群的最终效率。在 AI 训练从"能跑"到"高效跑"的演进中,调度器架构的不断革新将持续释放集群算力潜力。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部