PyTorch 分布式训练先跑通最小双进程闭环
2026/8/13 13:16:47 网站建设 项目流程

PyTorch 分布式训练先跑通最小双进程闭环

分布式训练不必一开始就上多机。两个本地进程已经足够暴露数据分片、梯度同步、进程退出和 checkpoint 协调中的不少问题。

1. 双进程闭环要验证什么

训练问题应拆成数值正确性、数据供给、显存使用和通信行为四部分。先以小规模、固定输入验证前向和反向结果,再观察多进程路径,避免把单一监控值当成整体结论。

先以固定输入对齐单卡,再加入第二个进程。模型或数据路径变化后,这条基线也要重新跑。

2. 一次只引入一个分布式变量

每次试验都应写清框架版本、设备类型、批量形状、随机种子和启动方式。发生偏差时,优先比较中间张量与梯度,而不是直接调整并行参数。

先验证 sampler,再验证梯度,最后处理恢复与异常退出。每一步保存配置和张量摘要,不复制原始训练数据。

3. 双进程启动示例

以下片段保留原有技术结构。运行前请替换为本地的非敏感示例,并根据依赖版本核对接口。

[磁盘 SSD] --(读取原始图片/文本)--> [CPU Host 内存] --(Transforms 预处理)--> [PCIe 总线] --(H2D 搬运)--> [GPU 显存]
import os import time import torch import torch.nn as nn import torch.distributed as dist from torch.nn.parallel import DistributedDataParallel as DDP from torch.utils.data import Dataset, DataLoader, DistributedSampler class SyntheticDataset(Dataset): """模拟的高效离线数据集""" def __init__(self, size: int = 10000, features: int = 128): self.size = size self.data = torch.randn(size, features) self.labels = torch.randint(0, 2, size=(size,)) def __len__(self) -> int: return self.size def __getitem__(self, idx: int): return self.data[idx], self.labels[idx] def cleanup_ddp(): """安全销毁分布式进程组""" if dist.is_initialized(): dist.destroy_process_group() def run_ddp_training_node(rank: int, world_size: int, epochs: int = 2): """ 单节点 DDP 最小化训练进程 """ # 1. 设置环境变量与初始化 NCCL 进程组 os.environ['MASTER_ADDR'] = 'localhost' os.environ['MASTER_PORT'] = '29500' try: dist.init_process_group(backend='nccl', rank=rank, world_size=world_size) torch.cuda.set_device(rank) print(f"[Rank {rank}] DDP 进程组初始化成功!") # 2. 构建模型并移至对应 GPU 设备 model = nn.Sequential( nn.Linear(128, 512), nn.ReLU(), nn.Linear(512, 2) ).to(rank) # 用 DDP 包装模型 model = DDP(model, device_ids=[rank]) criterion = nn.CrossEntropyLoss().to(rank) optimizer = torch.optim.AdamW(model.parameters(), lr=1e-3) scaler = torch.cuda.amp.GradScaler() # 混合精度 Scaler # 3. 配置分布式采样器与 DataLoader (关键性能点: pin_memory 与 num_workers) dataset = SyntheticDataset() sampler = DistributedSampler(dataset, num_replicas=world_size, rank=rank, shuffle=True) loader = DataLoader( dataset, batch_size=64, sampler=sampler, num_workers=2, # 启动 CPU 子进程并行预加载 pin_memory=True, # 启用锁页内存,加速 CPU 到 GPU 传输 drop_last=True ) # 4. 训练循环 model.train() for epoch in range(epochs): sampler.set_epoch(epoch) # 保证每个 Epoch 的 Shuffle 随机种子不同 epoch_loss = 0.0 t0 = time.time() for inputs, targets in loader: # 高效异步非阻塞传输到 GPU inputs = inputs.to(rank, non_blocking=True) targets = targets.to(rank, non_blocking=True) optimizer.zero_grad() # 开启 自动混合精度 (AMP) 前向计算 with torch.cuda.amp.autocast(): outputs = model(inputs) loss = criterion(outputs, targets) # AMP 梯度缩放与反向传播 scaler.scale(loss).backward() scaler.step(optimizer) scaler.update() epoch_loss += loss.item() elapsed = time.time() - t0 if rank == 0: print(f"[Epoch {epoch+1}/{epochs}] 完成! 平均 Loss: {epoch_loss/len(loader):.4f}, 耗时: {elapsed:.2f}s") except Exception as e: print(f"[Rank {rank}] 运行过程抛出异常: {e}") finally: cleanup_ddp() print(f"[Rank {rank}] 进程安全退出并销毁资源。") # 单进程启动模拟 (可由 torchrun 拉起) if __name__ == "__main__": if torch.cuda.is_available() and torch.cuda.device_count() >= 1: # 如果只有单卡,用 world_size=1 验证逻辑 world_size = min(2, torch.cuda.device_count()) print(f"[INFO] 开始在 {world_size} 个 GPU 设备上拉起 DDP 架构...") torch.multiprocessing.spawn( run_ddp_training_node, args=(world_size, 2), nprocs=world_size, join=True ) else: print("[WARNING] 未检测到可用的 GPU,跳过 DDP 物理执行。")
# 梯度累积场景下避免无谓的 DDP 通信 for i, (inputs, targets) in enumerate(loader): # 前 3 个 Step 使用 no_sync() 禁用梯度 AllReduce if (i + 1) % accum_steps != 0: with model.no_sync(): outputs = model(inputs) loss = criterion(outputs, targets) / accum_steps loss.backward() else: # 第 4 个 Step 触发真正的同步 outputs = model(inputs) loss = criterion(outputs, targets) / accum_steps loss.backward() optimizer.step() optimizer.zero_grad()

4. 扩到多机前复核

  • 每个 rank 的样本集合是否符合切分规则。
  • 同步后的梯度与单卡参考是否在容差内。
  • 任一进程失败时其他进程能否退出。
  • checkpoint 是否包含恢复所需的完整状态。

总结

双进程闭环跑稳后再扩节点,问题空间会小很多。跳过这一步,集群日志往往只会把根因埋得更深。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询