资讯动态

Ray Direct Transport (RDT) API 完全指南:绕过对象存储的 GPU 张量直传

发布时间:2026/9/19 10:10:43 来源:尧图企业网站定制
Ray Direct Transport (RDT) API 完全指南绕过对象存储的 GPU 张量直传【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/rayRay 对象默认存放在 CPU 侧的对象存储中访问时需要复制与反序列化对 GPU 数据而言这会造成 GPU→CPU→GPU 的多次昂贵拷贝。Ray Direct Transport (RDT) 让 Ray 可以直接在 Actor 之间传递张量配合ray.method(tensor_transport...)、ray.put、ray.get以及ray.experimental.collective系列 API 与 NIXL 高级 API将数据保存在 GPU 显存中直到真正需要传输。读完本文你将掌握 RDT 三种张量传输后端Gloo / NCCL / NIXL的完整用法、collective group 的创建与销毁、可变对象的同步陷阱wait_tensor_freed以及自定义张量传输TensorTransportManager的扩展接口。本文以 API 参考文档 为主线骨架并结合 RDT 主体指南、示例代码 与源码实现python/ray/experimental/rdt/展开讲解。RDT 是什么为 GPU 张量消除不必要的拷贝在传统 Ray 数据流中一个 CUDAtorch.Tensor从 Actor A 传递到 Actor B需要先拷贝到 CPU 内存、序列化写入对象存储再反序列化并拷回 GPU 显存。RDT 对这一流程做了三点改变GPU 数据停留在显存中直到真正需要传输时才移动避免昂贵的序列化与对象存储读写使用高效的传输库直接进行设备间传输GlooCPU 上的 PyTorch 集合通信库、NVIDIA NCCLNVIDIA GPU 集合通信库以及基于 RDMA 的点对点传输库NIXL在 AWS EFA 实例上由 LIBFABRIC 驱动其他环境使用 UCX。RDT 是对既有ray.ObjectRefAPI 的增强张量仍然以ObjectRef的形式传递但底层传输路径从对象存储换成了Actor 直连。注意RDT 目前处于alpha阶段仅支持torch.Tensor与 Ray Actor 任务后续版本可能存在破坏性 API 变更详见下文 Limitations 一节。与 Core API 的配合使用RDT 与 Ray Core 的接口集成点有三个ray.method装饰器、ray.put和ray.get。用ray.method(tensor_transport...)开启 RDT对需要返回张量的 Actor 任务在ray.method装饰器中传入tensor_transport参数即可import torch import ray ray.remote class MyActor: ray.method(tensor_transportgloo) def random_tensor(self): return torch.randn(1000, 1000)关键点装饰器只需加在返回张量的 Actor 任务上消费张量的 Actor 任务除非它也返回张量不需要加返回张量时Ray 只保存对张量的引用reference而不是复制进 CPU 内存当该ObjectRef被传给另一个任务时Ray 会用指定传输如 Gloo直接把张量传到目标任务该装饰器同样适用于返回张量嵌套在其他 Python 对象如 dict中的任务。tensor_transport支持的值与传输后端对应glooCPU、ncclNVIDIA GPU、nixlCPU 或 NVIDIA GPU无需预建 collective group。用ray.put(..., _tensor_transport...)创建 RDT 引用ray.put也支持传入_tensor_transport参数直接由某个 Actor 持有张量并生成ObjectRef例如在 NIXL 场景下批量把 GPU 张量交给一个 Actor 托管def produce(self, tensors): refs [] for t in tensors: refs.append(ray.put(t, _tensor_transportnixl)) return refs从源码看该参数在 ray_option_utils.py 与 worker.py 中被解析属于 RDT 预留的传输选项。用ray.get取回结果ray.get默认会复用ray.method中指定的传输后端来取回结果Gloo / NCCL 这类集合通信two-sided传输如果调用方不在 collective group 中ray.get会直接抛错提示ray.get is not allowed on RDT objects using the two-sided transport。此时必须显式指定_use_object_storeTrue让 Ray 从对象存储取回结果# 错误用法调用方不在 collective group 中会抛 ValueError # ray.get(tensor) # 正确用法显式走对象存储 print(ray.get(tensor, _use_object_storeTrue)) # torch.Tensor(...)NIXLone-sided传输ray.get可以直接用 NIXL 取回张量副本无需预建 groupprint(ray.get(tensor)) # 返回 torch.TensorCollective tensor transports创建与管理集合通信组Gloo 与 NCCL 都属于 two-sided集合通信传输发送方和接收方必须都参与传输。因此使用这类传输前必须先创建 collective group。相关 API 位于ray.experimental.collective模块API作用ray.experimental.collective.create_collective_group在一组 Actor 上创建集合通信组ray.experimental.collective.get_collective_groups查询当前已创建的集合通信组ray.experimental.collective.destroy_collective_group销毁集合通信组销毁后可基于同样的 Actor 重新创建创建 collective groupimport torch import ray from ray.experimental.collective import create_collective_group ray.remote class MyActor: ray.method(tensor_transportgloo) def random_tensor(self): return torch.randn(1000, 1000) def sum(self, tensor: torch.Tensor): return torch.sum(tensor) sender, receiver MyActor.remote(), MyActor.remote() # backend 必须与 ray.method 中的 tensor_transport 匹配 group create_collective_group([sender, receiver], backendtorch_gloo)backend与tensor_transport的对应关系为tensor_transportgloo↔backendtorch_glootensor_transportnccl↔backendnccl。创建完成后两个 Actor 即可通过 Gloo/NCCL 直接通信用destroy_collective_group(group)可销毁组之后可基于同样的 Actor 重建。Gloo 完整示例仅 CPUimport torch import ray from ray.experimental.collective import create_collective_group ray.remote class MyActor: ray.method(tensor_transportgloo) def random_tensor(self): return torch.randn(1000, 1000) def sum(self, tensor: torch.Tensor): return torch.sum(tensor) sender, receiver MyActor.remote(), MyActor.remote() group create_collective_group([sender, receiver], backendtorch_gloo) # 张量保存在 sender 的进程中而不是 Ray 的对象存储里 tensor sender.random_tensor.remote() result receiver.sum.remote(tensor) print(ray.get(result))注意receiver.sum没有加ray.method(tensor_transport...)装饰器所以它返回torch.sum(tensor)时走的是默认的对象存储路径但它作为 RDT 引用tensor的接收方时Ray 会自动用 Gloo 接收。RDT 同样支持张量嵌套在 Python 容器中以及返回多个张量的场景ray.remote class MyActor: ray.method(tensor_transportgloo) def random_tensor_dict(self): return {tensor1: torch.randn(1000, 1000), tensor2: torch.randn(1000, 1000)} def sum(self, tensor_dict: dict): return torch.sum(tensor_dict[tensor1]) torch.sum(tensor_dict[tensor2]) sender, receiver MyActor.remote(), MyActor.remote() group create_collective_group([sender, receiver], backendtorch_gloo) tensor_dict sender.random_tensor_dict.remote() result receiver.sum.remote(tensor_dict) print(ray.get(result))把 RDT 对象传回生产它的 ActorObjectRef可以传回生产它的 Actor此时 Ray 不做任何拷贝直接返回同一份torch.Tensor的引用tensor sender.random_tensor.remote() sum1 sender.sum.remote(tensor) # 返回引用无拷贝 sum2 receiver.sum.remote(tensor) # 发生一次 Gloo 传输 assert torch.allclose(*ray.get([sum1, sum2]))NCCL 示例NVIDIA GPU切换传输后端只需少量改动。相比 Gloo 示例差异有三点tensor_transportnccl、collective group 的backendnccl、张量通过.cuda()创建在 GPU 上import torch import ray from ray.experimental.collective import create_collective_group ray.remote(num_gpus1) class MyActor: ray.method(tensor_transportnccl) def random_tensor(self): return torch.randn(1000, 1000).cuda() def sum(self, tensor: torch.Tensor): return torch.sum(tensor) sender, receiver MyActor.remote(), MyActor.remote() group create_collective_group([sender, receiver], backendnccl) tensor sender.random_tensor.remote() result receiver.sum.remote(tensor) ray.get(result)NIXL无需 collective group 的点对点直传安装与后端选择NIXL 通过pip install nixl安装即可追求极致性能时可运行 install_gdrcopy.sh 脚本安装 GDRCopy例如install_gdrcopy.sh ${GDRCOPY_OS_VERSION} 12.8 x64。Ray 会根据硬件自动选择 NIXL 的底层后端无需配置AWS EFA 实例Ray 检测 EFA 设备检查/sys/class/net/efa*网络设备以及/sys/class/infiniband下绑定内核efa驱动的 rdma-verbs 设备——后者使容器/Kubernetes 环境下也能识别 EFA并验证一次真实规模的 CUDA 内存注册成功后才选择LIBFABRIC后端若验证失败通常是 GPUDirect 配置问题如缺少nvidia-peermem或 dmabuf 支持。使用前需先运行 EFA 安装器同时安装 EFA 驱动与 libfabric。其他所有环境使用UCX后端。当 Ray 运行 UCX 后端时可通过 UCX 环境变量控制传输选择仅对 UCX 生效对 LIBFABRIC 无效# 示例 UCX 配置按实际环境调整 export UCX_TLSall # 或指定具体传输如 rc,ud,sm,^cuda_ipc export UCX_NET_DEVICESall # 或指定网络设备如 mlx5_0:1,mlx5_1:1LIBFABRIC 后端则使用 libfabric 环境变量例如FI_PROVIDERefa固定 provider、FI_LOG_LEVELDebug输出诊断日志。基本用法NIXL 可以在 CPU 与 NVIDIA GPU 之间传输数据且不需要预先创建 collective group——只要相关 Actor 的环境里安装了 NIXL 即可。示例见 direct_transport_nixl.pyimport torch import ray ray.remote(num_gpus1) class MyActor: ray.method(tensor_transportnixl) def random_tensor(self): return torch.randn(1000, 1000).cuda() def sum(self, tensor: torch.Tensor): return torch.sum(tensor) sender, receiver MyActor.remote(), MyActor.remote() tensor sender.random_tensor.remote() result receiver.sum.remote(tensor) ray.get(result)与 Gloo 示例相比代码差异只有两点tensor_transportnixl且无需创建 collective group。NIXL 下的ray.put/ray.getNIXL 是 one-sided 传输ray.get可直接用 NIXL 取回结果也可以对ray.put(..., _tensor_transportnixl)产生的引用执行ray.gettensor1 torch.randn(1000, 1000).cuda() tensor2 torch.randn(1000, 1000).cuda() refs sender.produce.remote([tensor1, tensor2]) # 内部调用 ray.put(t, _tensor_transportnixl) ref1 receiver.consume_with_nixl.remote(refs) # 内部 ray.get(ref) 也走 NIXL print(ray.get(ref1))对象可变性陷阱与wait_tensor_freed与对象存储中不可变拷贝的语义不同RDT 对象是可变的Ray 只持有张量引用不复制直到请求传输时才移动。若返回张量的 Actor 自己也保留了引用并在 Ray 仍持有该引用期间对张量原地修改则接收方可能看到部分或全部修改。Ray 会在 RDT 对象同时被传回生产 Actor 与另一个 Actor 时打印警告UserWarning: GPU ObjectRef(...) is being passed back to the actor that created it Actor(MyActor, ...). Note that GPU objects are mutable. ...修复这类问题的方法是调用ray.experimental.wait_tensor_freed(tensor)它会阻塞直到所有依赖该张量的任务执行完毕、所有指向该张量的ObjectRef离开作用域之后 Actor 才能安全地原地写回张量。Ray 通过跟踪哪些任务把对应ObjectRef作为参数来判定依赖关系。正确的写法完整代码见 direct_transport_gloo.py 的__gloo_wait_tensor_freed_start段ray.remote class MyActor: ray.method(tensor_transportgloo) def random_tensor(self): self.tensor torch.randn(1000, 1000) return self.tensor def increment_and_sum_stored_tensor(self): # 1. 原地修改前先等 Ray 释放对该张量的所有引用 ray.experimental.wait_tensor_freed(self.tensor) self.tensor 1 return torch.sum(self.tensor) def increment_and_sum(self, tensor: torch.Tensor): return torch.sum(tensor 1) sender, receiver MyActor.remote(), MyActor.remote() group create_collective_group([sender, receiver], backendtorch_gloo) tensor sender.random_tensor.remote() tensor1 sender.increment_and_sum_stored_tensor.remote() # 2. 不要在这里 ray.get(tensor1)wait_tensor_freed 会阻塞到所有引用释放 # 此时 ray.get 会造成死锁 tensor2 receiver.increment_and_sum.remote(tensor) # 3. 删除 driver 对 tensor 的引用以解除 wait_tensor_freed 的阻塞 del tensor assert torch.allclose(ray.get(tensor1), ray.get(tensor2))三个关键变更sender 在原地修改前调用wait_tensor_freeddriver 跳过ray.get避免死锁driver 显式del tensor释放自己的引用。该 API 在 rdt_manager.py 中实现并被ray.experimental导出。Advanced APIsNIXL 内存管理与传输目标控制除上述核心 API 外ray.experimental还导出了一组面向 NIXL 与高级调优的实验性 API定义与导出见 python/ray/experimental/rdt/init.pyAPI用途ray.experimental.register_nixl_memory手动注册张量内存供 NIXL RDMA 传输使用ray.experimental.deregister_nixl_memory注销之前注册的 NIXL 内存ray.experimental.register_nixl_memory_pool注册 NIXL 内存池复用注册内存以降低注册开销ray.experimental.set_nixl_cuda_stream为 NIXL 设置 CUDA 流控制传输与计算流的同步ray.experimental.set_target_for_ref为指定的ObjectRef设置目标 Actorray.experimental.set_target_device_for_ref为指定的ObjectRef设置目标设备ray.experimental.wait_tensor_freed等待 Ray 释放对张量的所有引用上文已详述ray.experimental.register_tensor_transport在运行时注册自定义张量传输后端ray.experimental.TensorTransportManager自定义张量传输的抽象基类其中内存注册 / 内存池 / CUDA 流相关 API 直接服务于 NIXL 底层传输对应实现见 nixl_tensor_transport.py 与 nixl_memory_pool.py主要用于对注册时机、内存复用与流同步有特殊要求的性能敏感场景set_target_for_ref/set_target_device_for_ref用于将引用定向到指定 Actor 或设备。这些 API 属于实验性接口使用前建议先在python/ray/tests/rdt/下的测试用例如 test_rdt_nixl.py、test_nixl_backend_selection.py中验证行为。自定义张量传输TensorTransportManager与register_tensor_transportRay 允许在运行时注册自定义张量传输。做法是实现抽象接口ray.experimental.TensorTransportManager再用ray.experimental.register_tensor_transport注册。该抽象类的完整定义位于 tensor_transport_manager.py需要实现的核心方法包括tensor_transport_backend()返回传输后端名字Ray 用它匹配ray.method(tensor_transport...)参数is_one_sided()声明传输是 one-sided接收方主动发起如 NIXL、CUDA-IPC还是 two-sided收发双方都参与如 NCCL、GlooRay 不会对 one-sided 传输调用send_multiple_tensorscan_abort_transport()传输出错时能否安全中断进行中的 send/recv若为TrueRay 会调用abort_transport清理否则 Ray 直接杀掉涉及 Actor 以防止死锁actor_has_tensor_transport(actor)检查某个 Actor 是否具备该传输能力extract_tensor_transport_metadata(obj_id, rdt_object)在源 Actor 上、任务产出结果张量后立即调用记录张量 shape/dtype/device 并做传输相关的内存注册返回TensorTransportMetadataget_communicator_metadata(src_actor, dst_actor, backend)在 owner/driver 进程编排传输前调用返回双方协调所需的信息如 collective rankone-sided 传输通常返回空CommunicatorMetadatarecv_multiple_tensors(...)/send_multiple_tensors(...)分别在目的 Actor 与源 Actor 上执行实际收发garbage_collect(obj_id, ...)Ray 的分布式引用计数判定对象离开作用域后在源 Actor 上释放传输资源如注销内存abort_transport(...)系统错误时在收发双方 Actor 上中止进行中的传输。配套的元数据类TensorTransportMetadata携带张量 shape/dtype 列表与设备信息和CommunicatorMetadata传输协调信息同样可由自定义传输扩展。自定义传输需要在driver 进程中、在创建使用它的 Actor 之前完成注册。完整的分步实现指南含基于共享内存传输 numpy 数组的可运行示例见 custom-tensor-transport.rst示例代码在 direct_transport_custom.py。自定义传输的已知限制包括Actor 重启后无法访问自定义传输在 Actor 创建之后注册的传输对该 Actor 不可用异步 Actorout-of-order或跨进程提交任务时无法保证任务执行时传输已注册。错误处理RDT 的错误处理策略因错误层级而异应用层错误用户代码抛出的异常不会销毁 collective group像普通 Ray 对象一样传播给依赖任务Gloo / NCCL 集合操作的系统级错误销毁 collective group 并杀掉涉及 Actor防止挂死NIXL 传输的系统级错误Ray 或 NIXL 以异常中止传输并在依赖任务或对 NIXL ref 的ray.get上抛出该异常。系统级错误包括第三方传输库内部错误如 NCCL 网络错误、Actor 或节点故障、张量设备与传输后端不匹配如 NCCL 传 CPU 张量、RDT 对象拉取超时可用环境变量RAY_rdt_fetch_fail_timeout_milliseconds覆盖默认超时以及任何意外系统缺陷。总结与限制回顾 RDT 的核心要点使用集合通信类传输Gloo / NCCL时必须先创建 collective groupNIXL 只需所有涉及 Actor 安装 NIXLRDT 对象可变Ray 只持有引用而非拷贝应用代码需用wait_tensor_freed同步原地修改其他方面 Actor 用法与普通 Ray Core 一致。LimitationsRDT 处于 alpha 阶段当前限制如下未来版本可能解决仅支持torch.Tensor对象仅支持 Ray Actor 任务不支持普通 Ray task仅支持 GLOO、NCCL、NIXL 三种传输仅支持 CPU 与 NVIDIA GPURDT 对象可变语义差异见上文对 RDT 引用执行await暂不支持。对集合通信类 / two-sided 传输Gloo 和 NCCL还有额外限制只有创建 collective group 的进程能提交返回并传递 RDT 对象的 Actor 任务其他进程拿到 actor handle 后可正常提交任务但不能使用 RDT 对象创建组所在的进程不能把 RDTObjectRef序列化传给其他 Ray 任务或 Actor只能作为直接参数传给同组 Actor 任务每个 Actor 每种传输同一时间只能属于一个 collective group不支持ray.put不支持 out-of-order Actor如异步 Actor 或max_concurrency 1的 Actor。NIXL 存在一个已知问题暂不支持在同一 Actor 上存储张量集合有重叠但不相等的多个 GPU 对象若要支持该模式需确保第一个ObjectRef离开作用域后再把相同张量存入第二个对象复现示例见 direct_transport_nixl.py 的__nixl_limitations_start段相关行为在 test_rdt_nixl.py 中有测试覆盖。进一步阅读RDT 主体指南doc/source/ray-core/direct-transport/direct-transport.rst自定义张量传输实现指南doc/source/ray-core/direct-transport/custom-tensor-transport.rst完整示例代码doc/source/ray-core/doc_code/ 下的direct_transport_gloo.py、direct_transport_nccl.py、direct_transport_nixl.py、direct_transport_custom.py底层实现python/ray/experimental/rdt/rdt_manager.py、collective_tensor_transport.py、nixl_tensor_transport.py、nixl_memory_pool.py、cuda_ipc_transport.py、tensor_transport_manager.py测试用例python/ray/tests/rdt/Gloo / NCCL / NIXL / 自定义传输 / 后端选择【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价