Skip to content

필사 모드: Kubernetes AI 训练流水线:Volcano、Training Operator、Kueue 深度分析

中文
0%
정확도 0%
💡 왼쪽 원문을 읽으면서 오른쪽에 따라 써보세요. Tab 키로 힌트를 받을 수 있습니다.

1. Kubernetes 默认调度器的局限

Kubernetes 的默认调度器(kube-scheduler)以 Pod 为单位执行调度。根据单个 Pod 所请求的资源(CPU、Memory、GPU)将其放置到合适的 Node 上,是它的核心职责。这种结构适合 Web 服务、API 服务器这类常规工作负载,但面对 AI/ML 分布式训练这样的批处理工作负载时,存在根本性的局限。

1.1 缺少 Gang Scheduling

分布式训练中最大的问题是缺少 Gang Scheduling。例如在 PyTorch DDP(Distributed Data Parallel)训练中,4 个 Worker Pod 必须全部同时启动,训练才能推进。默认调度器逐个调度 Pod,因此可能出现 4 个中只有 3 个被调度、剩下 1 个因资源不足停留在 Pending 状态的情况。已经放置的 3 个 Pod 占用着 GPU 资源,训练却无法开始,形成 资源浪费 的状态。

1.2 死锁问题

更严重的是死锁。假设两个分布式训练 Job 各自需要 4 块 GPU,而集群总共只有 6 块 GPU。默认调度器可能给 Job A 分配 3 块、给 Job B 分配 3 块。两个 Job 都至少需要 4 块,于是哪一个都无法开始训练。这正是 Gang Scheduling 所要解决的 “All or Nothing” 问题。

1.3 缺少 Fair-share 与基于 Queue 的资源管理

默认调度器没有用于团队间资源共享的 Queue 概念。在 AI 组织中,多个研究团队共享同一套 GPU 集群,而仅凭默认调度器很难实现按团队设置 Quota、并把空闲资源共享出去的 Fair-share 策略。

为了解决这些局限,Volcano、Kubeflow Training Operator、Kueue 等项目应运而生。


2. Volcano Scheduler 官方文档分析

Volcano 是 CNCF(Cloud Native Computing Foundation)旗下的项目,是一个用于在 Kubernetes 上处理高性能批处理工作负载的调度器兼控制器。它是面向 AI/ML、BigData、HPC(High Performance Computing)等工作负载的批处理调度系统。

2.1 架构

Volcano 由三个核心组件构成。

  • Volcano Scheduler:替代 kube-scheduler 或与其并行运行的调度器。执行的是 Job/PodGroup 级别而非 Pod 级别的调度。
  • Volcano Controller Manager:管理 VolcanoJob、Queue、PodGroup 等 CRD 生命周期的控制器。
  • Volcano Webhook(Admission Controller):负责 CRD 创建时的有效性校验与默认值注入的 Admission Webhook。

Volcano Scheduler 与 Kubernetes 的 Scheduling Framework 兼容,通过 predicates 和 nodeorder 插件执行默认调度器的 PreFilter/Filter 与 Score 阶段。在此之上,以插件的形式追加 Gang Scheduling、Fair-share、Binpack 等高级调度策略。

2.2 Queue CRD

Queue 是把集群资源做逻辑切分、在多租户环境中提供资源隔离的 CRD。

apiVersion: scheduling.volcano.sh/v1beta1
kind: Queue
metadata:
  name: ml-research-queue
spec:
  weight: 4
  reclaimable: true
  capability:
    cpu: '64'
    memory: '256Gi'
    nvidia.com/gpu: '8'
  • weight:Queue 之间的资源分配比例。weight 越高,分配到的资源越多。
  • reclaimable:是否可以把空闲资源借给其他 Queue。设置为 true 时,该 Queue 的空闲资源可以被其他 Queue 使用。
  • capability:该 Queue 可以使用的资源上限。需要明确写出 CPU、Memory、GPU 等。

2.3 PodGroup CRD

PodGroup 是定义一组彼此强关联的 Pod 的 CRD,是 Gang Scheduling 的核心单位。

