5.2 K8s 训练 Operator 与 Ray Train:分布式训练的两套方案


文档摘要

5.2 K8s 训练 Operator 与 Ray Train:分布式训练的两套方案 第 2 章 2.4 节介绍了 Operator 范式与 Training Operator 全家桶,第 5.1 节讲了分布式训练的算法原理。本节回答工程问题:在 K8s 上具体怎么跑分布式训练?答案是两条路线——K8s 原生的 Training Operator,与灵活的 Ray Train。 5.2.1 训练编排的两条哲学路线 分布式训练要把一组互相绑定的 Pod(master + workers)编排起来,处理启动顺序、环境变量注入、拓扑感知、失败重试等复杂逻辑。

5.2 K8s 训练 Operator 与 Ray Train:分布式训练的两套方案

第 2 章 2.4 节介绍了 Operator 范式与 Training Operator 全家桶,第 5.1 节讲了分布式训练的算法原理。本节回答工程问题:在 K8s 上具体怎么跑分布式训练?答案是两条路线——K8s 原生的 Training Operator,与灵活的 Ray Train。

5.2.1 训练编排的两条哲学路线

分布式训练要把一组互相绑定的 Pod(master + workers)编排起来,处理启动顺序、环境变量注入、拓扑感知、失败重试等复杂逻辑。K8s 生态给出了两条主要路线,反映了两种不同的工程哲学:

路线 代表 哲学 适合
声明式 CRD Training Operator(PyTorchJob 等) K8s 原生,声明式,集成深 K8s 优先团队
编程式 Actor Ray Train(基于 Ray) Python 优先,灵活,跨语言 Python 优先团队

这两条路线不是互斥的——它们甚至可以协作(Ray Cluster 跑在 K8s Operator 之上)。理解二者的差异,是选型的前提。

5.2.2 Training Operator:声明式 CRD 范式

回顾第 2 章 2.4 节,Training Operator(Kubeflow 子项目)提供了一系列训练 CRD:PyTorchJob、MPIJob、DeepSpeedJob、TFJob、XGBoostJob、MXNetJob。每个 CRD 都遵循「声明 ReplicaSpecs → Controller 创建 Pod 群 → 注入框架环境变量 → 监控状态」的模式。

训练 Operator 的工作流程

以 PyTorchJob 为例:

Training Operator 的核心价值:

  1. 自动注入框架约定:PyTorchJob 自动注入 WORLD_SIZERANKMASTER_ADDRMASTER_PORT,开发者无需手工配置。
  2. Gang 启动语义:等所有 Pod 就绪后再启动训练(配合 Gang Scheduling,见第 3 章)。
  3. 失败重试:worker Pod 失败时按策略重启,master 失败通常整体重启。
  4. 状态机管理:CR 状态自动流转(Created → Running → Succeeded/Failed)。
  5. 与 K8s 生态深度集成:原生兼容 Kueue 排队、Prometheus 监控、节点亲和等。

各种 Job 类型的辨析

Job 类型 框架 启动器 典型用例
PyTorchJob PyTorch DDP torchrun PyTorch 数据并行
MPIJob Horovod / MPI mpirun + sshd 跨框架 AllReduce
DeepSpeedJob DeepSpeed deepspeed launcher ZeRO 大模型训练
TFJob TensorFlow tf.distributed TF 训练(逐渐式微)
XGBoostJob XGBoost xgboost rabit 分布式树模型

💡 选型经验:新 PyTorch 项目 → PyTorchJob;用 ZeRO 训大模型 → DeepSpeedJob;老 Horovod 项目 → MPIJob;分布式 XGBoost → XGBoostJob。三者的核心区别在于「启动器」与「环境变量约定」,底层都是把训练脚本跑在 K8s 编排的 Pod 群上。

5.2.3 Ray Train:基于 Actor 的编程式范式

Ray Train 是基于 Ray 框架的分布式训练库。理解 Ray Train 要先理解 Ray:

