PyTorch Lightning GPU 分布式训练实战(中级):DDP、DDP Spawn、DDP Notebook 与 TorchRun 多机扩展指南
【免费下载链接】pytorch-lightningPretrain, finetune ANY AI model of ANY size on 1 or 10,000+ GPUs with zero code changes.项目地址: https://gitcode.com/gh_mirrors/py/pytorch-lightning
本文聚焦 PyTorch Lightning 在 GPU 场景下的分布式训练策略,面向已能单卡训练、希望在多卡或多机间横向扩展的开发者。你将完整掌握DDP、DDP Spawn、DDP Notebook/Fork三种策略的原理、适用场景与取舍,学会用torchrun启动容错的多机训练,并能显式控制nccl等进程组后端以优化多机通信。文中所有结论均可对照仓库源码验证。
原文出处:docs/source-pytorch/accelerators/gpu_intermediate.rst,本文在此基础上结合
src/lightning/pytorch/strategies源码、Trainer 连接器 与 torchrun 集群文档 进行了纵深扩充。
三种分布式训练策略速览
PyTorch Lightning 的 Trainer 支持多种分布式训练方式,按启动进程的机制不同,主要分为三类:
- Regular(标准 DDP):
strategy="ddp",以子进程脚本方式启动,是生产环境的首选。 - Spawn:
strategy="ddp_spawn",基于torch.multiprocessing.spawn()启动,仅建议用于调试或迁移旧代码。 - Notebook / Fork:
strategy="ddp_notebook"(别名ddp_fork),基于进程fork机制,专为 Jupyter Notebook、Google Colab、Kaggle 等交互式环境设计。
关键行为:如果你请求了多个 GPU 或多个节点但没有显式设置 strategy,Lightning 会自动使用 DDP。该逻辑体现在 accelerator_connector.py 的_choose_strategy方法中:num_nodes > 1时返回"ddp",多设备且处于交互环境时返回"ddp_fork",其余情况返回"ddp"。
Distributed Data Parallel:标准 DDP 的工作原理与用法
DDP(torch.nn.parallel.DistributedDataParallel)在 Lightning 中的执行流程如下:
- 每个节点上的每张 GPU 各自拥有一个独立进程;
- 每张 GPU 只看到整个数据集的一个子集,且只处理这个子集(由分布式采样器保证划分);
- 每个进程各自初始化模型;
- 每个进程并行地执行完整的前向与反向传播;
- 所有进程的梯度被同步并求平均;
- 每个进程用平均后的梯度更新自己的优化器。
对应的 Trainer 配置非常简洁:
# 单机 8 卡 trainer = Trainer(accelerator="gpu", devices=8, strategy="ddp") # 4 节点共 32 卡 trainer = Trainer(accelerator="gpu", devices=8, strategy="ddp", num_nodes=4)注意 Lightning 2.x 中使用accelerator="gpu"与devices=N的组合,等价于早期版本的gpus=N。num_nodes必须为正整数,否则连接器会抛出ValueError(见 accelerator_connector.py)。
Lightning 是如何启动这些进程的
与用户手动spawn不同,Lightning 的 DDP 实现会在底层用正确的环境变量多次调用你的脚本。以 3 卡 DDP 为例,实际效果等价于:
MASTER_ADDR=localhost MASTER_PORT=random() WORLD_SIZE=3 NODE_RANK=0 LOCAL_RANK=0 python my_file.py --accelerator 'gpu' --devices 3 --etc MASTER_ADDR=localhost MASTER_PORT=random() WORLD_SIZE=3 NODE_RANK=0 LOCAL_RANK=1 python my_file.py --accelerator 'gpu' --devices 3 --etc MASTER_ADDR=localhost MASTER_PORT=random() WORLD_SIZE=3 NODE_RANK=0 LOCAL_RANK=2 python my_file.py --accelerator 'gpu' --devices 3 --etc该机制由 _SubprocessScriptLauncher 实现:主进程(LOCAL_RANK=0)通过subprocess.Popen为其余设备各启动一个子进程,并逐一设置MASTER_ADDR(主节点 IP)、MASTER_PORT(进程通信端口)、NODE_RANK(节点索引,0 到num_nodes-1)、LOCAL_RANK(节点内进程索引)与WORLD_SIZE(所有节点的进程总数,即num_processes * num_nodes)。源码中的 docstring 以python train.py --devices 4为例,说明它会额外创建LOCAL_RANK=1/2/3三个子进程。这个启动方式与torch.distributed.run的进程模型非常相似。
相比torch.multiprocessing.spawn()的优势
使用这种"重复调用脚本"的 DDP 方式,相比手动spawn有几个明确优点:
- 所有进程(包括主进程)都参与训练,主进程持有最新的模型与 Trainer 状态;
- 不存在 multiprocessing pickle 错误,因为模型无需通过队列回传;
- 天然支持多节点扩展,只要节点间网络互通即可。
重要限制:不能用于交互式环境
标准 DDP 依赖"重复执行整个脚本",因此无法在 Jupyter Notebook、Google Colab、Kaggle 等交互式环境中使用。Lightning 会在检测到交互环境时给出明确提示:若当前策略不兼容交互环境,会抛出MisconfigurationException,并建议改用Trainer(strategy="ddp_notebook")(见 accelerator_connector.py)。交互环境的判定来自sys.ps1是否存在或sys.flags.interactive是否为真(见 src/lightning/fabric/utilities/imports.py)。
DDP Spawn:调试与迁移场景的备选
ddp_spawn与标准 DDP 的唯一区别在于:它使用torch.multiprocessing.spawn()来启动训练进程。
# 单机 8 卡 trainer = Trainer(accelerator="gpu", devices=8, strategy="ddp_spawn")文档明确警告:强烈建议优先使用 DDP 以获得更快的速度与更好的性能。ddp_spawn只应在调试时使用,或用于迁移那些原本依赖 spawn 机制的旧代码库。不推荐它的原因(受 Python 与 PyTorch 机制限制):
.fit()结束后,只有模型的权重会被回传到主进程,Trainer 的其他状态(优化器状态、调度器状态、epoch 计数等)不会同步;- 不支持多节点训练;
- 总体而言比 DDP 更慢。
从源码看,spawn 与标准 DDP 共用同一个DDPStrategy类,只是start_method不同:ddp.py 中的register_strategies将"ddp"映射到start_method="popen"、将"ddp_spawn"映射到start_method="spawn",并在_configure_launcher中根据 start method 选择 _SubprocessScriptLauncher 或_MultiProcessingLauncher。因此ddp_spawn本质上就是"用 spawn 方式启动的 DDP",其功能边界由 Python 多进程的 spawn 语义决定。
DDP Notebook / Fork:交互式环境的多卡方案
DDP Notebook/Fork 是 Spawn 的替代方案,可在交互式 Python、Jupyter Notebook、Google Colab、Kaggle 等环境中使用。Trainer 在检测到此类环境时会默认启用它(对应_choose_strategy中len(self._parallel_devices) > 1 and _IS_INTERACTIVE时返回"ddp_fork"的分支)。
# Jupyter notebook 中训练 8 卡(自动启用 fork) trainer = Trainer(accelerator="gpu", devices=8) # 也可以显式指定 trainer = Trainer(accelerator="gpu", devices=8, strategy="ddp_notebook") # 非交互环境也可以显式使用 trainer = Trainer(accelerator="gpu", devices=8, strategy="ddp_fork")从注册表看,"ddp_notebook"与"ddp_fork"的start_method均为"fork",二者本质是同一个机制,notebook只是语义化的别名(见 ddp.py 与 ddp.py)。另外注意 accelerator_connector.py 的校验:如果平台不支持fork启动方法(如 Windows),选择 fork 别名会抛出ValueError并建议改用ddp_spawn。
在原生分布式策略中,标准 DDP(strategy="ddp")仍然是速度与稳定性俱佳的首选策略,只是它只能用于脚本形式运行。Fork/Notebook 在交互场景下还有一条额外限制:在调用Trainer.fit之前,不允许执行"把张量移到 GPU"或调用torch.cuda函数等 GPU 操作(见下文的对比表)。
DDP 三种变体的对比与取舍
| 特性 | DDP | DDP Spawn | DDP Notebook/Fork |
|---|---|---|---|
| 可在 Jupyter / IPython 环境中工作 | 否 | 否 | 是 |
| 支持多节点 | 是 | 是 | 是 |
| 支持的平台 | Linux、Mac、Win | Linux、Mac、Win | Linux、Mac |
| 要求所有对象可 pickle | 否 | 是 | 否 |
| 主进程中的限制 | 无 | 返回主进程后对象状态不是最新的(Trainer.fit()等),只有模型参数被回传 | 不允许在调用Trainer.fit之前执行 GPU 操作(如将张量移到 GPU、调用torch.cuda函数) |
| 进程创建时间 | 慢 | 慢 | 快 |
需要说明的是,表中"进程创建时间"的"慢/快"指的是子进程的启动开销:popen/spawn都需要重新初始化 Python 解释器,而fork通过复制当前进程内存实现,因此创建更快,这也是它适合 notebook 场景的原因之一。
用 TorchRun(TorchElastic)实现弹性多机调度
Lightning 支持 TorchRun(原 TorchElastic),用于实现容错(fault-tolerant)与弹性(elastic)的分布式任务调度。用法分两步:
第一步:在 Trainer 中指定 DDP 策略与 GPU 数量:
Trainer(accelerator="gpu", devices=8, strategy="ddp")第二步:用torchrun命令启动脚本。详细的参数说明见 集群文档,其核心命令模板为:
torchrun \ --nproc_per_node=<GPUS_PER_NODE> \ --nnodes=<NUM_NODES> \ --node_rank <NODE_RANK> \ --master_addr <MASTER_ADDR> \ --master_port <MASTER_PORT> \ train.py --arg1 --arg2各参数含义与约束:
--nproc_per_node:每个节点启动的进程数(默认 1),必须与Trainer(devices=...)中设置的数值一致;--nnodes:参与训练的节点/机器数(默认 1),必须与Trainer(num_nodes=...)中设置的数值一致;--node_rank:当前节点的索引(从 0 开始);--master_addr:rank 0 主节点的 IP 地址;--master_port:节点间通信端口,必须在每个节点的防火墙上开放 TCP 流量。
两节点各 8 卡的实际示例
假设主节点 IP 为10.10.10.16、通信端口为50000。第一个节点(node rank 0)上运行:
torchrun \ --nproc_per_node=8 --nnodes=2 --node_rank 0 \ --master_addr 10.10.10.16 --master_port 50000 \ train.py第二个节点(node rank 1)上运行:
torchrun \ --nproc_per_node=8 --nnodes=2 --node_rank 1 \ --master_addr 10.10.10.16 --master_port 50000 \ train.py注意两条命令唯一的区别就是--node_rank。除 Lightning 侧的配置外,你还需自行保证节点间的网络连通性(防火墙放行MASTER_PORT)、确定主节点(MASTER_ADDR)以及各节点的 rank。torchrun会从这些参数自动生成 PyTorch 分布式初始化所需的MASTER_ADDR、MASTER_PORT、RANK、WORLD_SIZE等环境变量。关于弹性、容错等更高级的配置,可查阅 torchrun 官方文档。
从源码角度印证:_choose_and_init_cluster_environment(见 accelerator_connector.py)会按优先级依次检测TorchElasticEnvironment、SLURMEnvironment、LSFEnvironment、MPIEnvironment,其中 TorchElastic 优先级最高,因为它也能在 SLURM 内部使用;当这些外部集群环境都不存在时才回退到LightningEnvironment。因此在torchrun启动时,Lightning 会自动识别 TorchElastic 提供的环境变量,_SubprocessScriptLauncher也会因creates_processes_externally为真而跳过重复创建子进程,直接把所有进程交给 torchrun 管理。
优化多机通信:显式指定进程组后端
默认情况下,在 GPU 上运行时 Lightning 会选择nccl后端而非gloo。PyTorch 的分布式包支持多种后端,各有适用的硬件与场景。
Lightning 允许你通过策略类(Strategy class)的构造参数process_group_backend显式指定后端;未指定时,Lightning 会根据当前硬件自动选择合适的后端:
from lightning.pytorch.strategies import DDPStrategy # 显式指定进程组后端 ddp = DDPStrategy(process_group_backend="nccl") # 将策略配置到 Trainer trainer = Trainer(strategy=ddp, accelerator="gpu", devices=8)源码级原理:默认后端是如何决定的
在 DDPStrategy.setup_distributed 中,Lightning 调用_get_process_group_backend():若用户显式传入了process_group_backend则直接使用,否则调用_get_default_process_group_backend_for_device按设备类型查表(见 src/lightning/fabric/utilities/distributed.py):该函数读取torch.distributed.Backend.default_device_backend_map,CUDA 设备对应nccl,其他未登记的设备类型回退到"gloo"。也就是说,文档所说的"GPU 上默认 nccl"实际上是委托给 PyTorch 的设备-后端映射表完成的。
确定后端后,setup_distributed会完成reset_seed()、通过set_world_ranks()计算global_rank = node_rank * num_processes + local_rank与world_size = num_nodes * num_processes,最后调用_init_dist_connection初始化进程组。若你想调整进程组超时时间,还可以在构造DDPStrategy时传入timeout(默认来自default_pg_timeout)。
DDP 构造参数速查
DDPStrategy的完整构造签名见 ddp.py,常用参数包括:
process_group_backend:进程组后端(如"nccl"、"gloo"),默认None表示按硬件自动选择;start_method:进程启动方式,"popen"(标准 DDP)、"spawn"(DDP Spawn)、"fork"/"forkserver"(Notebook/Fork),默认"popen";timeout:进程组通信超时,类型为datetime.timedelta;ddp_comm_state/ddp_comm_hook/ddp_comm_wrapper:DDP 通信钩子相关,用于自定义梯度通信(如 post-localSGD 模型平均化);model_averaging_period:与 post-localSGD 配合的模型平均周期;- 其他
**kwargs会直接透传给torch.nn.parallel.DistributedDataParallel(例如find_unused_parameters=True)。
此外,注册表还提供了若干快捷字符串策略,均可在Trainer(strategy=...)中直接使用(见 ddp.py):
ddp_find_unused_parameters_true/ddp_find_unused_parameters_false:显式控制find_unused_parameters;ddp_spawn_find_unused_parameters_true/..._false:spawn 版本;ddp_fork_find_unused_parameters_true/..._false与ddp_notebook_find_unused_parameters_true/..._false:fork 版本。
当模型存在未参与 loss 计算的参数时,DDP 的梯度归约会报错。此时若该行为是有意为之,应通过strategy='ddp_find_unused_parameters_true'或DDPStrategy(find_unused_parameters=True)显式启用未使用参数检测——这是DDPStrategy.on_exception中对这类错误的标准处置建议(见 ddp.py)。上述注册表条目均有对应的单元测试覆盖,可参考 tests/tests_pytorch/strategies/test_registry.py 中的参数化测试。
实战建议与排错清单
- 脚本训练选 DDP:只要以
python train.py方式运行,就用strategy="ddp"(或直接省略 strategy),兼顾速度与稳定性; - notebook 训练用 ddp_notebook:交互环境下 Lightning 会自动切换,也可显式指定;fork 启动前不要先执行 GPU 相关操作;
- 调试/迁移旧代码才用 ddp_spawn:接受"只有权重回传、不支持多节点、更慢"这三条限制;
- 多机训练确认三件事:
--nproc_per_node与devices一致、--nnodes与num_nodes一致、--master_port在节点防火墙放行; - 遇到梯度归约报错:确认模型是否有未参与 loss 计算的参数,按需开启
find_unused_parameters; - 希望自定义通信:用
DDPStrategy(process_group_backend=...)显式指定后端与超时。
通过本文的组合配置,你可以从单机单卡平滑扩展到单机多卡、再到多机多卡的分布式训练,同时兼顾交互式环境与弹性调度的实际需求。
【免费下载链接】pytorch-lightningPretrain, finetune ANY AI model of ANY size on 1 or 10,000+ GPUs with zero code changes.项目地址: https://gitcode.com/gh_mirrors/py/pytorch-lightning
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考