apiVersion: scheduling.volcano.sh/v1beta1
kind: PodGroup
metadata:
  name: pytorch-training-pg
  namespace: default
spec:
  minMember: 4
  minResources:
    cpu: '16'
    memory: '64Gi'
    nvidia.com/gpu: '4'
  priorityClassName: high-priority
  queue: ml-research-queue
  • minMember:PodGroup 内最少需要运行的 Pod 数量。若集群资源无法满足该数量,则 PodGroup 内的任何一个 Pod 都不会被调度。
  • minResources:运行 PodGroup 所需的最小资源。若可用资源无法满足,调度就会被搁置。
  • priorityClassName:PodGroup 的优先级。调度器在对 Queue 内的多个 PodGroup 排序时使用。
  • queue:PodGroup 所属的 Queue 名称。该 Queue 必须处于 Open 状态。

PodGroup 的 Status Phase 按如下方式迁移:Pending(资源未满足)-> Inqueue(校验完成,等待节点绑定)-> Running(已有不少于 minMember 的 Pod 在运行)-> Unknown(部分 Pod 无法调度)。创建 VolcanoJob 时,PodGroup 会被自动创建。


3. Volcano 的 Gang Scheduling 与 Fair-share 调度策略

3.1 Gang Scheduling

Gang Scheduling 是 Volcano 最核心的功能。它采用 “All or Nothing” 策略,只有当构成 Job 的全部 Pod(或 minMember 指定的最少 Pod 数量)能够同时被调度时才执行放置。

Gang Scheduling 在以下场景中不可或缺:

  • 基于 MPI 的分布式训练:所有 Worker 必须同时启动,allreduce 通信才能成立。
  • PyTorch DDP:Master 与 Worker 全部就绪后训练才会开始。
  • Spark/Big Data:Driver 与 Executor 全部就绪后 Job 才会执行。

Gang Plugin 会检查 PodGroup 的 minMember 字段,判断集群能否同时调度该数量的 Pod。若不可能,则一个 Pod 也不调度,从而避免资源浪费与死锁。

3.2 DRF(Dominant Resource Fairness)调度

DRF 是在多租户环境中实现公平资源分配的算法。它以每个 Job 的 主导资源(Dominant Resource)为基准计算 Fair-share。

例如,Job A 以 CPU 为主(CPU 8 核、Memory 4GB),Job B 以内存为主(CPU 2 核、Memory 32GB),那么就优先调度主导资源占比最低的 Job,以此确保公平性。DRF Plugin 同时考虑多维资源(CPU、Memory、GPU),执行不偏向特定资源类型的公平分配。

3.3 Proportion Plugin

Proportion Plugin 管理各 Queue 的资源比例。它以每个 Queue 的 weight 为基准按比例分配整个集群的资源,并在 Queue 的 capability 范围内进行分配。当多个团队共享 GPU 集群时,可以为每个团队创建 Queue 并设置 weight,从而调整资源分配比例。


4. Volcano Job CRD 与插件系统

4.1 VolcanoJob CRD

VolcanoJob 是定义批处理工作负载的 Volcano 核心 CRD。它可以把多个 Task(角色)打包成一个 Job 进行管理。

apiVersion: batch.volcano.sh/v1alpha1
kind: Job
metadata:
  name: pytorch-ddp-training
