ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

PyTorch 分布式 RPC 框架深入指南:跨机器模型训练的原语、RRef 与分布式自动微分

PyTorch 分布式 RPC 框架深入指南:跨机器模型训练的原语、RRef 与分布式自动微分 PyTorch 分布式 RPC 框架深入指南跨机器模型训练的原语、RRef 与分布式自动微分【免费下载链接】pytorchTensors and Dynamic neural networks in Python with strong GPU acceleration项目地址: https://gitcode.com/GitHub_Trending/py/pytorch本文以 PyTorch 官方文档 分布式 RPC 框架 为核心系统讲解torch.distributed.rpc提供的四大类分布式编程原语RPC 远程调用、RRef 远程引用、分布式自动微分Distributed Autograd与分布式优化器Distributed Optimizer并结合当前仓库源码torch/distributed/rpc/、torch/distributed/optim/optimizer.py、torch/distributed/nn/api/remote_module.py与两篇设计笔记RRef 协议、分布式自动微分设计剖析底层实现原理。读完本文你将掌握如何初始化 RPC 框架、如何选择同步/异步/远程引用三种调用方式、如何配置 TensorPipe 后端、如何利用 RRef 与 RemoteModule 组织跨节点模型以及如何用分布式自动微分和分布式优化器完成端到端的分布式训练。注意根据 rpc.mdRPC 包中的 API 已进入稳定维护模式stable and in maintenance modeCUDA 支持仍为beta功能且 RRefs、JIT 兼容、dist autograd、dist optimizer 与 profiling 等特性目前与 CUDA 支持不完全兼容官方不推荐在 CUDA 场景下使用这些特性。概述一套覆盖四种场景的分布式原语分布式 RPC 框架为多机模型训练提供了远程通信原语以及一个能够自动对跨多台机器切分的模型求导的高层 API。其全部能力可归纳为四组 APIRemote Procedure CallRPC在指定目标 worker 上以给定参数运行某个函数并取回返回值或创建指向返回值的引用。对应三个核心 API——rpc_sync同步、rpc_async异步、remote异步且返回远程引用的句柄。Remote ReferenceRRef一种指向本地或远端对象的分布式共享指针可在多个 worker 间共享引用计数由框架透明维护。每个 RRef 只有一个 owner所有者真实对象只存在于 owner 上非 owner worker 可通过显式请求从 owner 获取对象的拷贝。Distributed Autograd将参与前向传播的所有 worker 上的本地 autograd 引擎缝合起来在反向传播时自动触达各节点计算梯度用户无需关心梯度如何跨 RPC 边界传输、以及嵌套且相互依赖的 RPC 调用下本地引擎应以何种顺序启动。Distributed Optimizer构造函数接收一个torch.optim.Optimizer类如 SGD、Adagrad与一组参数 RRef 列表在每个不同的 RRef owner 上创建一个优化器实例并在调用step()时更新对应参数。当前向、反向跨多机分布时参数与梯度散落在多个 worker 上分布式优化器将这些本地优化器统一封装对外只暴露简洁的构造函数与step()API。RPC 基础初始化框架在使用 RPC 与分布式自动微分原语之前必须完成框架初始化。torch.distributed.rpc.init_rpc会一次性初始化 RPC 框架、RRef 框架与分布式自动微分三套子系统import torch.distributed.rpc as rpc rpc.init_rpc( name, # 当前节点的全局唯一名称 backendNone, # BackendType 枚举默认为 TENSORPIPE rank-1, # 当前节点的全局唯一 rank world_sizeNone, # 组内 worker 总数 rpc_backend_optionsNone, # RpcBackendOptions 的子类实例 )从 torch/distributed/rpc/init.py 的源码可以看到init_rpc的完整签名。其内部行为要点如下后端自动推断若只传rpc_backend_options而未显式指定backend框架会遍历BackendType枚举按选项类型反推后端见init.py#L141-L166若二者都未指定则默认使用BackendType.TENSORPIPE见init.py#L168-L175。初始化方法默认使用init_method env://即需要正确设置环境变量MASTER_ADDR与MASTER_PORT完成 rendezvous汇合init_rpc通过底层进程组完成握手并基于PrefixStore隔离多次调用见init.py#L177-L198。初始化顺序源码中先初始化分布式自动微分、再初始化 RPC agent并用同一store保证所有进程同步避免出现某些节点已初始化 autograd、另一些尚未初始化的竞态见init.py#L200-L210。超时store 的超时被设置为rpc_backend_options.rpc_timeoutRPC 默认超时时间为 60 秒常量定义见 torch/distributed/rpc/constants.py。框架是否可用可通过torch.distributed.rpc.is_available()检测其实现为检查torch._C._rpc_init是否存在见 torch/distributed/rpc/init.py#L24-L29。三大调用原语rpc_sync / rpc_async / remote这三个 API 允许用户在远端执行函数、或创建指向远端数据对象的引用RRef。选择原则非常直接如果用户代码没有返回值就无法继续用同步 API否则用异步 API 获取 Future待需要返回值时再等待。remote则适用于只在远端创建某物、但调用者永远不需要取回的场景——例如 driver 进程同时搭建参数服务器parameter server和 trainer它可以在参数服务器上创建 embedding 表并把该表的引用分享给 trainer而自身永远不会在本地使用这张表此时rpc_sync/rpc_async都不合适因为它们总是隐含返回值最终回到调用者。rpc_sync阻塞式远程调用rpc.rpc_sync(to, func, argsNone, kwargsNone, timeoutNone)to目标 worker可为字符串名称、rank 或WorkerInfo实例func可调用对象包括 Python 函数、内建算子如torch.add与带类型标注的 TorchScript 函数timeout本次调用的超时秒数0表示无限超时未提供时使用初始化时设定的默认值。实现定义见 torch/distributed/rpc/api.py#L772。该方法是线程安全的RPC 消息的收发与 Python 代码的执行并行进行。rpc_async非阻塞式远程调用rpc.rpc_async(to, func, argsNone, kwargsNone, timeoutNone)立即返回一个torch.futures.Future可后续wait()等待。实现见 torch/distributed/rpc/api.py#L846。文档同时给出警告由rpc_async返回的Future不应在shutdown()之后调用wait()。remote创建远程引用rpc.remote(to, func, argsNone, kwargsNone, timeoutNone)在目标 worker 上运行func并立即返回一个指向结果的RRef。目标 worker 成为该 RRef 的 owner发起remote调用的 worker 是 user用户。owner 维护其 RRef 的全局引用计数只有当全局不存在任何存活引用时owner 端的 RRef 才会被析构。实现见 torch/distributed/rpc/api.py#L555。张量传输的约束在这些 API 中当Tensor作为参数或返回值传递时目标 worker 会尝试创建具有相同 metashape、stride 等的Tensor。框架有意禁止直接传输 CUDA 张量——因为若源、目标 worker 的设备列表不一致可能导致崩溃。此时应用应显式在调用端将输入张量移到 CPU必要时再在 callee 端移到目标设备。辅助 API 与装饰器get_worker_info 与 shutdownget_worker_info(worker_nameNone)返回指定 worker 名称的WorkerInfo可避免每次调用都传递昂贵的字符串传None时返回当前 worker 的信息见 torch/distributed/rpc/api.py#L428。_to_worker_info辅助函数表明to参数同时兼容WorkerInfo、字符串与 intrank。shutdown(gracefulTrue, timeoutDEFAULT_SHUTDOWN_TIMEOUT)停止本地 agent 接收新请求并终止所有 RPC 线程。gracefulTrue时会阻塞等待所有本地与远端 RPC 进程都到达该方法、等待所有未完成任务完成gracefulFalse则是仅本地的快速关闭见 torch/distributed/rpc/api.py#L327。在 constants.py 中DEFAULT_SHUTDOWN_TIMEOUT 0。WorkerInfo描述单个 worker 的类包含name与id字段通常由用户为每个 worker 显式构造。async_execution 装饰器RPC 包还提供了装饰器允许应用指定某个函数在 callee 端应如何处理from torch.distributed.rpc import functions functions.async_execution def func(...): ...async_execution用于标注返回值保证是torch.futures.Future对象、且可以在 RPC callee 上异步执行的函数。callee 会提取被包装函数返回的Future并将后续处理步骤安装为对该Future的回调回调在Future完成时读取其值并通过 RPC 发送给调用方见 torch/distributed/rpc/functions.py#L5。这在实现批量 RPC 处理等需要聚合多个请求的场景中非常关键。后端注册机制backend_registry模块提供了两个后端相关函数见 torch/distributed/rpc/backend_registry.pybackend_registered(backend_name)检查某字符串名称是否已注册为 RPC 后端L51register_backend(backend_name, construct_rpc_backend_options_handler, init_backend_handler)注册新的 RPC 后端。注册通过动态扩展BackendType枚举完成L64-L99内置后端TENSORPIPE即通过此机制注册L433-L437。Backends可插拔的通信后端RPC 模块可通过不同后端完成节点间通信后端在init_rpc中通过BackendType枚举指定。无论使用哪个后端其余 RPC API 保持不变。每个后端还定义了自己的RpcBackendOptions子类其实例可传给init_rpc以配置后端行为。TensorPipe 后端默认TensorPipe agent 是默认后端基于 TensorPipe 库——一个专为机器学习设计的原生点对点通信原语从基础上解决了 Gloo 的部分局限。与 Gloo 相比其优势在于异步性允许大量传输同时进行、互不阻塞各自按自己的速度推进按需建连只在需要时于节点对之间打开 pipe某个节点故障时只有与其相连的 pipe 被关闭其余 pipe 正常工作多传输协议支持 TCP、共享内存、NVLink、InfiniBand 等并能自动探测可用性与协商每条 pipe 的最优传输方式高带宽内置基于 TCP 的传输能自动将大张量分块chunk并在多个 socket/线程上多路复用multiplex达到很高的带宽agent 会自动挑选最优传输无需人工干预。配置示例取自 rpc.mdimport os from torch.distributed import rpc os.environ[MASTER_ADDR] localhost os.environ[MASTER_PORT] 29500 rpc.init_rpc( worker1, rank0, world_size2, rpc_backend_optionsrpc.TensorPipeRpcBackendOptions( num_worker_threads8, rpc_timeout20, # 20 second timeout ) ) # worker2 上的 init_rpc 调用此处省略TensorPipeRpcBackendOptions 参数详解从 torch/distributed/rpc/options.py#L50-L109 的实现来看TensorPipeRpcBackendOptions继承自 C 侧的_TensorPipeRpcBackendOptionsBase支持以下参数参数类型默认值说明num_worker_threadsint16TensorPipeAgent 用于执行请求的线程池线程数。源码常量见 constants.pyrpc_timeoutfloat60 秒RPC 请求的默认超时超时未完成则抛出异常。单次调用可在rpc_sync/rpc_async中覆盖init_methodstrenv://用于 rendezvous 的分布式 store URL接受torch.distributed.init_process_group同参数的任何合法值device_mapsDict[str, Dict]None本 worker 到 callee 的设备放置映射键为 callee worker 名值为把本 worker 设备映射到 callee 设备的字典键值可为 int、str 或torch.devicedevicesListNone本 worker 上 RPC agent 使用的全部本地 CUDA 设备默认由其自身的device_maps与对端device_maps中对应设备推导。处理 CUDA RPC 请求时agent 会对列表中所有设备正确同步 CUDA stream此外TensorPipeRpcBackendOptions还提供两个增量配置方法options.py#L111-L179set_device_map(to, device_map)为每对 RPC 调用方/callee 设置设备映射可多次调用以增量添加配置映射必须是可逆的 1 对 1 映射否则抛出ValueError。set_devices(devices)设置 RPC agent 使用的本地设备列表。官方文档给出了set_device_map的完整示例worker0 通过device_maps{worker1: {0: 1}}把自身的cuda:0映射到 worker1 的cuda:1再通过options.set_device_map(worker1, {1: 2})把自身cuda:1映射到 worker1 的cuda:2发送参数与返回值时均按映射及其逆映射自动迁移设备。失败与重试策略RPC 框架不会自动重试任何rpc_sync、rpc_async与remote调用。原因在于框架无法判断某个操作是否幂等、重试是否安全。RPC 通信基于 TCP可能因网络故障或间歇性连接问题失败。因此应用需要自行处理失败并根据需要重试且应以合理的退避backoff策略避免激进重试压垮网络。RRef远程引用协议RRefRemote REFerence是指向远端 worker 上某个类型T如Tensor值的引用句柄。该句柄在 owner 上保持被引用值存活但不意味着值未来会传输到本地 worker。在多机训练中RRef 可用于持有存在于其他 worker 上的nn.Module引用并在训练期间调用相应函数检索或修改其参数。注意当前 RRef 不支持 CUDA 张量。RRef 协议的设计细节记录在 远程引用协议设计笔记 中核心要点如下概念模型每个 RRef 由remote调用创建callee worker 即 owner可被多个 user 使用。owner 存储真实数据并维护全局引用计数每个 RRef 由创建时分配的全局唯一RRefId标识。在 owner 上只存在一个包含真实数据的OwnerRRef实例而 user worker 上可以有任意多个不持有数据的UserRRef。UserRRef在作为rpc_sync/rpc_async/remote调用的参数或返回值时被创建owner 会收到通知以更新引用计数。当全局不存在UserRRef且 owner 本地也不再引用OwnerRRef时OwnerRRef及其数据才会被删除。设计假设瞬时网络故障协议通过消息重试处理瞬时网络故障但无法处理节点崩溃或永久网络分区——发生此类事故时应关闭所有 worker、回滚到最近 checkpoint 后恢复训练非幂等 UDF假定用户提供的函数不可重试但内部 RRef 控制消息是幂等的可重试乱序投递收发双方都使用多线程不假定任意节点对之间的消息有序。生命周期与两条保证协议的目标是在没有存活的UserRRef且用户代码也不持有OwnerRRef的恰当时机删除OwnerRRef。实现依赖两条关键保证G1任何UserRRef被删除时owner 都会被通知由UserRRef析构函数发送删除消息实现G2父 RRef 在子 RRef 被 owner 确认之前不会被删除父UserRRef在被 fork 时进入一个以新ForkId为键的上下文只有收到子节点的 ACK——子节点仅在得到 owner 确认后才会发送 ACK——才从上下文移除。由于 fork 图恒为树每次 fork 都会在 callee 上创建新的UserRRef实例即使 child 先于 parent 的消息到达 owner也至少有一个祖先存活由 G2 保证不会导致误删。设计笔记还以四张消息流程图分别分析了用户以返回值把 RRef 分享给 owner用户以参数分享给 ownerowner 分享给用户用户间分享四种场景核心结论是删除消息只在同时满足owner 已确认对应 ForkIdG2与Python GC 判定本地UserRRef可回收两个条件后才会发出。RemoteModule远程创建 nn.ModuleRemoteModule提供了一种在其他进程上远程创建nn.Module的简便方式真实模块驻留在远端 host本地 host 只持有该模块的句柄却能像调用普通nn.Module一样调用它——只是调用会转化为到远端的 RPC 请求且可通过额外 API 异步执行。同样地当前 RemoteModule 不支持 CUDA 张量。从 torch/distributed/nn/api/remote_module.py 的实现看其构造函数接收remote_device远端设备规格、module_cls模块类与可选的args/kwargs在指定远端节点创建模块若module_cls的forward签名为def forward(input: Tensor) - Tensor生成的RemoteModule将同时具有forward与异步版本forward_async。官方文档重点介绍的两个成员方法为remote_parameters(recurseTrue)返回指向远端模块参数的一组RRef可配合DistributedOptimizer使用remote_module.py#L277get_module_rref()返回指向远端模块的RRef[nn.Module]remote_module.py#L296。分布式自动微分框架该模块提供了基于 RPC 的分布式自动微分框架适用于模型并行训练等场景。简单说应用可以通过 RPC 发送和接收需要梯度记录的张量前向传播时框架记录需要梯度记录的张量何时通过 RPC 发送反向传播时则利用这些信息通过 RPC 执行分布式反向传播。当前分布式自动微分不支持 CUDA 张量。完整设计见 分布式自动微分设计笔记核心 API 为torch.distributed.autograd的context、backward、get_gradients与is_available。前向传播中的自动微分记录PyTorch 在前向传播中构建 autograd 图。对于分布式 autograd框架在执行 RPC 时向图中挂载send与recv函数send函数挂载在 RPC 源端其输出边指向 RPC 输入张量的 autograd 函数反向传播时其输入来自目标端对应recv函数的输出recv函数挂载在 RPC 目标端其输入取自目标端基于输入张量执行的算子其输出梯度在反向传播时被送回源端对应的send函数每对send-recv被赋予全局唯一的autograd_message_id用于反向传播时在远端节点查找对应函数对于 RRef每次调用RRef.to_here都会为涉及的张量挂载合适的send-recv对。下图展示了前述两节点示例省略t5.sum()对应的 autograd 图分布式 Autograd Context每次使用分布式自动微分的前向/反向过程都对应一个唯一的torch.distributed.autograd.context带有全局唯一的autograd_context_id按需在各节点创建。其作用有三多个分布式反向过程可能在同一张量上累积梯度为避免在运行优化器前.grad被各次反向混在一起每次反向的梯度独立累积在对应context中前向传播时把每次 autograd 过程的send/recv函数存入该 context既保持图中相关节点存活也便于反向传播时快速查找存放每次分布式 autograd 过程的元数据。用户视角的典型用法是上下文管理器形式import torch.distributed.autograd as dist_autograd with dist_autograd.context() as context_id: loss model.forward() dist_autograd.backward(context_id, loss)必须强调模型的 forward 必须在分布式 autograd 上下文管理器内执行因为只有存在合法 context所有send/recv函数才能被正确记录从而在所有参与节点上运行反向传播。分布式反向传播与依赖计算单机 autograd 引擎会先计算图中每个节点的依赖数dependencies以决定节点何时可执行。分布式场景下计算依赖困难得多某些 RPC 的结果可能并未参与 loss 计算对应send/recv在反向传播中并不需要执行。为此设计笔记提出两种算法FAST mode 算法当前唯一已实现核心假设是每个send函数在反向传播时依赖数均为 1即假定会从其他节点收到梯度。代价是应用需要知晓该限制——大多数应用不会执行未被使用的 RPC因此该假设通常成立。算法流程如下从持有反向传播根节点所有根必须本地的 worker 开始查找当前 Distributed Autograd Context 的所有send函数从给定根与所有send函数出发在本地计算依赖计算完成后以给定根启动本地 autograd 引擎引擎执行到recv函数时recv通过 RPC 把输入梯度发送给对应 worker——每个recv函数在前向阶段已记录目标 worker id同时携带autograd_context_id与autograd_message_id远端收到请求后用这两个 id 查找对应的send函数若该 worker 首次收到该autograd_context_id的请求则按 1-3 步在本地计算依赖把第 6 步找到的send函数入队到该 worker 的本地 autograd 引擎执行最后梯度不累积在张量的.grad字段上而是按 Distributed Autograd Context 独立累积在Dict[Tensor, Tensor]中可通过dist_autograd.get_gradientsAPI 取回。下图展示了分布式依赖计算面临的核心挑战send/recv跨节点形成依赖关系SMART mode 算法面向并非每个send/recv都参与反向传播的通用情形目前在文档中仅有设计思路可参考 PyTorch 相关 RFC尚未实现。完整代码示例设计笔记给出了分布式自动微分与分布式优化器的端到端示例以下代码置于dist_autograd_simple.py后可用MASTER_ADDRlocalhost MASTER_PORT29500 python dist_autograd_simple.py运行import torch import torch.multiprocessing as mp import torch.distributed.autograd as dist_autograd from torch.distributed import rpc from torch import optim from torch.distributed.optim import DistributedOptimizer def random_tensor(): return torch.rand((3, 3), requires_gradTrue) def _run_process(rank, dst_rank, world_size): name worker{}.format(rank) dst_name worker{}.format(dst_rank) # Initialize RPC. rpc.init_rpc( namename, rankrank, world_sizeworld_size ) # Use a distributed autograd context. with dist_autograd.context() as context_id: # Forward pass (create references on remote nodes). rref1 rpc.remote(dst_name, random_tensor) rref2 rpc.remote(dst_name, random_tensor) loss rref1.to_here() rref2.to_here() # Backward pass (run distributed autograd). dist_autograd.backward(context_id, [loss.sum()]) # Build DistributedOptimizer. dist_optim DistributedOptimizer( optim.SGD, [rref1, rref2], lr0.05, ) # Run the distributed optimizer step. dist_optim.step(context_id) def run_process(rank, world_size): dst_rank (rank 1) % world_size _run_process(rank, dst_rank, world_size) rpc.shutdown() if __name__ __main__: # Run world_size workers world_size 2 mp.spawn(run_process, args(world_size,), nprocsworld_size)分布式优化器torch.distributed.optim.DistributedOptimizer的工作方式见 分布式自动微分设计笔记 及 torch/distributed/optim/optimizer.py接收待优化的远程参数列表RRef也可以是包裹在本地RRef中的本地参数接收本地优化器类如torch.optim.SGD该优化器将运行在所有不同的 RRef owner 上分布式优化器在每个 worker 节点创建本地Optimizer实例并持有指向它们的RRef。从源码看构造函数会按param.owner()对参数 RRef 分组optimizer.py#L155-L165调用step()时分布式优化器通过 RPC 在所有相关远端 worker 上远程执行本地优化器step(context_id)必须传入分布式 autograd 的context_id本地优化器据此应用对应 context 中保存的梯度optimizer.py#L192若多个并发分布式优化器同时更新同一 worker 上的参数更新通过全局锁串行化对应_LocalOptimizer.global_lock见 optimizer.py#L53-L59。设计笔记与配套文档除本文外仓库还提供两篇深入的设计文档建议结合阅读分布式自动微分设计笔记详细讲解 RPC 版分布式自动微分的 send/recv 挂载、autograd context、FAST/SMART 依赖计算算法与端到端示例适用于模型并行训练等场景远程引用协议RRef设计笔记讲解 RRef 协议设计细节包括 G1/G2 保证、引用计数、消息流与四种分享场景适用于理解框架如何跨节点引用远程值。官方 RPC 教程系列还涵盖以下主题可从 PyTorch 官方教程站获取分布式 RPC 框架入门、用分布式 RPC 框架实现参数服务器、将 DistributedDataParallel 与 RPC 框架结合使用覆盖 RemoteModule、以及实现批量 RPC 处理batch RPC processing配合async_execution装饰器使用。在动手实践前请始终留意文首提到的维护模式状态与 CUDA beta 支持限制并遵循应用自行处理失败重试的约定以保证分布式训练的健壮性。【免费下载链接】pytorchTensors and Dynamic neural networks in Python with strong GPU acceleration项目地址: https://gitcode.com/GitHub_Trending/py/pytorch创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表