%%{init: {'theme': 'base'}}%%
graph TB
subgraph "Kubernetes + Volcano 架构"
API["K8s API Server"] --> VC["Volcano Controller"]
API --> VS["Volcano Scheduler"]
API --> VLA["Volcano Webhook"]
VC --> JOB["Volcano Job (vcjob)"]
VS --> GANG["Gang Scheduling"]
VS --> BIN["Bin Packing"]
JOB --> POD0["Pod (GPU 0-7)<br/>Rank 0"]
JOB --> POD1["Pod (GPU 8-15)<br/>Rank 1"]
JOB --> POD2["Pod (GPU 16-23)<br/>Rank 2"]
GANG -.->|全部就绪才启动| POD0
GANG -.->|全部就绪才启动| POD1
GANG -.->|全部就绪才启动| POD2
end
训练调度与资源管理
训练调度与资源管理
在一个拥有 1000 张 GPU 的集群上,1% 的利用率差距就是每年数百万美元的浪费。
章节导入:调度器是 AI 基础设施的”操作系统”
当你的团队只有 8 张 GPU 时,谁先用谁后用可以靠喊一嗓子解决。当你有 1000 张 GPU、5 个团队、20 个训练任务同时排队时,你需要一个调度器。
调度器解决的核心问题是:在有限资源下,最大化集群吞吐量和任务公平性。这听起来像操作系统课程的经典问题——事实上,AI 训练调度确实借鉴了很多传统操作系统的设计,但它有自己的独特挑战:任务之间有严格的时序依赖(必须同时启动所有 rank)、失败恢复代价高昂、资源需求差异巨大。
6.1 大规模集群调度器设计
三种主流调度器
| 调度器 | 起源 | 核心特点 | 适用场景 |
|---|---|---|---|
| Kubernetes | Google (Borg) | 容器化、生态丰富 | 通用云原生、混合负载 |
| Slurm | HPC 社区 | 批处理作业、GPU 原生支持 | 学术/研究集群 |
| Volcano | 华为 (CNCF) | K8s 上的批处理调度 | AI 训练专用 K8s 场景 |
Kubernetes + Volcano:AI 训练的主流选择
Volcano Job 定义
apiVersion: batch.volcano.sh/v1alpha1
kind: Job
metadata:
name: llm-pretraining
labels:
priority-class: high
spec:
minAvailable: 4 # Gang: 必须同时有 4 个 Pod 就绪
schedulerName: volcano # 使用 Volcano 调度器
policies:
- event: PodEvicted
action: RestartJob # Pod 被驱逐时重启整个任务
- event: TaskFailed
action: RestartJob
maxRetry: 3
plugins:
env: # 自动注入 RANK 环境变量
svc: # 自动创建 Headless Service
tasks:
- replicas: 4
name: worker
template:
spec:
schedulerName: volcano
containers:
- name: pytorch
image: pytorch/pytorch:2.4.0-cuda12.4
resources:
limits:
nvidia.com/gpu: 8
rdma/rdma-shared: 1
env:
- name: WORLD_SIZE
value: "32"
- name: NNODES
value: "4"
command:
- torchrun
- --nproc_per_node=8
- --nnodes=4
- --rdzv_backend=c10d
- --rdzv_endpoint=llm-pretraining:29500
- train.py
hostNetwork: true # 高性能网络
hostIPC: true # 共享内存(NCCL 需要)Slurm:HPC 世界的王者
#!/bin/bash
#SBATCH --job-name=llm-pretrain
#SBATCH --partition=gpu-h100
#SBATCH --nodes=4
#SBATCH --gres=gpu:h100:8
#SBATCH --cpus-per-task=64
#SBATCH --exclusive # 独占节点
#SBATCH --time=72:00:00
# Slurm 自动设置环境变量
echo "Job ID: $SLURM_JOB_ID"
echo "Nodes: $SLURM_JOB_NUM_NODES"
echo "Node list: $SLURM_JOB_NODELIST"
# 生成 hostfile
scontrol show hostnames $SLURM_JOB_NODELIST > hostfile
# 启动 PyTorch 分布式训练
torchrun \
--nproc_per_node=8 \
--nnodes=$SLURM_JOB_NUM_NODES \
--rdzv_backend=c10d \
--rdzv_endpoint=$(head -1 hostfile):29500 \
train.py --config configs/70b.yaml- 纯 AI 训练集群 → Slurm 更简单、GPU 管理更原生
- 混合负载(训练+推理+服务) → K8s + Volcano 更灵活
- 需要多租户隔离和弹性 → K8s 生态更成熟
- 已有 HPC 基础设施 → 保持 Slurm,不要为了”云原生”而迁移
6.2 GPU 资源调度策略
Gang Scheduling:AI 训练的刚需
分布式训练任务有一个硬约束:所有 rank 必须同时启动。如果 32 个 Pod 中只有 31 个获得 GPU,任务无法启动,那 31 个 Pod 占用的 GPU 就是浪费——这是死锁的经典场景。
%%{init: {'theme': 'base'}}%%
graph TB
subgraph "没有 Gang Scheduling"
A1["Job A: 需要 4 GPU"] --> A2["获得 3 GPU"]
A2 --> A3["❌ 无法启动<br/>浪费 3 GPU"]
B1["Job B: 需要 4 GPU"] --> B2["获得 1 GPU"]
B2 --> B3["❌ 无法启动<br/>浪费 1 GPU"]
C1["集群:8 GPU 总量"]
C1 --> A2
C1 --> B2
C2["利用率:0%"]
end
subgraph "有 Gang Scheduling"
D1["Job A: 需要 4 GPU"] --> D2{"4 GPU 可用?"}
D2 -->|是| D3["✅ 全部启动"]
D2 -->|否| D4["⏳ 等待"]
E1["Job B: 需要 4 GPU"] --> E2{"4 GPU 可用?"}
E2 -->|是| E3["✅ 全部启动"]
E2 -->|否| E4["⏳ 等待"]
end
Volcano 的 Gang Scheduling 通过 PodGroup 概念实现:
apiVersion: scheduling.volcano.sh/v1beta1
kind: PodGroup
metadata:
name: llm-training-pg
spec:
minMember: 4 # 至少 4 个 Pod 同时调度
priorityClassName: high
queue: training-queue # 指定队列Bin Packing vs Spreading
两种截然不同的资源分配哲学:
| 策略 | 做法 | 优点 | 缺点 |
|---|---|---|---|
| Bin Packing | 尽量填满一个节点再开下一个 | 空闲节点可休眠省电;通信延迟低(同节点) | 单点故障影响大 |
| Spreading | 尽量分散到不同节点 | 容错性好;负载均衡 | 跨节点通信多;无空闲节点 |
# Volcano 调度策略配置
apiVersion: scheduling.volcano.sh/v1beta1
kind: Configuration
data:
scheduling.sh: |
{
"cache": {
"type": "redis"
},
"actions": "allocate, backfill, preempt, reclaim",
"tiers": [
{
"plugins": [
{
"name": "priority"
},
{
"name": "gang",
"enablePreemptable": true
},
{
"name": "binpack", # 使用 Bin Packing
"arguments": {
"binpack.policy": "best"
}
}
]
},
{
"plugins": [
{
"name": "drf" # Dominant Resource Fairness
}
]
}
]
}对于 GPU 训练任务,Bin Packing 通常是更好的选择:同节点的 GPU 之间走 NVLink(900 GB/s),远快于跨节点的 IB(50-100 GB/s)。把一个训练任务的所有 Pod 放在尽量少的节点上,能显著减少通信开销。
拓扑感知调度
# 检查节点 GPU 拓扑
nvidia-smi topo -m
# 输出示例:
# GPU0 GPU1 GPU2 GPU3 GPU4 GPU5 GPU6 GPU7
# GPU0 X NV NV NV NV NV NV NV
# GPU1 NV X NV NV NV NV NV NV
# ...
# NV = NVLink 连接
# SYS = 跨 NUMA 节点(需要走 QPI,性能较差)理想情况下,同一个训练任务的 8 张卡应该在同一个 NUMA 节点内。
6.3 优先级、抢占与公平性
多租户公平性:DRF 算法
当多个团队共享集群时,如何公平分配资源?Dominant Resource Fairness (DRF) 是最常用的算法:
集群资源:100 GPU, 1000 CPU
Team A 任务:需要 4 GPU + 32 CPU(主导资源 = GPU)
Team B 任务:需要 1 GPU + 16 CPU(主导资源 = GPU)
Team C 任务:需要 0 GPU + 64 CPU(主导资源 = CPU)
DRF 分配:
- Team A: 主导份额 = GPU_share
- Team B: 主导份额 = GPU_share
- Team C: 主导份额 = CPU_share
目标:最小化 max(Team A 的 GPU_share, Team C 的 CPU_share)
``### K8s PriorityClass 和抢占
```yaml
# 定义优先级
apiVersion: scheduling.k8s.io/v1
kind: PriorityClass
metadata:
name: training-production
value: 1000000
globalDefault: false
preemptionPolicy: PreemptLowerPriority # 抢占低优先级任务
---
apiVersion: scheduling.k8s.io/v1
kind: PriorityClass
metadata:
name: training-experiment
value: 100000
preemptionPolicy: Never # 不抢占
Volcano Queue 配置
apiVersion: scheduling.volcano.sh/v1beta1
kind: Queue
metadata:
name: research-team
spec:
weight: 1 # 队列权重
reclaimable: true # 允许被抢占
capability:
cpu: "512"
memory: "2048Gi"
nvidia.com/gpu: "64" # GPU 配额上限
---
apiVersion: scheduling.volcano.sh/v1beta1
kind: Queue
metadata:
name: production-team
spec:
weight: 3 # 3 倍权重
reclaimable: false # 生产任务不可被抢占
capability:
cpu: "1024"
memory: "4096Gi"
nvidia.com/gpu: "128"抢占一个正在运行的训练任务意味着:丢失自上次 checkpoint 以来的所有计算、重新排队等待资源、重新加载数据。一个正在训练 70B 模型的任务被抢占,可能浪费数十万元的计算费用。生产环境务必设置 preemptionPolicy: Never 或高优先级。
6.4 弹性训练与动态扩缩容
传统分布式训练假设固定的 world_size。如果中途有节点故障或加入新节点,整个训练任务必须重启。弹性训练打破了这一限制。
PyTorch Elastic(Torchrun)
# Elastic 模式启动:允许节点动态加入和退出
torchrun \
--nproc_per_node=8 \
--nnodes=4:8 \
--rdzv_backend=c10d \
--rdzv_endpoint=master:29500 \
--max_restarts=3 \
train_elastic.py# 弹性训练代码需要处理 world_size 变化
import torch.distributed.elastic as elastic
def main():
# 每次重启都会重新调用 main()
world_size = dist.get_world_size()
rank = dist.get_rank()
# 关键:dataloader 需要根据当前 world_size 调整
sampler = DistributedSampler(
dataset,
num_replicas=world_size, # 动态!
rank=rank,
)
# 从 checkpoint 恢复
start_step = load_checkpoint_if_exists()
for step, batch in enumerate(loader, start=start_step):
loss = model(batch)
loss.backward()
# 定期保存 checkpoint
if step % 100 == 0:
save_checkpoint(model, optimizer, step)弹性训练的挑战
%%{init: {'theme': 'base'}}%%
graph LR
A["节点故障"] --> B["Rendezvous 重新协商"]
B --> C{参与节点数 >= min_nodes?}
C -->|是| D["重新分配 Rank"]
C -->|否| E["等待更多节点"]
D --> F["从最近 Checkpoint 恢复"]
F --> G["继续训练"]
E --> D
- Batch Size 变化:world_size 改变会导致有效 batch size 变化,影响学习率调度
- Sampler 重置:数据采样需要重新分配,可能重复处理部分数据
- Checkpoint 频率:需要更频繁的 checkpoint 以减少恢复损失
- 收敛保证:目前没有严格证明弹性训练能达到与静态训练相同的最终精度
对于关键的生产训练任务,弹性训练是”保险”,不是”常规操作”。
6.5 训练任务生命周期管理
任务的完整生命周期
%%{init: {'theme': 'base'}}%%
stateDiagram-v2
[*] --> Queued: 提交任务
Queued --> Scheduling: 资源匹配
Scheduling --> Pending: Gang 不满足
Pending --> Scheduling: 资源释放
Scheduling --> Initializing: 全部 Pod 就绪
Initializing --> Training: 进程启动 + 模型加载
Training --> Checkpointing: 定期保存
Checkpointing --> Training: 保存完成
Training --> Paused: 被抢占
Paused --> Queued: 重新排队
Training --> Completed: 训练完成
Training --> Failed: 不可恢复错误
Completed --> [*]
Failed --> [*]
任务编排与工作流
实际的训练不仅仅是 torchrun。一个完整的训练任务包括:
- 数据准备:分片、预编码、上传到共享存储
- 环境准备:镜像构建、依赖安装
- 训练:主训练循环
- 评估:在验证集上跑指标
- 保存:模型 artifact 上传
- 通知:通知下游任务或团队
# 使用 Argo Workflows 编排完整训练流程
apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
name: llm-training-pipeline
spec:
entrypoint: main
templates:
- name: main
steps:
- - name: prepare-data
template: data-prep
- - name: train
template: training
arguments:
parameters:
- name: data_path
value: "{{steps.prepare-data.outputs.parameters.data_path}}"
- - name: evaluate
template: eval
when: "{{steps.train.status}} == Succeeded"
- name: data-prep
container:
image: data-processor:latest
command: [python, prepare_data.py, --output, /data/processed]
outputs:
parameters:
- name: data_path
value: /data/processed
- name: training
inputs:
parameters:
- name: data_path
resource:
action: create
manifest: |
apiVersion: batch.volcano.sh/v1alpha1
kind: Job
# ... Volcano Job 定义
- name: eval
container:
image: evaluator:latest
command: [python, evaluate.py, --model, /checkpoints/latest]监控与可观测性
# 训练任务上报指标(Prometheus + Grafana)
from prometheus_client import Counter, Gauge, Histogram
# 定义指标
training_loss = Gauge('training_loss', 'Current training loss', ['model', 'run_id'])
gpu_utilization = Gauge('gpu_utilization', 'GPU utilization %', ['rank', 'gpu_id'])
grad_norm = Gauge('gradient_norm', 'Gradient norm')
tokens_per_second = Gauge('tokens_per_second', 'Training throughput')
checkpoint_save_duration = Histogram(
'checkpoint_save_duration_seconds',
'Time to save checkpoint'
)
# 在训练循环中上报
for step, batch in enumerate(loader):
loss = model(batch)
loss.backward()
optimizer.step()
if step % 10 == 0:
training_loss.labels(model="70b", run_id="run_001").set(loss.item())
grad_norm.set(calculate_grad_norm(model))
if step % 100 == 0:
tps = measure_tokens_per_second()
tokens_per_second.set(tps)训练任务必须监控的 5 个核心指标: 1. 吞吐量(tokens/sec 或 samples/sec)—— 最直接的效率指标 2. GPU 利用率——持续低于 90% 说明有瓶颈(通信、数据加载) 3. 显存使用——接近 100% 有 OOM 风险 4. Loss 曲线——异常跳跃可能意味着数据问题或梯度爆炸 5. 迭代时间分布——P99 和 P50 的差距揭示是否存在掉队者
小结
| 关注点 | 推荐方案 |
|---|---|
| 调度器 | K8s + Volcano(混合负载)/ Slurm(纯 HPC) |
| 调度策略 | Gang Scheduling + Bin Packing |
| 公平性 | DRF + PriorityClass + Queue |
| 弹性 | Torchrun Elastic(容错),不是常规操作 |
| 生命周期 | Argo Workflows / KubeFlow Pipelines |
| 监控 | Prometheus + Grafana + 自定义 exporter |
调度的本质是在资源受限下最大化整体产出。好的调度让每个团队都觉得自己拥有整个集群——坏的设计让所有人都在等待。
延伸阅读
- Borg / Omega / Kubernetes:Verma et al., “Large-scale cluster management at Google with Borg” (2015)
- Volcano:CNCF 项目文档 https://volcano.sh/en/docs/
- Gang Scheduling:Zhang et al., “Gang Scheduling for Distributed Training” (2020)
- PyTorch Elastic:PyTorch 官方文档 “Torchrun and Elastic Agent”
- DRF:Ghodsi et al., “Dominant Resource Fairness: Fair Allocation of Multiple Resource Types” (2011)
- MLSys 论文:“MLpedia: A Large-Scale Training System” (Meta, 2024)