spec:
  minAvailable: 5
  schedulerName: volcano
  queue: ml-research-queue
  plugins:
    env: []
    svc: []
  policies:
    - event: PodEvicted
      action: RestartJob
  tasks:
    - replicas: 1
      name: master
      template:
        spec:
          containers:
            - image: pytorch/pytorch:2.1.0-cuda12.1-cudnn8-runtime
              name: master
              command: ['torchrun']
              args: ['--nproc_per_node=1', '--nnodes=5', '--node_rank=0', 'train.py']
              resources:
                limits:
                  nvidia.com/gpu: 1
          restartPolicy: OnFailure
    - replicas: 4
      name: worker
      template:
        spec:
          containers:
            - image: pytorch/pytorch:2.1.0-cuda12.1-cudnn8-runtime
              name: worker
              command: ['torchrun']
              args: ['--nproc_per_node=1', '--nnodes=5', 'train.py']
              resources:
                limits:
                  nvidia.com/gpu: 1
          restartPolicy: OnFailure
  • minAvailable:Gang Scheduling 的最小 Pod 数量。与 PodGroup 的 minMember 作用相同。
  • schedulerName:指定为 volcano 以使用 Volcano Scheduler。
  • queue:Job 提交到的 Queue 名称。
  • plugins:要启用的 Volcano 插件列表。env 负责自动注入环境变量,svc 负责自动创建服务。
  • policies:基于事件的动作策略。可定义诸如 Pod 被 Evict 后重启 Job 之类的策略。
  • tasks:构成 Job 的 Task 列表。每个 Task 按角色(master、worker)定义 replicas 与 Pod Template。

4.2 插件系统详解

Volcano Scheduler 采用插件式架构,可以灵活扩展调度策略。主要插件如下。

Plugin功能适合的工作负载
GangAll or Nothing 调度分布式训练、MPI、Spark
Binpack尽量填满已有节点后再放置小规模 Job、最大化资源效率
DRF基于 Dominant Resource 的公平共享多租户批处理
Proportion管理各 Queue 的资源比例按团队做资源隔离
Priority基于优先级的调度优先级各异的多个 Job
Task-TopologyTask 间 Affinity/Anti-Affinity需要网络优化的分布式训练
TDM分时资源共享Kubernetes + YARN 混合环境
SLA等待时间上限有实时服务要求的批处理
Numa-aware感知 NUMA 拓扑的调度CPU 缓存亲和性重要的 HPC
Predicates基于 GPU 资源的前置过滤GPU 密集型 AI 工作负载
Nodeorder多维节点打分复合调度定制

在 Volcano Scheduler 的配置文件中,通过 actions 与 tiers 组合插件:

actions: 'enqueue, allocate, backfill'
tiers:
  - plugins:
      - name: priority
      - name: gang
      - name: conformance
  - plugins:
      - name: drf
      - name: predicates
      - name: proportion
      - name: nodeorder
      - name: binpack

5. Kubeflow Training Operator 官方文档分析

Kubeflow Training Operator 是在 Kubernetes 上管理分布式 AI 训练 Job 的 Operator。它支持 PyTorch、TensorFlow、MPI、XGBoost 等多种框架,并为每种框架提供对应的 CRD。

5.1 V1(Legacy)与 V2 架构

目前 Training Operator 有两个版本并存。

V1(Legacy):使用按框架划分的独立 CRD(PyTorchJob、TFJob、MPIJob 等)。每个 CRD 提供符合该框架分布式训练模式的 ReplicaSpec(Master、Worker、PS 等)。

V2(Current):统一为单一的 TrainJob CRD。框架相关配置通过 TrainingRuntime 与 ClusterTrainingRuntime 管理。它支持 PyTorch、MLX、HuggingFace、DeepSpeed、JAX、XGBoost 等更广泛的框架。

本文以实际生产中广泛使用的 V1 的 PyTorchJob 为中心展开分析,同时也一并观察 V2 的方向。

5.2 Training Operator 的职责

Training Operator 会自动处理以下事项:

  1. 自动设置环境变量:把 WORLD_SIZERANKMASTER_ADDRMASTER_PORT 等分布式训练所需的环境变量自动注入 Pod。
  2. torchrun CLI 配置:配置环境使 PyTorch 的 torchrun(此前的 torch.distributed.launch)能够正确运行。
  3. Job 状态管理:管理 Created -> Running -> Succeeded/Failed 等生命周期,并在失败时应用重启策略。
  4. 创建 Service:自动创建用于 Pod 间通信的 Headless Service。

6. PyTorchJob CRD 详解与 YAML 示例

PyTorchJob 是用于在 Kubernetes 上运行 PyTorch 分布式训练 Job 的 CRD。

6.1 PyTorchJob 结构

apiVersion: kubeflow.org/v1
kind: PyTorchJob
metadata:
  name: pytorch-ddp-mnist
  namespace: ml-training
