尧图精选

Lightning Fabric 集群环境插件(ClusterEnvironment)完全指南:自动检测、配置与分布式训练实战

🕒 发布时间:2026/9/19 5:28:27 📁 来源:尧图网络
Lightning Fabric 集群环境插件ClusterEnvironment完全指南自动检测、配置与分布式训练实战【免费下载链接】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导读本文围绕lightning.fabric.plugins.environments模块展开系统讲解 Lightning Fabric 如何在 SLURM、TorchElastic、MPI、LSF、Kubeflow、TPU Pod 以及裸机/单机等不同集群环境中自动识别并适配分布式训练上下文。读完本文你将掌握ClusterEnvironment抽象接口的设计意图、每种内置环境插件的检测逻辑与环境变量约定、手动指定环境插件的方法以及如何结合源码与测试验证环境配置的正确性从而在真实 HPC 或云集群上正确启动 Fabric 分布式训练任务。一、Environments 在 Lightning Fabric 中的定位在 Lightning Fabric 中ClusterEnvironment集群环境负责回答分布式训练中我是谁、我在哪、和谁通信这一类基础问题。它不关心模型如何定义、优化器如何配置只负责从运行环境中解析出四类关键信息进程拓扑world_size全局进程数、global_rank全局进程编号、local_rank节点内进程编号、node_rank节点编号通信端点main_address主节点地址、main_port主节点通信端口进程创建方式creates_processes_externally指示进程是由外部调度器如 SLURM、MPI、torchelastic创建还是需要 Lightning 自己 spawn环境识别与校验detect()判断当前环境是否匹配validate_settings()校验脚本配置与集群实际配置是否一致。该模块的公开 API 在 environments.rst 中通过 autosummary 列出 8 个符号涵盖 1 个抽象基类和 7 个内置实现类名适用场景ClusterEnvironment所有集群环境的抽象基类LightningEnvironment默认环境单机或非托管集群SLURMEnvironmentSLURM 调度器托管的 HPC 集群TorchElasticEnvironmenttorchelastic 弹性容错训练MPIEnvironment通过 MPImpi4py创建进程的集群LSFEnvironmentIBM LSF 资源管理器配合jsrunKubeflowEnvironmentKubeflowPyTorchJob算子XLAEnvironmentTPU PodPyTorch/XLA所有实现统一从 cluster_environment.py 的ClusterEnvironment继承并在 environments/init.py 中集中导出。二、ClusterEnvironment 抽象接口环境插件的统一契约ClusterEnvironment是一个纯抽象基类继承自abc.ABC它把任意集群环境收敛成一组必须实现的属性和方法见 cluster_environment.py。2.1 必须实现的抽象成员creates_processes_externally - bool环境是否由外部工具创建子进程。返回True表示进程已由调度器/启动器创建Lightning 不再自行 spawn返回False表示需要 Lightning 主动派生 worker 进程。main_address - str所有进程连接通信的主节点地址。main_port - int主节点上开放并配置好的通信端口。detect() - bool静态检测方法判断当前进程运行的环境是否与该实现匹配匹配返回True。world_size() - int所有设备、所有节点上的进程总数。set_world_size(size)/set_global_rank(rank)用于在需要 Lightning 自行创建进程的场景如LightningEnvironment回写拓扑信息外部调度器场景下通常被实现为忽略打印 debug 日志。global_rank() - int当前进程在所有节点和所有设备中的全局编号。local_rank() - int当前进程在所在节点内部的编号。node_rank() - int当前进程所在节点的编号。2.2 提供默认实现的可选钩子validate_settings(num_devices, num_nodes)将脚本中配置的devices/num_nodes与集群环境实际参数做一致性校验不一致时抛出异常基类默认空实现。teardown()训练结束后清理环境状态如清理临时设置的环境变量基类默认空实现。这一接口设计使得 Lightning Fabric 的策略层Strategy与如何解析集群信息完全解耦无论是 SLURM 还是 MPI策略代码只需面向ClusterEnvironment编程即可。三、环境自动检测与选择机制源码级环境插件最重要的特性是自动检测。在 Fabric 初始化时连接器connector会依次调用各环境插件的detect()命中即选中见 connector.pydef _choose_and_init_cluster_environment(self) - ClusterEnvironment: if isinstance(self._cluster_environment_flag, ClusterEnvironment): return self._cluster_environment_flag for env_type in ( # TorchElastic has the highest priority since it can also be used inside SLURM TorchElasticEnvironment, SLURMEnvironment, LSFEnvironment, MPIEnvironment, ): if env_type.detect(): return env_type() return LightningEnvironment()由源码可以明确两点关键设计优先级固定检测顺序为TorchElasticEnvironment → SLURMEnvironment → LSFEnvironment → MPIEnvironment → LightningEnvironment兜底。TorchElastic 优先级最高因为 torchelastic 可能被用作 SLURM 内部的启动器必须优先识别。兜底策略没有任何外部调度器特征时默认回退到LightningEnvironment保证单机训练开箱即用。手动覆盖如果用户显式传入了一个ClusterEnvironment实例self._cluster_environment_flag则直接使用该实例跳过自动检测——这正是 Kubeflow 环境的接入方式见第六节。四、LightningEnvironment单机与非托管集群的默认选择LightningEnvironment是默认且最常用的环境插件见 lightning.py。其文档明确说明它服务于单节点或自由非托管集群并支持两种运行模式用户只启动主进程如python train.py ...不设置任何环境变量Lightning 会在当前节点内自行 spawn 分布式训练的 worker 进程用户手动或借助torch.distributed.launch等工具启动所有进程此时需要设置相应环境变量最低限度必须设置LOCAL_RANK。4.1 关键实现细节进程创建判定creates_processes_externally直接返回LOCAL_RANK in os.environ。也就是说只要环境中存在LOCAL_RANK就认为用户自己扮演了进程启动器Lightning 不再派生进程。主地址读取MASTER_ADDR未设置时默认127.0.0.1。主端口读取MASTER_PORT若未设置则通过find_free_network_port()自动寻找空闲端口。世界大小/全局秩以内部状态为准初始world_size1、global_rank0由 Lightning spawn 逻辑通过set_world_size/set_global_rank写入其中set_global_rank还会同步更新rank_zero_only.rank保证 rank-0 日志过滤正确。节点内秩读取LOCAL_RANK未设置时为 0。节点编号优先读取NODE_RANK否则回退到GROUP_RANK再否则为 0。detect()恒为True它作为兜底环境永远匹配。teardown()若环境中存在WORLD_SIZE则在结束后将其删除避免污染后续进程。4.2 单机模式下的端口选择端口选择逻辑位于同一文件底部的find_free_network_port()lightning.py如果设置了STANDALONE_PORT环境变量则直接使用它否则绑定(, 0)让操作系统分配一个空闲端口。这适用于单节点训练中没有真实主节点但仍需设置MASTER_PORT的场景。五、SLURMEnvironmentHPC 集群的标准接入SLURM 是学术与工业 HPC 领域最常用的调度器SLURMEnvironment是 Lightning 支持最完善的环境插件之一见 slurm.py。5.1 构造参数SLURMEnvironment(auto_requeue: bool True, requeue_signal: Optional[signal.Signals] None)auto_requeue是否启用自动作业重提交job requeue默认为Truerequeue_signalSLURM 用于指示作业需要重新排队的信号Unix 下默认为SIGUSR1。5.2 环境变量解析约定SLURM 环境变量用途对应方法SLURM_NTASKS总任务数world sizeworld_size()SLURM_PROCID全局进程编号global_rank()SLURM_LOCALID节点内进程编号local_rank()SLURM_NODEID节点编号node_rank()SLURM_NODELIST节点列表用于推导主地址main_addressSLURM_JOB_ID作业 ID用于推导主端口main_portSLURM_JOB_NAME作业名用于交互模式判定job_name()MASTER_ADDR/MASTER_PORT显式覆盖主地址/端口main_address/main_port5.3 主地址与主端口的自动推导主地址优先读取MASTER_ADDR未设置时从SLURM_NODELIST中解析出第一个节点作为主节点并写回MASTER_ADDR。节点列表的解析由resolve_root_node_address()完成支持三种 SLURM 格式空格分隔的主机名列表如host0 host1 host3→ 取host0逗号分隔列表如host0,host1,host3→ 取host0方括号范围记法如host[5-9]→ 取host5。主端口使用SLURM_JOB_ID的后 4 位加上 15000 作为默认端口保证端口在 10k 范围避免撞车例如作业 ID0001234得到端口1234 15000 16234若设置了MASTER_PORT则优先使用并写回环境变量供所有进程共享。5.4 检测与常见错误防御detect()返回_is_srun_used()即环境中存在SLURM_NTASKS且作业名不是bash/interactive。若用户用srun --job-nameinteractive进入交互模式则不匹配 SLURM 环境从而允许其他环境接管。_validate_srun_used()检测到系统安装了srun但启动命令没有以srun开头时会发出PossibleUserWarning并给出形如srun python train.py ...的修正提示——这是 SLURM 下多卡/多节点任务挂起的常见原因。_validate_srun_variables()当SLURM_NTASKS 1却未设置SLURM_NTASKS_PER_NODE时直接抛RuntimeError提示改用--ntasks-per-node。validate_settings()非交互模式下校验SLURM_NTASKS_PER_NODE与脚本中devices是否一致、SLURM_NNODES与num_nodes是否一致不一致时抛出带修正建议的ValueError。5.5 典型 SLURM 启动脚本#!/bin/bash #SBATCH --nodes2 #SBATCH --ntasks-per-node4 #SBATCH --gresgpu:4 #SBATCH --job-nametrain_fabric srun python train_fabric.py --strategy ddp --devices 4 --num_nodes 2在仓库的 tests/tests_fabric/plugins/environments/test_slurm.py 中可以验证上述行为当未设置任何环境变量时main_address默认为127.0.0.1、main_port默认为12910、job_name()与job_id()均为None且访问world_size()、local_rank()、node_rank()会因缺少环境变量抛出KeyError设置SLURM_NODELIST、SLURM_JOB_ID、SLURM_NTASKS等变量后各属性即按上文规则解析。六、其余内置环境插件速览6.1 TorchElasticEnvironment弹性容错训练用于torch.distributed.elastictorchelastic启动的容错训练见 torchelastic.py。detect()通过torch.distributed.is_torchelastic_launched()判定。它要求WORLD_SIZE、RANK、LOCAL_RANK必须存在MASTER_ADDR/MASTER_PORT缺失时给出警告并回退到127.0.0.1/12910node_rank取自GROUP_RANK。其validate_settings()会校验devices * num_nodes是否等于world_size()不一致即抛错。6.2 MPIEnvironmentMPI 进程环境面向mpirun等 MPI 启动器前提是安装mpi4pympi.py构造时若mpi4py不可用则抛ModuleNotFoundErrordetect()在mpi4py可用且MPI.COMM_WORLD.Get_size() 1时返回Truempi4py 已安装但无 MPI 运行时则返回Falseworld_size、global_rank、local_rank分别取自COMM_WORLD.Get_size()、Get_rank()与本地通信子COMM_LOCAL.Get_rank()并使用lru_cache缓存main_address由 rank 0 通过bcast(socket.gethostname(), root0)广播得到main_port由 rank 0 广播一个空闲端口node_rank通过COMM_WORLD.gather收集各进程主机名、排序去重后定位本机索引并用COMM_WORLD.Split(colornode_rank)建立本地通信子。6.3 LSFEnvironmentIBM LSF 资源管理器面向 LSF 且要求通过 Job Step Managerjsrun执行见 lsf.py。它依赖四类环境变量LSB_JOBID作业 ID、LSB_DJOB_RANKFILEOpenMPI 兼容的 rank 文件、JSM_NAMESPACE_LOCAL_RANK本地秩、JSM_NAMESPACE_SIZE世界大小detect()检查这四个变量是否全部存在。其global_rank读取JSM_NAMESPACE_RANKnode_rank依据当前主机名在 rank 文件中的位置计算rank 文件第一个节点为 launch 节点会被剔除main_address取 rank 文件中的首个计算节点main_port默认由LSB_JOBID % 1000 10000计算得出可用MASTER_PORT覆盖并在初始化时把MASTER_ADDR/MASTER_PORT写回环境变量供 torch.distributed 进程组初始化使用。6.4 KubeflowEnvironmentKubeflow PyTorchJob面向 KubeflowPyTorchJob算子kubeflow.py。它有两个显著特点无法自动检测detect()直接raise NotImplementedError因此必须手动传入Fabric/Trainer 构造器这是与其它环境插件最大的不同进程布局约定creates_processes_externally为Truelocal_rank()恒为0node_rank()返回global_rank()——即 Kubeflow 的PyTorchJob每个 Pod 视为一个节点Pod 内部单进程。手动指定方式from lightning.fabric import Fabric from lightning.fabric.plugins.environments import KubeflowEnvironment fabric Fabric(strategyddp, plugins[KubeflowEnvironment()])6.5 XLAEnvironmentTPU PodPyTorch/XLA面向 TPU Pod 训练见 xla.py。其detect()委托给XLAAccelerator.is_available()构造时若torch_xla不可用则抛ModuleNotFoundError。与其它环境不同XLAEnvironment.creates_processes_externally返回False进程由 Lightning 管理main_address/main_port抛出NotImplementedErrorXLA 场景不使用这两者。world_size、global_rank、local_rank、node_rank在 PyTorch/XLA ≥ 2.1 时通过torch_xla.runtime的world_size()、global_ordinal()、local_ordinal()、host_index()获取旧版本则回退到torch_xla.core.xla_model与xla_env_vars且全部使用lru_cache缓存结果。在 Fabric 中XLA 策略会直接构造该环境见 strategies/xla.py 与 strategies/xla_fsdp.py。七、环境插件与策略、加速器的协作ClusterEnvironment是 Strategy 构造的必需依赖。从源码可以确认DDP 系策略strategies/ddp.py、strategies/fsdp.py、strategies/deepspeed.py都接收cluster_environment参数并从中读取地址、端口、rank 等来初始化 torch.distributed 进程组dp.pyDataParallel 策略则显式传入cluster_environmentNonestrategies/dp.py因为单进程多设备策略不需要集群拓扑XLA 策略固定使用XLAEnvironment与 TPU 加速器强绑定。因此理解环境的解析规则本质上就理解了 DDP/FSDP/DeepSpeed 等分布式策略在目标集群上的初始化前提。八、自定义集群环境实现自己的 ClusterEnvironment当目标平台不在内置列表时可以继承ClusterEnvironment实现自定义插件再手动传给 Fabricfrom lightning.fabric import Fabric from lightning.fabric.plugins.environments import ClusterEnvironment class MyClusterEnvironment(ClusterEnvironment): property def creates_processes_externally(self) - bool: return True property def main_address(self) - str: return master.host property def main_port(self) - int: return 29500 staticmethod def detect() - bool: return MY_CLUSTER in __import__(os).environ def world_size(self) - int: return 8 def global_rank(self) - int: return 0 def local_rank(self) - int: return 0 def node_rank(self) - int: return 0 fabric Fabric(strategyddp, plugins[MyClusterEnvironment()])实现时需要保证creates_processes_externally、main_address、main_port、detect、world_size、global_rank、local_rank、node_rank八个成员全部提供正确语义set_world_size/set_global_rank在外部创建进程的场景下可直接忽略validate_settings/teardown可按需覆盖。九、常见问题与排查建议SLURM 任务挂起启动命令没有以srun开头。检查日志中的PossibleUserWarning按提示将命令改为srun python ...。devices / num_nodes 不匹配SLURM 下--ntasks-per-node与脚本devices、--nodes与num_nodes必须一致否则validate_settings()会抛出带修正建议的ValueErrortorchelastic 下则是devices * num_nodes ! world_size()。Kubeflow 未生效Kubeflow 环境无法自动检测必须显式plugins[KubeflowEnvironment()]传入。MPI 环境未识别确认已安装mpi4py且使用mpirun启动mpi4py存在但无 MPI 运行时detect()返回False。端口冲突SLURM 默认端口由作业 ID 推导LSF 由作业 ID 取模推导如需固定显式设置MASTER_PORT可覆盖所有环境插件的默认推导。单机多卡优先使用默认LightningEnvironment直接python train.py --devices 4 --strategy ddp即可由 Fabric 自动选择空闲端口并 spawn 进程。十、参考资料模块 API 文档docs/source-fabric/api/environments.rst抽象基类实现src/lightning/fabric/plugins/environments/cluster_environment.py各环境实现目录src/lightning/fabric/plugins/environments/环境自动选择逻辑src/lightning/fabric/connector.py环境测试用例tests/tests_fabric/plugins/environments/【免费下载链接】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),仅供参考
上一篇/下一篇内容由系统自动关联 返回资讯列表 →