Ray 是一个通用的分布式计算框架,核心抽象是 Actor(常驻有状态对象)Task(无状态函数)。Ray 让你能用纯 Python 写分布式程序,由 Ray 的调度器自动把 Actor/Task 分布到集群的多个节点。

# Ray 风格伪代码(示意,非完整可运行) import ray from ray.train import Trainer trainer = Trainer(num_workers=8, use_gpu=True) trainer.start() results = trainer.run(train_func) trainer.shutdown()

Ray Train 把训练抽象为「启动 N 个 worker Actor,每个 Actor 跑训练函数」。它的核心特点:

  • Python 优先:完全用 Python API 编写,无 YAML。
  • Actor 模型:worker 是常驻 Actor,可在运行时动态调整数量(与弹性训练契合)。
  • 跨语言:Ray 支持 Python、Java、C++,可跨语言组合。
  • 与 Ray 生态融合:与 Ray Serve(推理)、Ray Tune(超参搜索)、Ray Data(数据处理)无缝集成。

Ray Cluster 在 K8s 上的部署

Ray Cluster 本身可以通过 Ray Operator(K8s Operator)部署——用 RayCluster CRD 声明 Ray 集群的 head 与 worker 节点,Operator 自动管理。所以「Ray Train 跑在 K8s 上」的完整路径是:

K8s → Ray Operator → RayCluster CR → Ray head + workers → Ray Train → 训练

这种方式让 Ray 既能利用 K8s 的资源管理(GPU 调度、弹性、健康检查),又保留了 Ray 的编程式灵活性。

5.2.4 两套方案的对比

把 Training Operator 与 Ray Train 放在一起对比:

维度 Training Operator Ray Train
范式 声明式 CRD(YAML) 编程式 Actor(Python)
集成深度 K8s 原生 K8s 上跑 Ray Cluster
灵活性 中(按 CRD schema) 高(Python 自由编程)
学习曲线 低(懂 K8s 即可) 中(要懂 Ray 抽象)
跨语言 否(主要 Python) 是(Python/Java/C++)
弹性训练 借助 Torchrun Elastic 原生 Actor 动态扩缩
生态 训练为主 训练+推理+数据+超参一体化
多框架 PyTorch/TF/MPI/XGBoost 主要 PyTorch/HuggingFace
适合 单一训练负载、K8s 优先 多负载一体化、Python 优先

5.2.5 选型决策

何时选 Training Operator,何时选 Ray Train?给一个工程决策树:

💡 判读:选型的核心问题是「你的负载单一还是多元」。如果一个集群只跑训练,Training Operator 更轻量直接;如果一个集群要跑训练 + 推理 + 数据处理 + 超参搜索(如大模型团队的「ML 平台」),Ray 的统一抽象优势明显,用一套 API 管多种负载。这也是为什么 OpenAI、Uber、Shopify 等团队都采用 Ray 作为 ML 平台底座。

5.2.6 与第 3 章调度器的协作

无论选哪套方案,训练编排都要与第 3 章的批处理调度器协作:

  • Training Operator + Kueue:PyTorchJob 提交时声明 LocalQueue,Kueue 按配额排队、Gang 准入。这是 K8s 官方推荐组合。
  • Training Operator + Volcano:PodGroup 标签让 Volcano 做 Gang Scheduling。
  • Ray + Kueue:Ray Cluster 的 worker Pod 也可走 Kueue 排队,但 Ray 自身的 Actor 调度与 K8s 调度交互较复杂,常需自定义集成。

这种「训练编排 + 调度排队」的分层架构,让训练作业的「生命周期」与「何时调度」分离,是云原生 AI 训练栈的设计精髓。

5.2.7 检查点与失败恢复的工程实践

无论用哪套方案,长时训练的检查点与失败恢复都是核心议题(回顾第 3 章 3.4 节)。两套方案的检查点实践略有不同:

Training Operator 的检查点

  • 训练脚本内调 torch.save(model.state_dict(), path)
  • path 指向 PVC 或对象存储(跨节点可读)。
  • 配合 CRD 的 restartPolicy,失败重启后从检查点恢复。
  • 弹性训练用 Torchrun Elastic + rendezvous。