spec:
  pytorchReplicaSpecs:
    Master:
      replicas: 1
      restartPolicy: OnFailure
      template:
        metadata:
          annotations:
            sidecar.istio.io/inject: 'false'
        spec:
          containers:
            - name: pytorch
              image: my-registry/pytorch-training:latest
              command:
                - 'torchrun'
                - '--nproc_per_node=2'
                - '--nnodes=3'
                - '--node_rank=0'
                - '--master_addr=$(MASTER_ADDR)'
                - '--master_port=$(MASTER_PORT)'
                - 'train.py'
                - '--epochs=100'
                - '--batch-size=256'
              env:
                - name: NCCL_DEBUG
                  value: 'INFO'
              resources:
                limits:
                  nvidia.com/gpu: 2
                  memory: '32Gi'
                  cpu: '8'
              volumeMounts:
                - name: training-data
                  mountPath: /data
                - name: checkpoints
                  mountPath: /checkpoints
          volumes:
            - name: training-data
              persistentVolumeClaim:
                claimName: training-data-pvc
            - name: checkpoints
              persistentVolumeClaim:
                claimName: checkpoints-pvc
    Worker:
      replicas: 2
      restartPolicy: OnFailure
      template:
        metadata:
          annotations:
            sidecar.istio.io/inject: 'false'
        spec:
          containers:
            - name: pytorch
              image: my-registry/pytorch-training:latest
              command:
                - 'torchrun'
                - '--nproc_per_node=2'
                - '--nnodes=3'
                - '--master_addr=$(MASTER_ADDR)'
                - '--master_port=$(MASTER_PORT)'
                - 'train.py'
                - '--epochs=100'
                - '--batch-size=256'
              env:
                - name: NCCL_DEBUG
                  value: 'INFO'
              resources:
                limits:
                  nvidia.com/gpu: 2
                  memory: '32Gi'
                  cpu: '8'
              volumeMounts:
                - name: training-data
                  mountPath: /data
                - name: checkpoints
                  mountPath: /checkpoints
          volumes:
            - name: training-data
              persistentVolumeClaim:
                claimName: training-data-pvc
            - name: checkpoints
              persistentVolumeClaim:
                claimName: checkpoints-pvc

6.2 核心字段说明

  • Master:必须始终为 replicas: 1。它代表训练 Job 的状态,并充当模型参数聚合的基准点。Training Operator 会把该 Pod 的地址以 MASTER_ADDR 注入其他 Pod。
  • Worker:真正执行训练的那些 Pod。通过调整 replicas 数量来设定分布式训练规模。
  • restartPolicy:设置为 OnFailure 时,Pod 失败后会自动重启。
  • sidecar.istio.io/inject: "false":在安装了 Istio 的集群中禁用 Sidecar 注入,使 PyTorchJob 能够正确运行。这是官方文档明确推荐的配置。

6.3 Job 状态监控

kubectl get pytorchjobs -n ml-training
kubectl describe pytorchjob pytorch-ddp-mnist -n ml-training
kubectl get pods -l training.kubeflow.org/job-name=pytorch-ddp-mnist -n ml-training

PyTorchJob 的状态可以在 conditions 字段中确认,按 Created -> Running -> Succeeded 的顺序迁移。


7. MPIJob、TFJob 支持

7.1 MPIJob

MPIJob 是用于基于 MPI(Message Passing Interface)的分布式训练的 CRD。它主要用于 Horovod 这类基于 allreduce 的框架。

apiVersion: kubeflow.org/v2beta1
kind: MPIJob
metadata:
  name: horovod-training
spec:
  slotsPerWorker: 2
  mpiReplicaSpecs:
    Launcher:
      replicas: 1
      template:
        spec:
          containers:
            - name: mpi-launcher
              image: horovod/horovod:latest
              command:
                - mpirun
                - --allow-run-as-root
                - -np
                - '8'
                - -bind-to
                - none
                - -map-by
                - slot
                - python
                - train.py
              resources:
                limits:
                  cpu: '2'
                  memory: '4Gi'
    Worker:
      replicas: 4
      template:
        spec:
          containers:
            - name: mpi-worker
              image: horovod/horovod:latest
              resources:
                limits:
                  nvidia.com/gpu: 2
                  cpu: '8'
                  memory: '32Gi'

