进阶 pytorch.org 2026-10-07 23:07:35 · 6 阅读

第2章 结合分布式 DataParallel 与分布式 RPC 框架

将 DistributedDataParallel 与 Distributed RPC Framework 结合使用 创建日期:2020 年 7 月 28 日 | 最后更新:2023 年 6 月 6 日 | 最后验证:未验证 作者:Pritam Damania 和 Yi Wang 注意 可以在 github 上查看并编辑本教程。 本教程通过一个简单的示例,演示如何将 DistributedDataParallel (DDP) 与 Distributed RPC framework 结合,把分布式数据并行和分布式模型并行结合起来训练一个简单模型。示例源代码可以在这里找到。 之前的教程《Getting Started With Distributed Data Parallel》和《Getting Started with Distributed RPC Framework》分别介绍了如何进行分布式数据并行训练和分布式模型并行训练。但在某些训练场景中,你可能需要将这两种技术结合使用,例如: 如果模型包含稀疏部分(大型 embedding 表)和稠密部分(全连接层),我们可以把 embedding 表放在参数服务器上,同时用 DistributedDataParallel 在多个训练节点之间复制全连接层。此时可以用 Distributed RPC framework 在参数服务器上执行 embedding 查询。 启用混合并行,如 PipeDream 论文所述。我们可以用 Distributed RPC framework 将模型的各个阶段(pipeline stage)分配到多个 worker 上形成流水线,并用 DistributedDataParallel 复制每个阶段(如有需要)。
本教程将介绍上述第一种情况。整个环境中共有 4 个 worker,如下所示: 1 个 Master,负责在参数服务器上创建 embedding 表(nn.EmbeddingBag),并驱动两个 trainer 上的训练循环。 1 个 Parameter Server,负责在内存中保存 embedding 表,并响应来自 Master 和 trainer 的 RPC 请求。 2 个 Trainer,各自保存一个全连接层(nn.Linear),并通过 DistributedDataParallel 在彼此之间复制。trainer 还负责执行前向传播、反向传播和优化器更新。
整个训练流程如下执行: 主节点创建一个 RemoteModule,该模块在参数服务器(Parameter Server)上托管嵌入表。 随后,主节点在训练器上启动训练循环,并将远程模块传递给训练器。 训练器创建一个 HybridModel。该模型首先使用主节点提供的远程模块执行嵌入查找,然后执行包裹在 DDP 内部的 FC 层。 训练器执行模型的前向传播,并基于损失函数使用分布式自动微分(Distributed Autograd)执行反向传播。 在反向传播过程中,首先计算 FC 层的梯度,并通过 DDP 中的 allreduce 同步到所有训练器。 接着,分布式自动微分将梯度传播到参数服务器,在此处更新嵌入表的梯度。 最后,使用分布式优化器(Distributed Optimizer)更新所有参数。 注意 如果你同时使用 DDP 和 RPC,反向传播时应始终使用分布式自动微分(Distributed Autograd)。 现在,让我们详细分解每个部分。首先,我们需要在执行任何训练之前初始化所有工作进程。我们创建 4 个进程,其中 rank 0 和 1 是训练器,rank 2 是主节点,rank 3 是参数服务器。 我们使用 TCP init_method 在所有 4 个工作进程上初始化 RPC 框架。 RPC 初始化完成后,主节点使用 RemoteModule 在参数服务器上创建一个托管 EmbeddingBag 层的远程模块。 然后,主节点遍历每个训练器,通过调用每个训练器的 _run_trainer(使用 rpc_async)来启动训练循环。 最后,主节点等待所有训练结束后再退出。 训练器首先使用 init_process_group 初始化 DDP 的 ProcessGroup,设置 world_size=2(对应两个训练器)。 接下来,它们使用 TCP init_method 初始化 RPC 框架。请注意,RPC 初始化和 ProcessGroup 初始化使用的端口不同,这是为了避免两个框架初始化时的端口冲突。 初始化完成后,训练器只需等待主节点发出的 _run_trainer RPC 调用。 参数服务器仅初始化 RPC 框架,并等待来自训练器和主节点的 RPC 调用。 ```python def run_worker(rank, world_size): r""" 封装函数,用于初始化 RPC、调用目标函数以及关闭 RPC。 """ # 在 TCP init_method 中,init_rpc 和 init_process_group 需要不同的端口号以避免冲突。 rpc_backend_options = TensorPipeRpcBackendOptions() rpc_backend_options.init_method = "tcp://localhost:29501" # Rank 2 是主节点,3 是参数服务器(ps),0 和 1 是训练器。 if rank == 2: rpc.init_rpc( "master", rank=rank, world_size=world_size, rpc_backend_options=rpc_backend_options, ) remote_emb_module = RemoteModule( "ps", torch.nn.EmbeddingBag, args=(NUM_EMBEDDINGS, EMBEDDING_DIM), kwargs={"mode": "sum"}, ) # 在训练器上运行训练循环。 futs = [] for trainer_rank in [0, 1]: trainer_name = "trainer{}".format(trainer_rank) fut = rpc.rpc_async( trainer_name, _run_trainer, args=(remote_emb_module, trainer_rank) ) futs.append(fut) # 等待所有训练完成。 for fut in futs: fut.wait() elif rank <= 1: # 在训练器上初始化分布式数据并行(DDP)的进程组。 dist.init_process_group( backend="gloo", rank=rank, world_size=2, init_method="tcp://localhost:29500" ) # 初始化 RPC。 trainer_name = "trainer{}".format(rank) rpc.init_rpc( trainer_name, rank=rank, world_size=world_size, rpc_backend_options=rpc_backend_options, ) # 训练器仅等待来自主节点的 RPC 调用。 else: rpc.init_rpc( "ps", rank=rank, world_size=world_size, rpc_backend_options=rpc_backend_options, ) # 参数服务器不执行其他操作 pass # 阻塞直到所有 RPC 完成 rpc.shutdown() if __name__ == "__main__": # 2 个训练器,1 个参数服务器,1 个主节点。 world_size = 4 mp.spawn(run_worker, args=(world_size,), nprocs=world_size, join=True) ``` 在讨论训练器的细节之前,让我们先介绍训练器使用的 HybridModel。如下所述,HybridModel 使用一个托管在参数服务器上的远程模块(remote_emb_module,包含嵌入表)以及用于 DDP 的设备进行初始化。模型的初始化过程将 nn.Linear 层包裹在 DDP 中,以便在所有训练器之间复制和同步该层。 模型的前向方法非常直观。它使用 RemoteModule 的 forward 方法在参数服务器上执行嵌入查找,并将其输出传递给 FC 层。 ```python class HybridModel(torch.nn.Module): r""" 模型由稀疏部分和稠密部分组成。 1) 稠密部分是一个 nn.Linear 模块,通过 DistributedDataParallel 在所有训练器之间复制。 2) 稀疏部分是一个 Remote Module,托管在参数服务器上的 nn.EmbeddingBag。 此远程模型可以获取参数服务器上嵌入表的远程引用(Remote Reference)。 """ def __init__(self, remote_emb_module, device): super(HybridModel, self).__init__() self.remote_emb_module = remote_emb_module self.fc = DDP(torch.nn.Linear(16, 8).cuda(device), device_ids=[device]) self.device = device def forward(self, indices, offsets): emb_lookup = self.remote_emb_module.forward(indices, offsets) return self.fc(emb_lookup.cuda(self.device)) ``` 接下来,让我们看看训练器的设置。训练器首先使用托管在参数服务器上的嵌入表远程模块及其自身的 rank 创建上述的 HybridModel。 现在,我们需要获取所有希望用 DistributedOptimizer 进行优化的参数的 RRef 列表。 为了从参数服务器获取嵌入表的参数,我们可以调用 RemoteModule 的 remote_parameters 方法。该方法基本上遍历嵌入表的所有参数并返回一个 RRef 列表。训练器通过 RPC 在参数服务器上调用此方法,以获取所需参数的 RRef 列表。由于 DistributedOptimizer 总是接收一个需要优化的参数 RRef 列表,因此即使是 FC 层的本地参数,我们也需要为它们创建 RRef。这是通过遍历 model.fc.parameters(),为每个参数创建 RRef 并追加到 remote_parameters() 返回列表中实现的。 请注意,我们不能使用 model.parameters(),因为它会递归调用 model.remote_emb_module.parameters(),而 RemoteModule 不支持该操作。 最后,我们使用所有 RRef 创建 DistributedOptimizer,并定义一个 CrossEntropyLoss 函数。 ```python def _run_trainer(remote_emb_module, rank): r""" 每个训练器运行前向传播,包括在参数服务器上执行嵌入查找以及在本地运行 nn.Linear。在反向传播期间,DDP 负责聚合稠密部分(nn.Linear)的梯度,分布式自动微分确保梯度更新传播到参数服务器。 """ # 设置模型。 model = HybridModel(remote_emb_module, rank) # 将所有模型参数作为 rrefs 获取,用于 DistributedOptimizer。 # 获取嵌入表的参数。 model_parameter_rrefs = model.remote_emb_module.remote_parameters() # model.fc.parameters() 仅包含本地参数。 # 注意:不能在此处调用 model.parameters(), # 因为这会调用 remote_emb_module.parameters(), # 该模块支持 remote_parameters() 但不支持 parameters()。 for param in model.fc.parameters(): model_parameter_rrefs.append(RRef(param)) # 设置分布式优化器 opt = DistributedOptimizer( optim.SGD, model_parameter_rrefs, lr=0.05, ) criterion = torch.nn.CrossEntropyLoss() ``` 现在,我们准备好介绍在每个训练器上运行的主训练循环了。get_next_batch 只是一个辅助函数,用于生成训练的随机输入和目标。我们运行多个 epoch 的训练循环,对于每个批次: * 为 Distributed Autograd 设置 Distributed Autograd Context。 * 运行模型的前向传播并获取输出。 * 使用损失函数基于输出和目标计算损失。 * 使用 Distributed Autograd 基于损失执行分布式反向传播。 * 最后,运行分布式优化器步骤以优化所有参数。 ```python def get_next_batch(rank): for _ in range(10): num_indices = random.randint(20, 50) indices = torch.LongTensor(num_indices).random_(0, NUM_EMBEDDINGS) # 生成 offsets。 offsets = [] start = 0 batch_size = 0 while start < num_indices: offsets.append(start) start += random.randint(1, 10) batch_size += 1 offsets_tensor = torch.LongTensor(offsets) target = torch.LongTensor(batch_size).random_(8).cuda(rank) yield indices, offsets_tensor, target # 训练 100 个 epoch for epoch in range(100): # 创建分布式自动微分上下文 for indices, offsets, target in get_next_batch(rank): with dist_autograd.context() as context_id: output = model(indices, offsets) loss = criterion(output, target) # 运行分布式反向传播 dist_autograd.backward(context_id, [loss]) # 运行分布式优化器 opt.step(context_id) # 无需清零梯度,因为每次迭代都会创建一个新的 # 分布式自动微分上下文,其中托管不同的梯度 print("Training done for epoch {}".format(epoch)) ``` 完整示例的源代码可以在这里找到。

评论 (0)