
Ray Train 分布式训练入门将 PyTorch 脚本改造成 TorchTrainer 全流程指南【免费下载链接】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本文是一份以 doc/source/train/getting-started-pytorch.rst 为主体的实战指南。它面向已有 PyTorch 单机训练脚本的开发者演示如何利用 Ray Train 的TorchTrainer将脚本改造成可扩展的分布式训练程序涵盖模型与数据加载器的分布式改造、指标与检查点上报、ScalingConfig资源/GPU 配置、RunConfig持久化存储设置以及训练结果的访问与模型加载。读完本文你将掌握从单机 PyTorch到多 worker 分布式训练的最小改造路径并理解每一步在 Ray 源码中的实现原理。Ray Train 与 TorchTrainer 的角色分工Ray Train 是 Ray 提供的分布式训练框架。它把分布式训练抽象成几个可组合的部件训练函数training function在每个分布式训练 worker 上执行的纯 Python/PyTorch 代码即train_funcScalingConfig定义分布式训练 worker 的数量以及是否使用 GPU对应源码 ScalingConfigTorchTrainer负责编排并启动整个分布式训练任务对应源码 torch_trainer.py。借助这三个部件你可以继续使用熟悉的 PyTorch APIDataLoader、DistributedDataParallel、torch.save等由 Ray Train 替你处理进程组初始化、数据分片、设备放置与检查点持久化等分布式细节。本教程对应的仓库源码位于 python/ray/train/torch/核心实现 train_loop_utils.py、torch_trainer.py、config.py测试用例可参考 test_torch_trainer.py。快速开始最小的分布式训练骨架把现有训练脚本接入 Ray Train最小改动只有三步from ray.train.torch import TorchTrainer from ray.train import ScalingConfig def train_func(): # Your PyTorch training code here. ... scaling_config ScalingConfig(num_workers2, use_gpuTrue) trainer TorchTrainer(train_func, scaling_configscaling_config) result trainer.fit()三个元素的分工如下train_func在每个分布式训练 worker 上执行的 Python 代码即你的训练逻辑本体ScalingConfig定义分布式 worker 的数量与是否使用 GPU见后文配置扩展与资源TorchTrainer启动分布式训练任务返回一个包含指标与检查点的Result对象。对比同样的脚本改造前后有何不同下面用一份 FashionMNIST ResNet18 的训练脚本完整对比纯 PyTorch与PyTorch Ray Train两种写法。两份代码逻辑一致同样的模型、同样的优化器、同样的 10 个 epoch 训练循环。改造后PyTorch Ray Trainimport 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, Loss, Optimizer model resnet18(num_classes10) model.conv1 torch.nn.Conv2d( 1, 64, kernel_size(7, 7), stride(2, 2), padding(3, 3), biasFalse ) # [1] Prepare model. model ray.train.torch.prepare_model(model) # model.to(cuda) # This is done by prepare_model criterion CrossEntropyLoss() optimizer Adam(model.parameters(), lr0.001) # Data 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) # [2] Prepare dataloader. train_loader ray.train.torch.prepare_data_loader(train_loader) # Training 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: # This is done by prepare_data_loader! # images, labels images.to(cuda), labels.to(cuda) outputs model(images) loss criterion(outputs, labels) optimizer.zero_grad() loss.backward() optimizer.step() # [3] Report metrics and checkpoint. 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) # [4] Configure scaling and resource requirements. scaling_config ray.train.ScalingConfig(num_workers2, use_gpuTrue) # [5] Launch distributed training job. trainer ray.train.torch.TorchTrainer( train_func, scaling_configscaling_config, # [5a] If running in a multi-node cluster, this is where you # should configure the runs persistent storage that is accessible # across all worker nodes. # run_configray.train.RunConfig(storage_paths3://...), ) result trainer.fit() # [6] Load the trained model. 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)改造前纯 PyTorchimport 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 # Model, Loss, Optimizer model resnet18(num_classes10) model.conv1 torch.nn.Conv2d( 1, 64, kernel_size(7, 7), stride(2, 2), padding(3, 3), biasFalse ) model.to(cuda) criterion CrossEntropyLoss() optimizer Adam(model.parameters(), lr0.001) # Data transform Compose([ToTensor(), Normalize((0.28604,), (0.32025,))]) train_data FashionMNIST(root./data, trainTrue, downloadTrue, transformtransform) train_loader DataLoader(train_data, batch_size128, shuffleTrue) # Training for epoch in range(10): for images, labels in train_loader: images, labels images.to(cuda), labels.to(cuda) outputs model(images) loss criterion(outputs, labels) optimizer.zero_grad() loss.backward() optimizer.step() metrics {loss: loss.item(), epoch: epoch} checkpoint_dir tempfile.mkdtemp() checkpoint_path os.path.join(checkpoint_dir, model.pt) torch.save(model.state_dict(), checkpoint_path) print(metrics)两版代码的关键差异集中在五处下文逐一展开关注点纯 PyTorchPyTorch Ray Train设备放置手动model.to(cuda)、images.to(cuda)prepare_model/prepare_data_loader自动完成数据分片无每个进程读全量数据prepare_data_loader注入DistributedSampler并行包装手动DistributedDataParallelprepare_model自动包装指标/检查点本地打印、写临时目录ray.train.report统一上报与持久化启动方式手动多进程/多卡调度TorchTrainer.fit()一键启动设置训练函数train_func改造的第一步是把训练代码整体包进一个训练函数中每个分布式 worker 都会执行这个函数def train_func(): # Your model training code here. ...通过 train_loop_config 传参训练函数可以接收一个字典参数config该字典通过 Trainer 的train_loop_config传入def train_func(config): lr config[lr] num_epochs config[num_epochs] config {lr: 1e-4, num_epochs: 10} trainer ray.train.torch.TorchTrainer(train_func, train_loop_configconfig, ...)不要用 train_loop_config 传大数据对象train_loop_config需要被序列化并分发到每个 worker因此应避免通过它传递大数据对象如数据集、模型以减少序列化/反序列化开销。更推荐的做法是在train_func内部直接初始化这些大对象def load_dataset(): # Return a large in-memory dataset ... def load_model(): # Return a large in-memory model instance ... # 不推荐大数据对象随 config 序列化 config {data: load_dataset(), model: load_model()} def train_func(config): data config[data] model config[model] ... # 推荐在 train_func 内部加载 def train_func(config): data load_dataset() model load_model() ... trainer ray.train.torch.TorchTrainer(train_func, train_loop_configconfig, ...)设置模型prepare_model在训练函数中用ray.train.torch.prepare_model这一工具函数替换手动的分布式改造它替你完成两件事把模型移动到正确的设备CPU 或 GPU将模型包装进DistributedDataParallel当 worker 数大于 1 时。对应的 diff 改造如下-from torch.nn.parallel import DistributedDataParallel import ray.train.torch def train_func(): ... # Create model. model ... # Set up distributed training and device placement. - device_id ... # Your logic to get the right device. - model model.to(device_id or cpu) - model DistributedDataParallel(model, device_ids[device_id]) model ray.train.torch.prepare_model(model) ...源码级原理从 train_loop_utils.py 的实现可以看到prepare_model的关键行为move_to_deviceTrue时自动执行model.to(device)也可以传一个torch.device实例指定目标设备python/ray/train/torch/train_loop_utils.py#L418-L433parallel_strategy支持ddp默认、fsdpFullyShardedDataParallel要求 torch1.11.0 且 GPU 可用、None不包装三种取值仅当world_size 1时才实际执行 DDP/FSDP 包装采用 DDP 且使用 GPU 时会自动注入device_ids与output_devicepython/ray/train/torch/train_loop_utils.py#L469-L495。这也解释了为何同一份train_func可以做到无论单机单卡还是多机多卡代码完全一致。仓库自带的端到端测试 test_torch_e2e 对prepare_model开/关两种路径都做了验证。关于设备get_device如果你需要在训练函数里手动获取当前 worker 的设备可以使用ray.train.torch.get_device()见 train_loop_utils.py。例如model.to(ray.train.torch.get_device())它假定CUDA_VISIBLE_DEVICES已被正确设置返回当前进程对应的 torch 设备。当每个 worker 申请了多张 GPU 时get_device()返回其中索引最小的设备如需完整列表使用get_devices()train_loop_utils.py。设置数据加载器prepare_data_loaderray.train.torch.prepare_data_loader负责把普通的DataLoader改造成分布式可用的加载器它做两件事为DataLoader注入DistributedSampler实现跨 worker 的数据分片自动把每个 batch 移动到正确的设备替代手写的X.to(device)。from torch.utils.data import DataLoader import ray.train.torch def train_func(): ... dataset ... data_loader DataLoader(dataset, batch_sizeworker_batch_size, shuffleTrue) data_loader ray.train.torch.prepare_data_loader(data_loader) for epoch in range(10): if ray.train.get_context().get_world_size() 1: data_loader.sampler.set_epoch(epoch) for X, y in data_loader: - X X.to_device(device) - y y.to_device(device) ...关于 batch_size 的口径注意DataLoader接收的batch_size是每个 worker 各自的 batch size。全局 batch size 与单 worker batch size 的换算关系为global_batch_size worker_batch_size * ray.train.get_context().get_world_size()其中ray.train.get_context().get_world_size()返回参与训练的 worker 总数上下文 API 定义见 context.py还提供get_world_rank、get_local_rank等。源码级原理与注意事项从 train_loop_utils.py 的实现看prepare_data_loader仅在同时满足以下条件时才注入DistributedSampler训练 worker 数大于 1world_size 1用户没有手动设置过DistributedSampler—— 若已手动设置它会尊重现有 sampler 的配置不会重复添加数据集不是IterableDatasetsampler 对迭代型数据集无效。此外还有几点值得注意set_epoch 的必要性DistributedSampler需要跨 epoch 重新打乱数据。当 worker 数大于 1 时在每个 epoch 开始前调用data_loader.sampler.set_epoch(epoch)是必须的否则所有 epoch 的数据顺序将保持一致shuffle 失效。IterableDataset 的替代方案DistributedSampler无法配合包装了IterableDataset的DataLoader。如果数据来自迭代器iterator建议改用 Ray Data 作为数据入口——它为大规模数据集提供了高性能的流式数据摄入。详见 Working with PyTorchRay Data 数据摄入。可调参数prepare_data_loader(data_loader, add_dist_samplerTrue, move_to_deviceTrue, auto_transferTrue)。其中auto_transfer在 GPU 场景下会创建额外的 CUDA 流将数据从主机内存到显存的拷贝与默认 CUDA 流上的训练计算重叠见 train_loop_utils.py 与_WrappedDataLoader的预取实现 train_loop_utils.py设备为 CPU 时该选项自动失效。上报指标与检查点ray.train.report为监控训练进度可用ray.train.report上报中间指标与检查点。检查点先写入本地临时目录再通过ray.train.Checkpoint.from_directory包装成 Ray Train 检查点对象import os import tempfile import ray.train def train_func(): ... with tempfile.TemporaryDirectory() as temp_checkpoint_dir: torch.save( model.state_dict(), os.path.join(temp_checkpoint_dir, model.pt) ) metrics {loss: loss.item()} # Training/validation metrics. # Build a Ray Train checkpoint from a directory checkpoint ray.train.Checkpoint.from_directory(temp_checkpoint_dir) # Ray Train will automatically save the checkpoint to persistent storage, # so the local temp_checkpoint_dir can be safely cleaned up after. ray.train.report(metricsmetrics, checkpointcheckpoint) ...Checkpoint.from_directory的实现在 python/ray/train/_checkpoint.py它把本地目录封装为一个基于本地文件系统的Checkpoint对象。上报后Ray Train 会负责把检查点同步到持久化存储因此函数退出后临时目录即可安全清理。更多细节监控与日志上报参见 monitoring-logging.rst检查点机制与生命周期参见 checkpoints.rst。配置扩展与资源ScalingConfig在训练函数之外创建ray.train.ScalingConfig对象来配置num_workers—— 分布式训练 worker 进程Ray actor的数量use_gpu—— 每个 worker 是否使用 GPU否则使用 CPU。from ray.train import ScalingConfig scaling_config ScalingConfig(num_workers2, use_gpuTrue)更完整的资源配置从 ScalingConfig 的源码看它还支持以下参数trainer_resources分配给训练协调者coordinator的资源。协调者负责启动 worker 组并在每个 worker 上执行训练函数该进程不需要 GPU且总是与 rank 0 worker 调度在同一节点。默认分配 1 个 CPUresources_per_worker每个 worker 额外预留的资源可用于覆盖默认的每 worker 1 CPU / 1 GPU也可申请自定义资源placement_strategyworker 所在 placement group 的放置策略accelerator_type实验性限定集群中指定类型的加速器节点来启动协调者与 worker。资源键对大小写敏感CPU、GPU为大写memory为小写memory单位为字节如1e9表示 1 GB。一个完整的示例scaling_config ScalingConfig( # Number of distributed workers. num_workers2, # Turn on/off GPU. use_gpuTrue, # Assign extra CPU/GPU/custom resources per worker. resources_per_worker{GPU: 1, CPU: 1, memory: 1e9, custom: 1.0}, # Try to schedule workers on different nodes. placement_strategySPREAD, )从源码python/ray/air/config.py可以推断当use_gpuTrue时每个 worker 默认申请{GPU: 1}外加 1 个 CPUuse_gpuFalse时默认申请{CPU: 1}。若resources_per_worker与use_gpu冲突如use_gpuFalse却显式申请了 GPU会在校验阶段直接报错。配置持久化存储RunConfig创建ray.train.RunConfig对象指定训练结果含检查点与产物的保存路径from ray.train import RunConfig # Local path (/some/local/path/unique_run_name) run_config RunConfig(storage_path/some/local/path, nameunique_run_name) # Shared cloud storage URI (s3://bucket/unique_run_name) run_config RunConfig(storage_paths3://bucket, nameunique_run_name) # Shared NFS path (/mnt/nfs/unique_run_name) run_config RunConfig(storage_path/mnt/nfs, nameunique_run_name)多节点集群必须使用共享存储对于单节点集群指定共享存储位置是可选的但对于多节点集群共享存储云存储或 NFS是必需项。如果多节点集群使用本地路径在检查点持久化阶段会触发报错Ray Train 会检测到 worker 之间无法访问彼此的本地存储详见 persistent-storage.rst。启动分布式训练TorchTrainer 与 fit()把以上各部件组合起来即可用TorchTrainer启动分布式训练任务from ray.train.torch import TorchTrainer trainer TorchTrainer( train_func, scaling_configscaling_config, run_configrun_config ) result trainer.fit()TorchTrainer 的完整构造参数从 torch_trainer.py 的构造函数可以看到TorchTrainer继承自DataParallelTrainer完整支持以下参数train_loop_per_worker训练函数必填train_loop_config传给训练函数的字典torch_configTorchConfig对象用于配置 PyTorch 分布式进程组scaling_config扩展与资源配置run_config运行配置含存储路径datasets/dataset_config以字典形式传入的 Ray Data 数据集可在训练函数内通过ray.train.get_dataset_shard(name)访问默认按 worker 平均分片metadata随训练与检查点暴露的元数据字典须 JSON 可序列化resume_from_checkpoint用于断点续训的检查点可在训练函数内通过ray.train.get_checkpoint()访问。底层TorchConfig 与进程组初始化TorchTrainer未显式传入torch_config时会默认构造一个TorchConfig()torch_trainer.py。从 config.py 看TorchConfig包含三个字段backendtorch.distributed使用的通信后端默认None表示自动选择——使用 GPU 时选nccl否则选gloo即源码注释所述 GPU training uses NCCL and CPU training uses Glooinit_method初始化方式env环境变量或tcp默认envtimeout_s进程组操作的超时秒数默认1800。在启动阶段_setup_torch_process_group会调用torch.distributed.init_process_group建立进程组并在使用 NCCL 时默认设置TORCH_NCCL_ASYNC_ERROR_HANDLING1使集合通信超时能及时暴露为错误config.py。训练结束时_shutdown_torch负责销毁进程组并清空 CUDA 缓存。fit() 的返回值fit()定义在BaseTrainer中base_trainer.py它会将 Trainer 转为可调度的 Trainable 并交给 Ray Tune 执行返回一个Result对象训练失败时抛出TrainingFailedError。访问训练结果Result训练完成后返回的ray.train.Result对象封装了整个训练运行的信息包括训练期间上报的指标与检查点。定义见 python/ray/air/result.py常用属性如下result.metrics # The metrics reported during training. result.checkpoint # The latest checkpoint reported during training. result.path # The path where logs are stored. result.error # The exception that was raised, if training failed. result.metrics_dataframe # Full history of metrics, indexed by iteration. result.best_checkpoints # Best checkpoints with their associated metrics.result.metrics最近一次上报的指标字典result.checkpoint最近一次上报的检查点result.path结果目录在持久化存储上的路径本地或 S3 等远程位置result.error若训练失败保存所抛出的异常result.metrics_dataframe以迭代iteration为索引的完整指标历史 DataFrame。更完整的用法示例参见 results.rst。加载训练好的模型拿到result.checkpoint后可以用as_directory()上下文管理器将其内容暴露为本地目录远程检查点会自动下载到临时目录退出上下文后清理本地检查点则直接返回路径见 _checkpoint.py然后按常规方式反序列化模型权重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 脚本到 Ray Train 的改造后可以继续深入用户指南查阅 User Guides 了解特定任务的完整用法例如 results.rst结果检查、checkpoints.rst检查点、persistent-storage.rst持久化存储、using-accelerators.rst加速器与扩展配置、monitoring-logging.rst监控与日志端到端示例浏览 train/examples 目录查看 Ray Train 的完整实战示例API 参考翻阅 train/api 了解本教程所用类与方法的详细签名Ray Data 数据摄入当数据集规模变大时改用 Working with PyTorchRay Data 提供的流式数据摄入方案替代 PyTorchDataLoader。【免费下载链接】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),仅供参考