MPIJob 的特点:

  • Launcher:执行 mpirun 命令的 Pod。它本身不执行训练,而是把任务分发给 Worker。
  • Worker:真正执行训练的 Pod。通过 slotsPerWorker 指定每个 Worker 的进程(GPU)数量。
  • 由于与框架无关,Horovod、PyTorch、TensorFlow、MXNet 等所有支持 MPI 的框架都可以使用。

7.2 TFJob

TFJob 是面向 TensorFlow 的 Parameter Server 分布式训练策略的 CRD。

apiVersion: kubeflow.org/v1
kind: TFJob
metadata:
  name: tf-ps-training
spec:
  tfReplicaSpecs:
    PS:
      replicas: 2
      template:
        spec:
          containers:
            - name: tensorflow
              image: tensorflow/tensorflow:2.14.0-gpu
              command: ['python', 'train.py']
              resources:
                limits:
                  cpu: '4'
                  memory: '16Gi'
    Worker:
      replicas: 4
      template:
        spec:
          containers:
            - name: tensorflow
              image: tensorflow/tensorflow:2.14.0-gpu
              command: ['python', 'train.py']
              resources:
                limits:
                  nvidia.com/gpu: 1
                  cpu: '8'
                  memory: '32Gi'

TFJob 的特点:

  • PS(Parameter Server):保存模型参数,接收来自各个 Worker 的梯度并更新参数。
  • Worker:处理训练数据的一部分并计算梯度,然后发送给 PS。
  • Training Operator 会自动设置 TF_CONFIG 环境变量,使 TensorFlow 的分布式策略正确运行。

8. Multi-node Multi-GPU 分布式训练配置(PyTorch DDP on K8s)

8.1 PyTorch DDP 的工作原理

PyTorch DDP(Distributed Data Parallel)是数据并行方式的分布式训练策略。它在每块 GPU 上放置一份模型副本,把 mini batch 按 GPU 数量切分后进行训练,再通过 allreduce 运算同步梯度。

在 Kubernetes 上构建 Multi-node Multi-GPU DDP 时必须理解的核心概念:

  • WORLD_SIZE:进程总数。(节点数)x(每节点 GPU 数)。例:3 个节点、每个 2 块 GPU = WORLD_SIZE 6。
  • RANK:每个进程的唯一 ID(从 0 到 WORLD_SIZE-1)。
  • LOCAL_RANK:节点内部的进程 ID(从 0 到每节点 GPU 数-1)。
  • MASTER_ADDR:Master 节点的地址。由 Training Operator 自动设置。
  • MASTER_PORT:Master 节点的通信端口。由 Training Operator 自动设置。

8.2 训练代码示例

使用 Training Operator 时,训练代码本身使用 PyTorch 原生的 Distributed API:

import torch
import torch.distributed as dist
from torch.nn.parallel import DistributedDataParallel as DDP

def main():
    # Training Operator 会自动设置环境变量,
    # 因此 torchrun 会自动处理 init_process_group。
    dist.init_process_group(backend="nccl")

    local_rank = int(os.environ["LOCAL_RANK"])
    torch.cuda.set_device(local_rank)

    model = MyModel().cuda(local_rank)
    model = DDP(model, device_ids=[local_rank])

    # 用 DistributedSampler 按 GPU 数量切分数据
    train_sampler = torch.utils.data.distributed.DistributedSampler(train_dataset)
    train_loader = torch.utils.data.DataLoader(
        train_dataset,
        batch_size=64,
        sampler=train_sampler
    )

    for epoch in range(num_epochs):
        train_sampler.set_epoch(epoch)
        for batch in train_loader:
            loss = model(batch)
            loss.backward()
            optimizer.step()
            optimizer.zero_grad()

        # 检查点仅在 rank 0 上保存
        if dist.get_rank() == 0:
            torch.save({
                'epoch': epoch,
                'model_state_dict': model.module.state_dict(),
                'optimizer_state_dict': optimizer.state_dict(),
            }, f'/checkpoints/checkpoint_epoch_{epoch}.pt')

    dist.destroy_process_group()

