
人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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 Train 是 Ray 生态中面向大规模分布式训练与微调的可扩展机器学习库它以极简的 API 将模型训练代码从单机无缝扩展到云上多机集群并抽象掉分布式计算的全部复杂性。阅读本文后你将掌握 Ray Train 的四大核心概念训练函数、Worker、ScalingConfig、Trainer、一条命令完成安装的方法以及 PyTorch、PyTorch Lightning、Hugging Face Transformers、XGBoost、JAX 五大主流框架从单机脚本改造为分布式训练任务的完整实战路径。Ray Train 是什么面向大模型与大数据的分布式训练方案无论你面对的是大模型还是大数据集Ray Train 都致力于成为最简单的分布式训练解决方案。它将训练代码从单机扩展到云端集群中的多台机器并自动处理分布式计算带来的复杂性包括进程组初始化、设备放置、数据分片、指标上报与检查点持久化等。Ray Train 对主流机器学习框架提供了开箱即用的支持覆盖 PyTorch 生态与更多框架两大阵营PyTorch 生态更多框架PyTorchTensorFlowPyTorch LightningKerasHugging Face TransformersHorovodHugging Face AccelerateXGBoostDeepSpeedLightGBM其中每个框架都有对应的独立入门指南本文后续章节将逐一展开。安装 Ray Train安装 Ray Train 只需一条命令$ pip install -U ray[train]这条命令通过 Ray 的 Python 包安装机制完成。从仓库中的 python/setup.py 可以看出train这一 extras 依赖组被定义为 Tune extras 的超集setup_spec.extras[train] list(setup_spec.extras[tune])也就是说ray[train]会一并安装 Ray Tune超参数调优所需依赖——这正是 Ray Train 能与 Ray Tune 无缝集成做超参搜索的基础。需要说明的是ray[train]安装的是 Ray 核心训练框架本身各框架的第三方库需要按需额外安装。例如在 Hugging Face Transformers 入门指南中完整的安装命令是pip install ray[train] torch transformers[torch] datasets evaluate numpy scikit-learn更多关于 Ray 及各库的安装细节可参见仓库中的 安装文档本文按 doc/source/train/train.md 的指引指向官方安装章节。核心概念训练函数、Worker、ScalingConfig 与 Trainer要高效使用 Ray Train首先需要理解 overview.rst 中定义的四个核心概念训练函数Training Function一个包含完整模型训练循环逻辑的 Python 函数Worker运行训练函数的进程伸缩配置ScalingConfig对 Worker 数量与计算资源如 CPU 或 GPU的配置Trainer将训练函数、Worker 与伸缩配置三者绑定、用于执行分布式训练任务的 Python 类。训练函数train_func训练函数是用户自定义的 Python 函数包含端到端的模型训练循环逻辑。当启动分布式训练任务时每个 Worker 都会执行这个函数。Ray Train 文档的通用约定是train_func为用户定义的训练函数它被传入 Trainer 的train_loop_per_worker参数def train_func(): 在每个分布式 Worker 进程上运行的用户自定义训练函数。 该函数通常包含加载模型、加载数据集、训练模型、 保存检查点以及记录指标等逻辑。 ...WorkerRay Train 将模型训练计算分发到集群中独立的 Worker 进程上执行每个 Worker 就是一个运行train_func的进程。Worker 的数量决定了训练任务的并行度并通过ray.train.ScalingConfig配置。ScalingConfig定义训练规模ScalingConfig是定义训练任务规模的机制最基础的两个参数用于控制 Worker 并行度与计算资源num_workers分布式训练任务启动的 Worker 数量use_gpu每个 Worker 是否使用 GPU。from ray.train import ScalingConfig # 单 Worker CPU scaling_config ScalingConfig(num_workers1, use_gpuFalse) # 单 Worker GPU scaling_config ScalingConfig(num_workers1, use_gpuTrue) # 多 Worker每个 Worker 使用一块 GPU scaling_config ScalingConfig(num_workers4, use_gpuTrue)从仓库源码 python/ray/train/v2/api/config.py 可以看到ScalingConfig的字段远比表面丰富除num_workers、use_gpu外还包括resources_per_worker每个 Worker 所需的资源字典例如{CPU: 4}或{TPU: 4}accelerator_type实验性指定加速器类型如 TPU 代次确保任务调度到目标资源上use_tpu实验性Ray 2.49.0 新增字段置 True 时训练在 TPU 上执行topology实验性定义 TPU 芯片的物理排布如4x4多主机训练时必填弹性训练支持num_workers可传(min, max)二元组实现弹性伸缩elastic_resize_monitor_interval_s控制扩缩容监控间隔默认 60 秒本地模式当num_workers0时训练函数将在同一进程内直接运行见 config.py 中的本地模式日志逻辑方便调试。Trainer绑定一切并启动训练Trainer 将上述三个概念绑定在一起用于启动分布式训练任务。Ray Train 为不同框架提供了对应的 Trainer 类详见 API 参考。调用fit()方法执行训练任务其内部流程为按scaling_config的定义启动 Worker在所有 Worker 上搭建框架对应的分布式环境如 PyTorch 的进程组在所有 Worker 上运行train_func。from ray.train.torch import TorchTrainer trainer TorchTrainer(train_func, scaling_configscaling_config) trainer.fit()在源码层面所有框架 Trainer 的共同基类是 python/ray/train/base_trainer.py 中的BaseTrainer其构造参数正是train_loop_per_worker、scaling_config、run_config与datasets四件套fit()方法base_trainer.py在内部将 Trainer 转换为 Tune 可执行的 Trainable并通过 Ray Tune 的 Tuner 完成实际运行调度返回一个Result对象。快速上手五大框架入门指南主文档 train.md 提供了 Overview、PyTorch、PyTorch Lightning、Hugging Face Transformers、JAX 五张入门卡片并附 XGBoost 指南。下面逐一深入。PyTorch改造现有训练脚本入门指南 getting-started-pytorch.rst 演示了如何将现有 PyTorch 脚本改造为 Ray Train 分布式训练最终代码结构如下from ray.train.torch import TorchTrainer from ray.train import ScalingConfig def train_func(): # 你的 PyTorch 训练代码 ... scaling_config ScalingConfig(num_workers2, use_gpuTrue) trainer TorchTrainer(train_func, scaling_configscaling_config) result trainer.fit()关键改动集中在训练函数内部核心是三个工具函数1. 准备模型ray.train.torch.prepare_model(model)该函数完成两件事把模型移动到正确的设备CPU/GPU上并包上DistributedDataParallel。它替代了手写的model.to(device)与DistributedDataParallel(model, device_ids[device_id])-from torch.nn.parallel import DistributedDataParallel import ray.train.torch def train_func(): ... # 创建模型。 model ... # 设置分布式训练与设备放置。 - device_id ... # 自行获取正确设备的逻辑。 - model model.to(device_id or cpu) - model DistributedDataParallel(model, device_ids[device_id]) model ray.train.torch.prepare_model(model) ...2. 准备数据加载器ray.train.torch.prepare_data_loader(data_loader)该函数为 DataLoader 添加DistributedSampler并将 batch 移动到正确设备。注意DataLoader的batch_size是每个 Worker 的批大小全局批大小按global_batch_size worker_batch_size * ray.train.get_context().get_world_size()换算。若你已手动配置了DistributedSampler该函数不会重复添加但DistributedSampler不适用于包装IterableDataset的 DataLoader此类场景建议改用 Ray Data 做高性能流式数据摄入。3. 上报指标与检查点ray.train.report()metrics {loss: loss.item()} # 训练/验证指标 # 从目录构建 Ray Train 检查点 checkpoint ray.train.Checkpoint.from_directory(temp_checkpoint_dir) # Ray Train 会自动把检查点保存到持久化存储 # 因此本地临时目录 temp_checkpoint_dir 之后可以安全清理。 ray.train.report(metricsmetrics, checkpointcheckpoint)下面是以 FashionMNIST ResNet18 为例的完整可运行版本含训练完成后加载模型import os import tempfile import torch from torch.nn import CrossEntropyLoss from torch.optim import Adam from torch.utils.data import DataLoader from torchvision.models import resnet18 from torchvision.datasets import FashionMNIST from torchvision.transforms import ToTensor, Normalize, Compose import ray.train.torch def train_func(): # 模型、损失、优化器 model resnet18(num_classes10) model.conv1 torch.nn.Conv2d( 1, 64, kernel_size(7, 7), stride(2, 2), padding(3, 3), biasFalse ) model ray.train.torch.prepare_model(model) # 移动设备 DDP 包装 criterion CrossEntropyLoss() optimizer Adam(model.parameters(), lr0.001) # 数据 transform Compose([ToTensor(), Normalize((0.28604,), (0.32025,))]) data_dir os.path.join(tempfile.gettempdir(), data) train_data FashionMNIST(rootdata_dir, trainTrue, downloadTrue, transformtransform) train_loader DataLoader(train_data, batch_size128, shuffleTrue) train_loader ray.train.torch.prepare_data_loader(train_loader) # 训练 for epoch in range(10): if ray.train.get_context().get_world_size() 1: train_loader.sampler.set_epoch(epoch) for images, labels in train_loader: # 设备放置由 prepare_data_loader 完成 outputs model(images) loss criterion(outputs, labels) optimizer.zero_grad() loss.backward() optimizer.step() # 上报指标与检查点 metrics {loss: loss.item(), epoch: epoch} with tempfile.TemporaryDirectory() as temp_checkpoint_dir: torch.save( model.module.state_dict(), os.path.join(temp_checkpoint_dir, model.pt) ) ray.train.report( metrics, checkpointray.train.Checkpoint.from_directory(temp_checkpoint_dir), ) if ray.train.get_context().get_world_rank() 0: print(metrics) scaling_config ray.train.ScalingConfig(num_workers2, use_gpuTrue) trainer ray.train.torch.TorchTrainer( train_func, scaling_configscaling_config, # 多节点集群中应在此配置所有 Worker 节点都能访问的持久化存储 # 例如 run_configray.train.RunConfig(storage_paths3://...) ) result trainer.fit() # 加载训练好的模型 with result.checkpoint.as_directory() as checkpoint_dir: model_state_dict torch.load(os.path.join(checkpoint_dir, model.pt)) model resnet18(num_classes10) model.conv1 torch.nn.Conv2d( 1, 64, kernel_size(7, 7), stride(2, 2), padding(3, 3), biasFalse ) model.load_state_dict(model_state_dict)PyTorch Lightning利用 Ray 专属策略与回调入门指南 getting-started-pytorch-lightning.rst 指出Ray Train 已在每个 Worker 上搭好分布式进程组你只需对 Lightning Trainer 做少量改动。核心改动一使用 Ray 提供的分布式策略。Ray Train 为 Lightning 提供三个保留基类同名参数的子类策略内部会配置根设备与分布式采样器参数ray.train.lightning.RayDDPStrategyray.train.lightning.RayFSDPStrategyray.train.lightning.RayDeepSpeedStrategy核心改动二配置 Ray 集群环境插件。ray.train.lightning.RayLightningEnvironment负责配置 Worker 的 local/global/node rank 与 world sizetrainer pl.Trainer( - devices[0, 1, 2, 3], - strategyDDPStrategy(), - plugins[LightningEnvironment()], devicesauto, acceleratorauto, strategyray.train.lightning.RayDDPStrategy(), plugins[ray.train.lightning.RayLightningEnvironment()] ) trainer ray.train.lightning.prepare_trainer(trainer)注意devices与accelerator应始终使用auto因为 Ray TorchTrainer 已为你正确配置了CUDA_VISIBLE_DEVICES。核心改动三添加RayTrainReportCallback回调用于持久化检查点与上报指标trainer pl.Trainer( max_epochs10, devicesauto, acceleratorauto, strategyray.train.lightning.RayDDPStrategy(), plugins[ray.train.lightning.RayLightningEnvironment()], callbacks[ray.train.lightning.RayTrainReportCallback()], enable_checkpointingFalse, # 可选关闭默认检查点行为 ) trainer ray.train.lightning.prepare_trainer(trainer) trainer.fit(model, train_dataloaderstrain_dataloader)RayTrainReportCallback还支持两个进阶选项checkpoint_upload_modeCheckpointUploadMode.ASYNC将检查点上传卸载到 Ray Train 管理的后台线程避免阻塞 Lightning 训练循环validation参数则启动异步 Ray 任务校验检查点替代在训练 Worker 中同步执行validation_step注意此模式与 Lightning 的AsyncCheckpointIO插件不兼容因为 Ray Train 需要控制上传线程以等待其完成后提交检查点。上报指标与检查点是支持容错训练与超参数优化的前提。版本与迁移说明Ray Train 官方测试覆盖pytorch_lightning的1.6.5与2.1.2版本建议使用pytorch_lightning1.6.5使用 Lightning 2.x 时请改用lightning.pytorch.xxx导入路径。API 层面Ray 2.4 曾引入LightningTrainerLightningConfigBuilder的黑盒式 APIRay 2.7 起统一为透明的TorchTrainer方案让你直接控制原生 Lightning 代码两种写法的完整对照见指南中的迁移章节。Hugging Face Transformers文本模型分布式训练入门指南 getting-started-transformers.rst 展示了如何将现有 Transformers 训练脚本迁移到 Ray Train。核心思路同样是把所有逻辑数据集构造与预处理、模型初始化、Transformers Trainer 定义放入训练函数然后做两处关键修改1. 添加上报回调import transformers from ray.train.huggingface.transformers import RayTrainReportCallback def train_func(): ... trainer transformers.Trainer(...) trainer.add_callback(RayTrainReportCallback()) ...2. 准备 Trainer将 Transformers Trainer 传入ray.train.huggingface.transformers.prepare_trainer(trainer)以校验配置并启用 Ray Data 集成随后即可启动训练。一个重要的序列化注意事项使用 Hugging Face Datasets 或 Evaluate 时务必在训练函数内部调用datasets.load_dataset与evaluate.load切勿把已加载的数据集和指标从函数外部传入——否则在把对象传输给 Worker 时可能引发序列化错误。训练完成后可通过result.checkpoint.as_directory()找到RayTrainReportCallback.CHECKPOINT_NAME指定的检查点文件用AutoModelForSequenceClassification.from_pretrained直接加载模型。同样地Transformers 集成也存在 API 迁移Ray 2.1 引入的TransformersTrainertrainer_init_per_worker接口正在被 Ray 2.7 统一的TorchTrainer取代新 API 与标准 Transformers 脚本风格更贴近且支持通过torch_batches迭代 Ray Data 分片ray.train.get_dataset_shard作为训练数据。XGBoost梯度提升树的分布式训练入门指南 getting-started-xgboost.rst 讲解了 XGBoost 脚本的改造路径。Ray Train 会自动完成分布式 XGBoost 训练所需的 Worker 通信设置你只需1. 包装训练函数并配置参数。通过train_loop_config字典向train_func传参def train_func(config): label_column config[label_column] num_boost_round config[num_boost_round] ... config {label_column: y, num_boost_round: 10} trainer ray.train.xgboost.XGBoostTrainer(train_func, train_loop_configconfig, ...)注意避免通过train_loop_config传递大数据对象如数据集、模型以减少序列化/反序列化开销应在train_func内部直接初始化大对象。2. 添加上报回调import xgboost from ray.train.xgboost import RayTrainReportCallback def train_func(): ... bst xgboost.train( ..., callbacks[ RayTrainReportCallback(metrics[eval-logloss], frequency1) ], ) ...3. 数据分片。分布式训练时每个 Worker 应使用数据集的不同分片。最简单的方式是预先切分数据集并按ray.train.get_world_rank()分配文件更灵活的方式是用 Ray Data 在运行时自动分片先整体加载数据集在训练函数内通过ray.train.get_dataset_shard(dataset_name)取本 Worker 的分片并转为xgboost.DMatrix最后把数据集以字典形式传给 Trainerkey 必须与get_dataset_shard的调用一致trainer XGBoostTrainer(..., datasets{train: train_dataset, eval: eval_dataset}) trainer.fit()4. 配置规模与 GPU。ScalingConfig可配置num_workers、use_gpu、resources_per_worker# 4 个节点每节点 8 CPU scaling_config ScalingConfig(num_workers4, resources_per_worker{CPU: 8}) # 4 个 Worker每个使用 1 块 GPU scaling_config ScalingConfig(num_workers4, use_gpuTrue)使用 GPU 时还需在 XGBoost 参数中设置device: cuda。另有一个实战要点同时使用 Ray Data 时不要用resources_per_worker占用集群全部 CPU——Ray Data 需要 CPU 并行执行预处理否则会成为性能瓶颈。例如每节点 8 CPU 时可为训练 Worker 分配 6 CPU、留 2 CPU 给 Ray Data。5. 配置持久化存储与读取结果。通过RunConfig(storage_path..., name...)指定结果含检查点与工件的保存路径支持本地路径、S3 云存储与 NFS。对单节点集群共享存储位置是可选的但对多节点集群共享存储云存储或 NFS是必需的使用本地路径会在检查点阶段报错。训练完成后返回的Result对象源码定义见 python/ray/train/v2/api/result.py包含result.metrics # 训练期间上报的指标 result.checkpoint # 训练期间上报的最新检查点 result.path # 日志存放路径 result.error # 训练失败时抛出的异常JAXGPU 与 TPU 上的 SPMD 训练入门指南 getting-started-jax.rst 介绍 Ray Train 的ray.train.v2.jax.JaxTrainer。JAX 是面向加速器的数组计算与程序变换库利用 XLA 编译器为 GPU/TPU 生成高度优化的代码其核心威力在于jax.grad、jax.jit、jax.vmap等变换的可组合性。JaxTrainer 遵循单程序多数据SPMD范式让训练代码在多个 Worker 上同时执行TPU 场景每个 Worker 运行在 TPU 切片内的一台独立 TPU 虚拟机上Ray Train 自动完成 TPU 切片的原子性预留GPU 场景Ray 自动在 CUDA 设备上搭建 JAX 分布式系统。TPU 伸缩配置需要用到 Ray 2.49.0 在 V2ScalingConfig中新增的两个字段from ray.train import ScalingConfig tpu_scaling_config ScalingConfig( num_workers4, # TPU 切片中 TPU VM 总数 use_tpuTrue, # 初始化 JAX TPU 后端 topology4x4, # TPU 芯片物理排布多主机训练必填 accelerator_typeTPU-V6E, # TPU 代次确保调度到目标切片 placement_strategySPREAD, # 将 Worker 分散放置 )各字段含义use_tpu布尔开关topology描述 TPU 芯片排布如4x4多主机训练必须设置num_workers应等于所有切片 TPU VM 总数例如一个 v4-32 切片含 2x2x4 拓扑共 4 台 VM则设为 4两个 v4-32 切片则设为 8resources_per_worker通常请求每台 VM 的芯片数如{TPU: 4}accelerator_type指定 TPU 代次。GPU 场景则保持经典写法ScalingConfig(num_workers4, use_gpuTrue)默认每 Worker 一块 GPU。改造 JAX 脚本只需把训练逻辑包进一个函数并传递给JaxTrainer同时把print替换为ray.train.report({loss: ..., epoch: ...})上报指标。训练完成后同样通过result.metrics读取最终指标如result.metrics[loss]。更多框架支持如果你的框架不在上述五大指南中可参考 more-frameworks.rst 索引下的集成指南Hugging Face Accelerate 指南DeepSpeed 指南TensorFlow 与 Keras 指南LightGBM 指南Horovod 指南仓库中对应的可直接运行示例代码位于 doc/source/train/doc_code/ 目录包括tf_starter.pyTensorFlow 起步、lightgbm_quickstart.pyLightGBM 快速入门、xgboost_quickstart.pyXGBoost 快速入门与hvd_trainer.pyHorovod等。深入进阶User Guides 使用指南全景主文档将 User Guides 定位为常见训练任务的如何做指南集合覆盖以下主题每篇均配有独立文档数据加载与预处理如何用 Ray Data 摄入、分片与预处理训练数据使用加速器CPU/GPU/TPU 资源配置细节持久化存储跨节点结果与检查点的可靠存储监控与日志训练过程的日志与指标监控检查点检查点的保存、管理与加载异步校验检查点的异步验证模式实验追踪与 MLflow 等工具集成结果访问解读Result对象中的指标、检查点与最佳检查点容错训练Worker 故障恢复与容错策略弹性训练训练中动态增删 Worker监控你的应用面向应用层的监控方案本地模式num_workers0时在单进程内调试运行可复现性固定随机种子、复现训练结果超参数优化与 Ray Tune 集成做超参搜索。这些指南与主文档相互呼应例如 hyperparameter-optimization.rst 对应上报指标与检查点能力ray.train.reportfault-tolerance.rst 与 elastic-training.rst 则分别依赖RunConfig(max_failures...)与ScalingConfig的弹性num_workers元组语法这些都可在 python/ray/train/v2/api/config.py 的RunConfig定义与ScalingConfig源码中找到对应字段。更多学习资源教程、示例、基准与 API除了入门指南与使用指南train.md 还提供以下进阶资源Tutorials教程覆盖从视觉到推荐系统等 ML 负载模式的动手实践教程Examples端到端示例仓库 doc/source/train/examples/ 目录下按框架组织了大量可运行示例例如 PyTorch 的 Fashion MNIST 示例 与回归示例、Transformers 的文本分类 notebook、Lightning 的 MNIST 示例、DeepSpeed 的 GPT-J 微调、JAX 的 intro_to_jax_trainer 等覆盖模型微调、强化学习RLHF、推理等真实场景Benchmarks基准测试benchmarks.rst 提供 Ray Train 的规模基准数据API 参考doc/source/train/api/api.md 提供所有 Trainer 类、ScalingConfig、RunConfig、Checkpoint、Result等类的完整 API 描述以及ray.train.report、ray.train.get_context等训练循环工具函数的用法。总结本文以 doc/source/train/train.md 为主线完整覆盖了 Ray Train 的安装方式pip install -U ray[train]、四大核心概念训练函数、Worker、ScalingConfig、Trainer以及 PyTorch、PyTorch Lightning、Hugging Face Transformers、XGBoost、JAX 五大框架的分布式改造实战并延伸介绍了更多框架集成与全部使用指南。无论你是要把现有单机脚本扩展到集群还是要训练超大模型、处理超大数据集Ray Train 都提供了统一、简洁的入口定义一个训练函数配好ScalingConfig选择对应的 Trainer调用fit()即可获得包含指标与检查点的Result。下一步建议结合 examples 目录 中的端到端示例动手实践并在 API 参考 中查阅更多细节。赞分享人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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点击查看免费下载相关推荐Oumi 训练指南从单卡微调到大规模分布式训练的端到端实践Oumi 训练指南从单卡微调到大规模分布式训练的端到端实践 导读 本文围绕开源项目 Oumi GitHub_Trending/ou/oumi 的官方训练文人工智能大模型预训练微调强化学习模型评测模型推理服务MCP 服务LLMOpsseq2seq模型训练完全手册从单机到分布式训练的终极指南seq2seq模型训练完全手册从单机到分布式训练的终极指南 想要掌握seq2seq模型的完整训练流程吗这份终极指南将带你从基础的单机训练到高效的分布式训练深度学习NLPUltralytics YOLO26 Train 模式完全指南从单卡训练到多硬件分布式调优Ultralytics YOLO26 Train 模式完全指南从单卡训练到多硬件分布式调优 本指南以 Ultralytics YOLO26 的 Train训人工智能计算机视觉深度学习机器学习预训练上一篇100行代码搞定智能会议纪要基于internlm_20b_base_ms的企业级解决方案下一篇Jekyll邮件订阅终极指南如何快速搭建Newsletter集成方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考