- 1. Kubernetes 默认调度器的局限
- 2. Volcano Scheduler 官方文档分析
- 3. Volcano 的 Gang Scheduling 与 Fair-share 调度策略
- 4. Volcano Job CRD 与插件系统
- 5. Kubeflow Training Operator 官方文档分析
- 6. PyTorchJob CRD 详解与 YAML 示例
- 7. MPIJob、TFJob 支持
- 8. Multi-node Multi-GPU 分布式训练配置(PyTorch DDP on K8s)
- 9. 用 PVC/PV 管理训练数据与检查点
- 10. Kueue:下一代 Job 排队系统
- 11. Kueue 的 ResourceFlavor、ClusterQueue、LocalQueue
- 12. 在 Spot/Preemptible Instance 上训练(Fault Tolerance)
- 13. References
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 | 功能 | 适合的工作负载 |
|---|---|---|
| Gang | All or Nothing 调度 | 分布式训练、MPI、Spark |
| Binpack | 尽量填满已有节点后再放置 | 小规模 Job、最大化资源效率 |
| DRF | 基于 Dominant Resource 的公平共享 | 多租户批处理 |
| Proportion | 管理各 Queue 的资源比例 | 按团队做资源隔离 |
| Priority | 基于优先级的调度 | 优先级各异的多个 Job |
| Task-Topology | Task 间 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 会自动处理以下事项:
- 自动设置环境变量:把
WORLD_SIZE、RANK、MASTER_ADDR、MASTER_PORT等分布式训练所需的环境变量自动注入 Pod。 - torchrun CLI 配置:配置环境使 PyTorch 的
torchrun(此前的torch.distributed.launch)能够正确运行。 - Job 状态管理:管理 Created -> Running -> Succeeded/Failed 等生命周期,并在失败时应用重启策略。
- 创建 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 | 性能特性 | 适合的用途 |
|---|---|---|---|
| NFS | RWX/ROX | 中等带宽、兼容性高 | 训练数据共享 |
| CephFS | RWX | 高带宽、分布式存储 | 大规模数据集 |
| Lustre / GPFS | RWX | 极高带宽 | HPC 级别 I/O |
| Local SSD | RWO | 最高 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 过程如下:
- 用户创建 Job 后,Kueue 将其转换为 Workload。
- Workload 被提交到 LocalQueue。
- LocalQueue 引用 ClusterQueue。
- ClusterQueue 在 ResourceFlavor 定义的资源池中检查 Quota。
- 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 官方文档
- Volcano 官方网站
- Volcano Introduction
- Volcano PodGroup
- Volcano Plugins
- Volcano Unified Scheduling
- Volcano GitHub Repository
Kubeflow Training Operator 官方文档
- Kubeflow Trainer Overview
- PyTorchJob (Legacy V1)
- Distributed Training Reference
- MPIJob
- TFJob
- Migrating to Trainer V2
- Kubeflow Trainer GitHub Repository
- Volcano Integration for Job Scheduling
Kueue 官方文档
- Kueue 官方网站
- Kueue Overview
- Kueue Concepts
- ClusterQueue
- ResourceFlavor
- Run a PyTorchJob with Kueue
- Kueue GitHub Repository
其他参考资料
현재 단락 (1/583)
Kubernetes 的默认调度器(kube-scheduler)以 Pod 为单位执行调度。根据单个 Pod 所请求的资源(CPU、Memory、GPU)将其放置到合适的 Node 上,是它的核心职...