8.3 NCCL 通信优化

在 Multi-node GPU 通信中,NCCL(NVIDIA Collective Communication Library)的性能会直接影响训练速度。Kubernetes 环境下的 NCCL 优化要点:

env:
  - name: NCCL_DEBUG
    value: 'INFO' # 调试时查看 NCCL 通信日志
  - name: NCCL_IB_DISABLE
    value: '0' # 可使用 InfiniBand 时启用
  - name: NCCL_SOCKET_IFNAME
    value: 'eth0' # 指定通信所用的网络接口
  - name: NCCL_P2P_LEVEL
    value: 'NVL' # 使用 NVLink 时设置 P2P 通信级别

在拥有 InfiniBand 或 RoCE 之类 RDMA 网络的环境中,也可以给 Pod 设置 hostNetwork: true 来最大化网络性能。


9. 用 PVC/PV 管理训练数据与检查点

在分布式训练中,存储设计是训练数据访问、检查点保存、模型产物管理的核心。

9.1 训练数据用 PVC

训练数据需要被所有 Worker 读取,因此需要 ReadOnlyMany(ROX)ReadWriteMany(RWX) Access Mode。

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: training-data-pvc
  namespace: ml-training
spec:
  accessModes:
    - ReadOnlyMany
  storageClassName: nfs-storage
  resources:
    requests:
      storage: 500Gi

9.2 检查点用 PVC

检查点只由 Rank 0 进程保存,因此 ReadWriteOnce(RWO) 也够用,但若考虑 Fault Tolerance,ReadWriteMany(RWX) 更加灵活。

apiVersion: v1
kind: PersistentVolumeClaim
metadata:
  name: checkpoints-pvc
  namespace: ml-training
spec:
  accessModes:
    - ReadWriteMany
  storageClassName: fast-ssd
  resources:
    requests:
      storage: 100Gi

9.3 存储后端选择

存储后端Access Mode性能特性适合的用途
NFSRWX/ROX中等带宽、兼容性高训练数据共享
CephFSRWX高带宽、分布式存储大规模数据集
Lustre / GPFSRWX极高带宽HPC 级别 I/O
Local SSDRWO最高 IOPS检查点、本地缓存
S3 (CSI)RWX高带宽、容量无上限基于对象存储的数据集

对于大规模训练数据集,直接在 NFS 或 CephFS 上访问可能出现 I/O 瓶颈。此时可以考虑使用在训练开始前把数据复制到 Local SSD 的 Init Container 模式,或者通过 CSI Driver 挂载 S3 兼容对象存储。


10. Kueue:下一代 Job 排队系统

Kueue 是在 Kubernetes SIG(Special Interest Group)之下开发的云原生 Job 排队系统。它为 Batch、HPC、AI/ML 工作负载提供 Job 级别的 Admission Control。

10.1 Kueue 与 Volcano 的对比

Kueue 和 Volcano 工作在不同的层次上。

  • Volcano:Scheduler 层。决定把 Pod 放到哪个 Node 上。它替代 kube-scheduler。
  • Kueue:Admission Control 层。决定 Job 何时可以开始(admit)。它与 kube-scheduler 协同工作,并不替代调度器。

Kueue 决定“Job 是否应该开始”,而 Volcano(或 kube-scheduler)决定“已开始的 Job 的 Pod 要放在哪里”。两套系统可以互补使用。

10.2 Kueue 的核心概念

Kueue 的 Admission 过程如下:

  1. 用户创建 Job 后,Kueue 将其转换为 Workload
  2. Workload 被提交到 LocalQueue
  3. LocalQueue 引用 ClusterQueue
  4. ClusterQueue 在 ResourceFlavor 定义的资源池中检查 Quota。
  5. Quota 充足时就 Admit 该 Workload,实际的 Pod 随之被创建。

11. Kueue 的 ResourceFlavor、ClusterQueue、LocalQueue

11.1 ResourceFlavor

