ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

在 Kubernetes 上使用 RayJob 分布式训练 Fashion MNIST:PyTorch + Ray Train 端到端实战指南

在 Kubernetes 上使用 RayJob 分布式训练 Fashion MNIST:PyTorch + Ray Train 端到端实战指南 在 Kubernetes 上使用 RayJob 分布式训练 Fashion MNISTPyTorch 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 QuickstartRayJobKubeRay 提供的 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 --imagekindest/node:v1.26.0说明Kind 适合快速验证与本地开发生产环境请使用托管的 Kubernetes 服务如 EKS、GKE、ACK 等或自建集群命令与本文保持一致。Step 2安装 KubeRay operatorKubeRay 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?refv1.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创建 RayJobRayJob 由两部分组成一个 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状态。部署前请依据自己的机器资源调整以下字段字段所在位置含义与约束replicasrayClusterSpec.workerGroupSpecsKubeRay 调度到集群的 worker Pod 数量。示例中每个 worker Pod 请求 3 个 CPUhead Pod 请求 1 个 CPU见template字段submitter Pod 还需 1 个 CPU。例如机器有 8 个 CPUreplicas最大取 2才能保证所有 Pod 都进入Running状态。NUM_WORKERSspec.runtimeEnvYAML要启动的 Ray actor 数量对应 ScalingConfig 的num_workers。每个 Ray actor 必须由集群中的一个 worker Pod 承载因此NUM_WORKERS必须小于等于replicas。CPUS_PER_WORKERspec.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取replicas2、NUM_WORKERS2则总需求为 head(1) worker×2(3×26) 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会变为SUCCEEDEDkubectl 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 因默认shutdownAfterJobFinishesfalse仍保持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 pid1138, ip10.244.0.18) # 0%| | 0/26421880 [00:00?, ?it/s] # (RayTrainWorker pid1138, ip10.244.0.18) # 0%| | 32768/26421880 [00:0001: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, # filesystemlocal, # checkpointNone # ) # ...日志解读要点前几行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.8748checkpoint 结果目录位于 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_workers2, use_gpuFalse)中通过ScalingConfig(num_workers..., use_gpu...)声明训练使用的 Ray actor 数量。关于NUM_WORKERS与 ScalingConfig 的对应关系源码 python/ray/train/v2/api/config.py 中ScalingConfig.num_workers的定义为要启动的 workerRay actor数量默认 1也可传入(min, max)元组表示弹性范围use_gpu为 True 时每个 worker 预留 1 块 GPU。因此纯 CPU 场景NUM_WORKERS即num_workers受限于集群中可调度 CPU 总量每个 actor 占用CPUS_PER_WORKER个 CPUGPU 场景将num_workers设为 GPU 数量即可做到每个 worker 独占 1 块 GPU本示例为 CPU 版未启用 GPU。常见调整与注意事项机器资源不足导致 PodPending优先减小replicas再同步减小NUM_WORKERS确保NUM_WORKERS replicasactor 调度失败检查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),仅供参考
返回列表