Ray Train 的检查点

  • Ray Train 提供 ray.train.save_checkpoint() API。
  • 检查点自动同步到外部存储。
  • Ray 的容错机制(Actor 重启)与训练恢复结合。

两套方案都建议把检查点写到对象存储(S3/OSS)而非本地盘,避免节点故障导致检查点丢失。同时保留多个历史检查点,支持回滚到任意保存点。

5.2.8 训练监控与可观测性

训练过程的可观测性也很重要:

  • 资源利用率:GPU 利用率、显存占用、网络吞吐(通过 DCGM Exporter,详见第 8 章)。
  • 训练指标:loss、learning rate、梯度范数(通过 MLflow/W&B 自动上报)。
  • 吞吐指标:samples/s、tokens/s(衡量训练效率)。
  • 通信开销:AllReduce 耗时占比(判断是否通信瓶颈)。

一个健康的训练,GPU 利用率应稳定在 80%+,loss 平稳下降,AllReduce 通信占比合理(<30%)。如果 GPU 利用率低、通信占比高,往往说明并行策略有问题(如 TP 跨节点了)。第 8 章会详细讲这些指标的监控。

⚠️ 现实坑:分布式训练最常见的工程坑是「GPU 利用率低但找不到原因」。可能的原因有:数据加载慢(CPU/IO 瓶颈)、通信开销大(拓扑差或并行策略错)、检查点保存阻塞(同步保存)、Python GIL(数据预处理单线程)。要靠监控指标逐项排查,详见第 8 章。

5.2.9 训练之外:Ray 的全栈价值

最后强调一下 Ray 在云原生 AI 中的全栈价值。Ray 不只是训练框架,它是一个统一的分布式计算平台:

  • Ray Data:分布式数据处理与预处理。
  • Ray Train:分布式训练。
  • Ray Tune:分布式超参搜索。
  • Ray Serve:推理服务(与第 6 章 KServe 等可类比)。
  • RLlib:分布式强化学习。

用一个统一的 Ray Cluster 跑通数据→训练→推理的全链路,是 Ray 的核心卖点。这也是为什么许多 ML 平台团队选 Ray——它避免了「数据用 Spark、训练用 PyTorchJob、推理用 KServe、超参用 Optuna」的多语言割裂。

但 Ray 的代价是「学习曲线与运维复杂度」——Ray 自己是一套独立的分布式系统,要在 K8s 上运维好 Ray Cluster,需要专门的 Ray 知识。团队若不准备深度投入 Ray,Training Operator 反而更简单。

本节小结

  • K8s 上跑分布式训练有两条路线:声明式 CRD(Training Operator)与编程式 Actor(Ray Train)。
  • Training Operator 是 K8s 原生 CRD 范式,自动注入框架约定、Gang 启动、状态机管理,与 K8s 生态深度集成。
  • Ray Train 基于 Ray Actor 模型,Python 优先、灵活、跨语言,与 Ray Serve/Tune/Data 一体化。
  • PyTorchJob/DeepSpeedJob/MPIJob 的核心区别在于启动器与环境变量约定,底层都是 K8s 编排 Pod 群。
  • 选型核心:单一训练负载 + K8s 优先 → Training Operator;多元负载一体化 + Python 优先 → Ray Train。
  • 无论哪套方案都要与 Kueue/Volcano 协作做排队 Gang 准入,分层架构是设计精髓。
  • 检查点要写对象存储(S3/OSS),保留多个历史版本支持回滚。
  • 训练监控看 GPU 利用率(>80%)、loss、吞吐、AllReduce 占比,GPU 利用率低要逐项排查。
  • Ray 是统一分布式平台(Data/Train/Tune/Serve/RLlib),全栈价值大但学习与运维成本也高。

下一节《5.3 Kubeflow Pipelines 的 DAG 编排》将讲清如何把训练组织成可版本、可重跑的流水线。


发布者: 作者: 灏天文库 转发
评论区 (0)
U