ResourceFlavor 定义集群中可用资源的类型。它用于区分节点的特性(价格、架构、可用性等)。

apiVersion: kueue.x-k8s.io/v1beta1
kind: ResourceFlavor
metadata:
  name: gpu-a100
spec:
  nodeLabels:
    cloud.google.com/gke-accelerator: nvidia-tesla-a100
---
apiVersion: kueue.x-k8s.io/v1beta1
kind: ResourceFlavor
metadata:
  name: gpu-t4-spot
spec:
  nodeLabels:
    cloud.google.com/gke-accelerator: nvidia-tesla-t4
    cloud.google.com/gke-provisioning: spot
  tolerations:
    - key: cloud.google.com/gke-spot
      operator: Equal
      value: 'true'
      effect: NoSchedule

在上面的示例中,gpu-a100 表示 On-demand A100 GPU 节点,gpu-t4-spot 表示 Spot T4 GPU 节点。通过 Toleration 使 Pod 能够被调度到 Spot 节点上。

11.2 ClusterQueue

ClusterQueue 定义集群范围的资源池,并管理 Quota 与 Fair Sharing 规则。

apiVersion: kueue.x-k8s.io/v1beta1
kind: ClusterQueue
metadata:
  name: ml-cluster-queue
spec:
  cohort: ml-team-cohort
  namespaceSelector: {}
  preemption:
    withinClusterQueue: LowerPriority
    reclaimWithinCohort: LowerPriority
    borrowWithinCohort:
      policy: LowerPriority
      maxPriorityThreshold: 100
  resourceGroups:
    - coveredResources: ['cpu', 'memory', 'nvidia.com/gpu']
      flavors:
        - name: gpu-a100
          resources:
            - name: 'cpu'
              nominalQuota: 64
            - name: 'memory'
              nominalQuota: 256Gi
            - name: 'nvidia.com/gpu'
              nominalQuota: 8
              borrowingLimit: 4
              lendingLimit: 2
        - name: gpu-t4-spot
          resources:
            - name: 'cpu'
              nominalQuota: 128
            - name: 'memory'
              nominalQuota: 512Gi
            - name: 'nvidia.com/gpu'
              nominalQuota: 16

核心字段说明:

  • cohort:属于同一个 Cohort 的多个 ClusterQueue 之间可以互相借用空闲 Quota。
  • nominalQuota:默认得到保障的资源量。
  • borrowingLimit:可以从 Cohort 内其他 ClusterQueue 借用的最大资源量。
  • lendingLimit:可以借给其他 ClusterQueue 的最大资源量。
  • preemption.withinClusterQueue:是否可以在同一个 ClusterQueue 内 Preempt 优先级更低的 Workload。设置为 LowerPriority 时,高优先级 Workload 可以中断低优先级 Workload。
  • preemption.reclaimWithinCohort:是否可以回收(Preempt)被 Cohort 内其他 ClusterQueue 借走的资源。
  • flavorFungibility:当存在多个 ResourceFlavor 时,决定 Borrowing 与 Preemption 的优先顺序。

11.3 LocalQueue

LocalQueue 是用户在命名空间范围内提交 Job 的入口。

apiVersion: kueue.x-k8s.io/v1beta1
kind: LocalQueue
metadata:
  name: ml-research-queue
  namespace: ml-training
spec:
  clusterQueue: ml-cluster-queue

用户通过给 Job 添加 kueue.x-k8s.io/queue-name 标签来提交到 LocalQueue:

apiVersion: kubeflow.org/v1
kind: PyTorchJob
metadata:
  name: my-training-job
  namespace: ml-training
  labels:
    kueue.x-k8s.io/queue-name: ml-research-queue
spec:
  # ... PyTorchJob spec

11.4 Kueue 的排队策略

Kueue 支持两种排队策略:

  • StrictFIFO:严格的先进先出。若队列前端的 Workload 无法被 Admit,后面的 Workload 也不会被 Admit。
  • BestEffortFIFO:即使前面的 Workload 因资源不足无法被 Admit,只要后面的 Workload 资源需求较小,也可以先被 Admit。

