训练调度与资源管理

训练调度与资源管理

在一个拥有 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 训练的主流选择

%%{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

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
TipSlurm vs K8s 的选择
  • 纯 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
            }
          ]
        }
      ]
    }
TipAI 训练推荐 Bin Packing

对于 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"
Warning抢占的代价

抢占一个正在运行的训练任务意味着:丢失自上次 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

Warning弹性训练的隐藏成本
  1. Batch Size 变化:world_size 改变会导致有效 batch size 变化,影响学习率调度
  2. Sampler 重置:数据采样需要重新分配,可能重复处理部分数据
  3. Checkpoint 频率:需要更频繁的 checkpoint 以减少恢复损失
  4. 收敛保证:目前没有严格证明弹性训练能达到与静态训练相同的最终精度

对于关键的生产训练任务,弹性训练是”保险”,不是”常规操作”。

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。一个完整的训练任务包括:

  1. 数据准备:分片、预编码、上传到共享存储
  2. 环境准备:镜像构建、依赖安装
  3. 训练:主训练循环
  4. 评估:在验证集上跑指标
  5. 保存:模型 artifact 上传
  6. 通知:通知下游任务或团队
# 使用 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)
Tip关键监控指标

训练任务必须监控的 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)