在 Kubernetes 上使用 RayJob 分布式训练 Fashion MNIST:PyTorch + Ray Train 端到端实战指南
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
导读
本文基于 Ray 官方文档中的 MNIST 训练示例,完整演示如何利用 KubeRay 的 RayJob 自定义资源,在 Kubernetes 集群上以 CPU 资源端到端运行 PyTorch 模型的分布式训练任务。你将学会创建 Kind 本地集群、安装 KubeRay operator、编写并提交 RayJob、核对 worker/head/submitter Pod 状态、阅读训练日志与结果,以及根据机器资源正确配置replicas、NUM_WORKERS、CPUS_PER_WORKER等关键参数,最终掌握"提交即训练、训练完即查看结果"的云原生 AI 工作负载落地路径。
背景:为什么用 RayJob 跑分布式训练
本示例属于 Ray 官方"在 Kubernetes 上跑 Ray 训练任务"系列(示例源码)。它把两个层面的内容串在一起:
- 训练本身:用 Ray Train 的
TorchTrainer对 Fashion MNIST 数据集做多 worker 分布式训练; - 集群编排:用 KubeRay 的 RayJob 自定义资源,让 KubeRay operator 自动创建 RayCluster、等待集群就绪后自动提交 Ray job,训练结束后由你决定是否回收集群。
在开始之前,需要区分三个容易混淆的概念(详见 RayJob Quickstart):
- RayJob:KubeRay 提供的 Kubernetes 自定义资源(CRD),统管"集群创建"与"作业提交"两件事;
- Ray job:一个打包好的 Ray 应用,可以被提交到远端 Ray 集群执行;
- Submitter(提交器):一个 Kubernetes Job,负责执行
ray job submit把 Ray job 提交到 RayCluster。
RayJob 的价值在于:你只需描述"训练入口命令 + 需要的 worker 数量 + 运行环境",KubeRay operator 会替你完成集群的拉起与作业的投递,无需手动先建 RayCluster 再单独提交任务。
Step 1:创建 Kubernetes 集群
本示例使用 Kind 在本地创建一个单节点 Kubernetes 集群。如果你已有可用的 Kubernetes 集群,可以跳过这一步:
kind create cluster --image=kindest/node:v1.26.0说明:Kind 适合快速验证与本地开发;生产环境请使用托管的 Kubernetes 服务(如 EKS、GKE、ACK 等)或自建集群,命令与本文保持一致。
Step 2:安装 KubeRay operator
KubeRay operator 负责监听 RayJob / RayCluster 等自定义资源并执行对应生命周期操作。官方推荐使用 Helm 安装(详见 KubeRay Operator Installation),也可以使用 Kustomize:
# 方法一:Helm(推荐) helm repo add kuberay https://ray-project.github.io/kuberay-helm/ helm repo update kubectl create namespace ray-system helm install kuberay-operator kuberay/kuberay-operator --version 1.7.0 -n ray-system # 方法二:Kustomize # kubectl create namespace ray-system # kubectl create -k "github.com/ray-project/kuberay/ray-operator/config/default?ref=v1.7.0" -n ray-system安装完成后验证 operator 运行状态:
kubectl get pods -n ray-system # NAME READY STATUS RESTARTS AGE # kuberay-operator-6bc45dd644-gwtqv 1/1 Running 0 24sStep 3:创建 RayJob
RayJob 由两部分组成:一个 RayCluster 自定义资源(描述 head/worker Pod 规格),以及一个可提交到该集群的 Ray job。KubeRay 会在 RayCluster 就绪后自动提交 job。
首先下载示例 YAML 文件:
# 下载 `ray-job.pytorch-mnist.yaml` curl -LO https://raw.githubusercontent.com/ray-project/kuberay/master/ray-operator/config/samples/pytorch-mnist/ray-job.pytorch-mnist.yaml该文件位于 KubeRay 仓库的
ray-operator/config/samples/pytorch-mnist/目录下,你也可以直接从 KubeRay 仓库对应路径获取后按需修改。
部署前必须理解的三个关键字段
示例 YAML 中的资源需求量较大,直接应用到小机器上会导致 Pod 一直处于Pending状态。部署前请依据自己的机器资源调整以下字段:
| 字段 | 所在位置 | 含义与约束 |
|---|---|---|
replicas | rayClusterSpec.workerGroupSpecs | KubeRay 调度到集群的 worker Pod 数量。示例中每个 worker Pod 请求 3 个 CPU,head Pod 请求 1 个 CPU(见template字段),submitter Pod 还需 1 个 CPU。例如机器有 8 个 CPU,replicas最大取 2,才能保证所有 Pod 都进入Running状态。 |
NUM_WORKERS | spec.runtimeEnvYAML | 要启动的 Ray actor 数量(对应 ScalingConfig 的num_workers)。每个 Ray actor 必须由集群中的一个 worker Pod 承载,因此NUM_WORKERS必须小于等于replicas。 |
CPUS_PER_WORKER | spec.runtimeEnvYAML | 必须小于等于(每个 worker Pod 的 CPU 资源请求量) - 1。例如示例中 worker Pod 请求 3 CPU,则CPUS_PER_WORKER必须设为 2 或更小。 |
原因:KubeRay 会在 worker Pod 内先预留部分 CPU 给 Ray 运行时自身(如 GCS 客户端、调度相关组件),若
CPUS_PER_WORKER把 Pod 的 CPU 全部占满,Ray actor 会因资源不足而无法调度。
资源核算示例:若机器为 8 CPU,取replicas=2、NUM_WORKERS=2,则总需求为 head(1) + worker×2(3×2=6) + submitter(1) = 8 CPU,正好可以全部Running。
提交 RayJob 并核对状态
# `replicas` 和 `NUM_WORKERS` 均设为 2。 # 创建 RayJob。 kubectl apply -f ray-job.pytorch-mnist.yaml # 检查现有 Pod:根据 `replicas`,应有 2 个 worker Pod。 # 确保所有 Pod 都处于 `Running` 状态。 kubectl get pods # NAME READY STATUS RESTARTS AGE # kuberay-operator-6dddd689fb-ksmcs 1/1 Running 0 6m8s # rayjob-pytorch-mnist-raycluster-rkdmq-small-group-worker-c8bwx 1/1 Running 0 5m32s # rayjob-pytorch-mnist-raycluster-rkdmq-small-group-worker-s7wvm 1/1 Running 0 5m32s # rayjob-pytorch-mnist-nxmj2 1/1 Running 0 4m17s # rayjob-pytorch-mnist-raycluster-rkdmq-head-m4dsl 1/1 Running 0 5m32s从上到下依次是:KubeRay operator、两个 worker Pod、submitter Pod(名字与 RayJob 同名)、head Pod。确认 RayJob 进入RUNNING状态:
kubectl get rayjob # NAME JOB STATUS DEPLOYMENT STATUS START TIME END TIME AGE # rayjob-pytorch-mnist RUNNING Running 2024-06-17T04:08:25Z 11mStep 4:等待 RayJob 完成并查看训练结果
训练需要几分钟时间。等待 RayJob 完成后,JOB_STATUS会变为SUCCEEDED:
kubectl get rayjob # NAME JOB STATUS DEPLOYMENT STATUS START TIME END TIME AGE # rayjob-pytorch-mnist SUCCEEDED Complete 2024-06-17T04:08:25Z 2024-06-17T04:22:21Z 16m训练完成后 submitter Pod 会变为Completed(不再占用 CPU),而 RayCluster 的 head/worker Pod 因默认shutdownAfterJobFinishes=false仍保持Running(详见 RayJob Quickstart 中对集群回收行为的说明):
# 查看 Pod 名称与状态。 kubectl get pods # NAME READY STATUS RESTARTS AGE # kuberay-operator-6dddd689fb-ksmcs 1/1 Running 0 113m # rayjob-pytorch-mnist-raycluster-rkdmq-small-group-worker-c8bwx 1/1 Running 0 38m # rayjob-pytorch-mnist-raycluster-rkdmq-small-group-worker-s7wvm 1/1 Running 0 38m # rayjob-pytorch-mnist-nxmj2 0/1 Completed 0 38m # rayjob-pytorch-mnist-raycluster-rkdmq-head-m4dsl 1/1 Running 0 38m查看训练日志:
kubectl logs -f rayjob-pytorch-mnist-nxmj2 # 2024-06-16 22:23:01,047 INFO cli.py:36 -- Job submission server address: http://rayjob-pytorch-mnist-raycluster-rkdmq-head-svc.default.svc.cluster.local:8265 # 2024-06-16 22:23:01,844 SUCC cli.py:60 -- ------------------------------------------------------- # 2024-06-16 22:23:01,844 SUCC cli.py:61 -- Job 'rayjob-pytorch-mnist-l6ccc' submitted successfully # 2024-06-16 22:23:01,844 SUCC cli.py:62 -- ------------------------------------------------------- # ... # (RayTrainWorker pid=1138, ip=10.244.0.18) # 0%| | 0/26421880 [00:00<?, ?it/s] # (RayTrainWorker pid=1138, ip=10.244.0.18) # 0%| | 32768/26421880 [00:00<01:27, 301113.97it/s] # ... # Training finished iteration 10 at 2024-06-16 22:33:05. Total running time: 7min 9s # ╭───────────────────────────────╮ # │ Training result │ # ├───────────────────────────────┤ # │ checkpoint_dir_name │ # │ time_this_iter_s 28.2635 │ # │ time_total_s 423.388 │ # │ training_iteration 10 │ # │ accuracy 0.8748 │ # │ loss 0.35477 │ # ╰───────────────────────────────╯ # Training completed after 10 iterations at 2024-06-16 22:33:06. Total running time: 7min 10s # Training result: Result( # metrics={'loss': 0.35476621258825347, 'accuracy': 0.8748}, # path='/home/ray/ray_results/TorchTrainer_2024-06-16_22-25-55/TorchTrainer_122aa_00000_0_2024-06-16_22-25-55', # filesystem='local', # checkpoint=None # ) # ...日志解读要点:
- 前几行
cli.py输出来自 submitter 的ray job submit,确认 job 已成功提交到 head 服务(端口 8265 为 Ray Dashboard / job 提交地址); (RayTrainWorker pid=..., ip=...)前缀表明训练循环运行在 Ray Train 的 worker actor 上,tqdm进度条显示的26421880是 10 个 epoch 的总体样本迭代量;- 表格与
Result(...)展示的是TorchTrainer.fit()返回的训练结果:10 个迭代后 loss ≈ 0.3548、accuracy ≈ 0.8748,checkpoint 结果目录位于 head Pod 的/home/ray/ray_results/...。
清理资源
删除 RayJob 即可。由于示例默认不开启自动回收,head/worker Pod 会随 RayJob 的删除一并清理:
kubectl delete -f ray-job.pytorch-mnist.yaml若希望训练结束后自动回收 RayCluster,可在 RayJob 中设置shutdownAfterJobFinishes: true,并配合ttlSecondsAfterFinished控制延迟回收时间;相关行为细节可参考 RayJob Quickstart(其中ray-job.shutdown.yaml示例设置了shutdownAfterJobFinishes: true与ttlSecondsAfterFinished: 10,即 job 结束后 10 秒删除 RayCluster,而 submitter 因包含 job 日志会被保留,直至 RayJob 本身被删除)。
深入原理:训练脚本如何与 RayJob 协作
RayJob 的entrypoint指向的训练脚本位于仓库 python/ray/train/examples/pytorch/torch_fashion_mnist_example.py,其文档说明见 Train a PyTorch model on Fashion MNIST。理解脚本结构有助于你按需修改runtimeEnvYAML中的配置:
- 数据准备:
get_dataloaders使用torchvision下载 Fashion MNIST 数据集(含Normalize((0.28604,), (0.32025,))归一化),并用FileLock防止多 worker 并发下载冲突; - 分布式数据加载:
ray.train.torch.prepare_data_loader会对 DataLoader 做分片(shard),使每个 worker 只处理自己那部分数据; - 分布式模型包装:
ray.train.torch.prepare_model自动用 PyTorchDistributedDataParallel包装模型并移动到正确的设备(CPU/GPU); - 指标上报:每个 epoch 结束后调用
ray.train.report(metrics={"loss": test_loss, "accuracy": accuracy}),这正是最终Training result表格中accuracy与loss的来源; - 多 worker 配置:
train_fashion_mnist(num_workers=2, use_gpu=False)中通过ScalingConfig(num_workers=..., use_gpu=...)声明训练使用的 Ray actor 数量。
关于NUM_WORKERS与 ScalingConfig 的对应关系,源码 python/ray/train/v2/api/config.py 中ScalingConfig.num_workers的定义为"要启动的 worker(Ray actor)数量"(默认 1,也可传入(min, max)元组表示弹性范围),use_gpu为 True 时每个 worker 预留 1 块 GPU。因此:
- 纯 CPU 场景:
NUM_WORKERS(即num_workers)受限于集群中可调度 CPU 总量,每个 actor 占用CPUS_PER_WORKER个 CPU; - GPU 场景:将
num_workers设为 GPU 数量即可做到每个 worker 独占 1 块 GPU(本示例为 CPU 版,未启用 GPU)。
常见调整与注意事项
- 机器资源不足导致 Pod
Pending:优先减小replicas,再同步减小NUM_WORKERS,确保NUM_WORKERS <= replicas; - actor 调度失败:检查
CPUS_PER_WORKER是否超过(worker Pod CPU 请求量 - 1),预留 CPU 给 Ray 运行时; - 数据集下载:Fashion MNIST 数据会在每个 worker 上首次下载到
~/data(受FileLock保护),若网络受限可预先在镜像或 PVC 中准备数据; - 自动回收集群:如需省钱省资源,设置
shutdownAfterJobFinishes: true并配合ttlSecondsAfterFinished; - 扩展阅读:更多 RayJob 场景(批推理、Kueue 优先级调度、Gang 调度等)可参考 RayJob Quickstart 末尾的示例导航,以及仓库 doc/source/cluster/kubernetes/examples 目录下的其他示例文档。
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考