12. 在 Spot/Preemptible Instance 上训练(Fault Tolerance)

12.1 Spot Instance 的优点与风险

云上的 Spot Instance(AWS)、Preemptible VM(GCP)、Spot VM(Azure)相比 On-demand 可以节省高达 60~90% 的成本。AI 训练的成本大部分来自 GPU 计算费用,因此利用 Spot Instance 非常有吸引力。

但 Spot Instance 随时可能被云服务商回收(Preemption),通常只提供 30 秒到 2 分钟的短暂事先通知。若训练中的 Job 没有检查点就被 Preempt,全部进度都会丢失。

12.2 Fault-tolerant 训练策略

基于 Checkpoint 的恢复

最基本的策略是周期性保存 Checkpoint。Rank 0 进程按一定周期(以 epoch 或 step 为单位)把模型参数、优化器状态、训练进度保存到持久化存储。Pod 重启后,从最后一个 Checkpoint 继续训练。

# 在 preStop Hook 中保存紧急检查点
# 配置 Pod 的 lifecycle.preStop 调用该脚本。
import signal
import sys

def graceful_shutdown(signum, frame):
    """收到 Spot Instance 回收通知时保存紧急检查点"""
    if dist.get_rank() == 0:
        torch.save({
            'epoch': current_epoch,
            'step': current_step,
            'model_state_dict': model.module.state_dict(),
            'optimizer_state_dict': optimizer.state_dict(),
        }, '/checkpoints/emergency_checkpoint.pt')
    dist.barrier()
    sys.exit(0)

signal.signal(signal.SIGTERM, graceful_shutdown)

Kubernetes preStop Hook 配置

lifecycle:
  preStop:
    exec:
      command: ['/bin/sh', '-c', 'python save_checkpoint.py --emergency']

TorchElastic(Elastic Training)

PyTorch 的 TorchElastic(现已整合进 torchrun)是针对 Spot Instance 环境优化的分布式训练框架。

主要特点:

  • 弹性伸缩:即使 Worker 被增加或移除,训练也不会中断。
  • 自动重启:Worker 发生故障时自动恢复训练。
  • 部分执行:即使只拿到所请求 GPU 数量的一部分,也可以开始训练。
apiVersion: kubeflow.org/v1
kind: PyTorchJob
metadata:
  name: elastic-training
spec:
  elasticPolicy:
    rdzvBackend: etcd
    rdzvHost: etcd-service
    rdzvPort: 2379
    minReplicas: 2
    maxReplicas: 8
  pytorchReplicaSpecs:
    Worker:
      replicas: 4
      restartPolicy: OnFailure
      template:
        spec:
          nodeSelector:
            cloud.google.com/gke-provisioning: spot
          tolerations:
            - key: cloud.google.com/gke-spot
              operator: Equal
              value: 'true'
              effect: NoSchedule
          containers:
            - name: pytorch
              image: my-registry/elastic-training:latest
              resources:
                limits:
                  nvidia.com/gpu: 1

12.3 Hybrid 集群策略

在实际生产中会采用混合 On-demand 与 Spot 节点的 Hybrid 策略:

  • On-demand 节点:放置集群核心组件(etcd、Kueue Controller、Training Operator)。设置 Taint 以避免普通工作负载被放置到这里。
  • Spot 节点:放置真正的训练 Worker Pod。用 Kueue 的 ResourceFlavor 单独管理 Spot 节点,并选择性价比高的 GPU。

与 Kueue 的 Preemption 策略结合后,Spot 节点被回收时,Kueue 可以自动把 Workload 重新放置到其他可用资源上。


13. References

本文参考了以下官方文档与资料撰写。

Volcano 官方文档

Kubeflow Training Operator 官方文档

Kueue 官方文档

其他参考资料

현재 단락 (1/583)

Kubernetes 的默认调度器(kube-scheduler)以 Pod 为单位执行调度。根据单个 Pod 所请求的资源(CPU、Memory、GPU)将其放置到合适的 Node 上,是它的核心职...

작성 글자: 0원문 글자: 19,556작성 단락: 0/583