通信实战:包装库进程、Ray Collective 与内嵌 HTTP 服务器)
Ray Actor 带外Out-of-Band通信实战包装库进程、Ray Collective 与内嵌 HTTP 服务器【免费下载链接】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 的 Actor 之间默认通过方法调用通信、通过分布式对象存储共享数据但在包装第三方分布式库、追求高性能集合通信、以及向集群外部暴露服务等场景下绕过 Ray 控制面的带外通信Out-of-band Communication反而更高效。本文基于 out-of-band-communication.rst 展开结合仓库中ray.util.collective的实现与示例代码讲清楚三类带外通信的适用场景、配置方法与底层机制读完你可以在自己的 Ray 应用中独立实现与 Horovod/Spark 等库进程协同、CPU/GPU 上的 NCCL/GLOO 集合通信以及 actor 内嵌 HTTP 服务。为什么需要带外通信Ray 的常规 Actor 通信模型是方法调用 对象存储客户端通过actor.method.remote()提交任务数据以 ObjectRef 的形式放入分布式对象存储由 Ray 的调度器与引用计数系统统一管理。这套模型的好处是容错、引用管理、内存生命周期都由 Ray 负责无需用户操心。但存在两类场景让带外通信更有价值库进程协同很多分布式计算库如 Horovod、Spark早已拥有成熟、高性能的内部通信栈NCCL、MPI、内部 RPC 等它们把 Ray 当作语言级 Actor 调度器让 Ray 负责进程的创建与编排而实际的数据交换交给这些库自身的通信栈完成——这就是带外通信。性能与语义诉求集合通信allreduce、broadcast 等经过对象存储中转会产生额外序列化、拷贝与调度开销直接让参与进程通过 NCCL/GLOO 等在带外互连性能更高。此外向 Ray 集群之外的客户端暴露服务HTTP/gRPC 端点本质上也是带外通信。场景一包装第三方库进程许多库已经具备成熟的内部通信机制Ray 只承担语言集成式 Actor 调度器的角色库进程之间的实际通信大多走既有通信栈。文档中给出的典型例子Horovod-on-Ray使用 NCCL 或 MPI 进行集合通信Ray 只负责拉起并管理参与训练的进程RayDPRay Spark使用 Spark 自身的内部 RPC 和对象管理器Object Manager进行数据交换。仓库中也有对应的验证性测试例如 horovod_example.py 与 test_horovod.py以及 test_raydp.py可以看到 Ray 在这些组合中扮演调度与生命周期管理角色、而把通信让位给第三方栈的协作方式。这种模式的核心收益是不必把成熟库的通信协议翻译成 Ray 的方法调用避免性能损耗同时仍能享受 Ray 的弹性调度、故障恢复与资源管理能力。场景二Ray Collective 集合通信库Ray 自带的集合通信库ray.util.collective是文档重点推荐的带外通信方案用于在分布式 CPU 或 GPU 之间进行高效的集合通信与点对点通信。其完整说明见 ray-collective.rst实现集中在 collective.py。架构与后端从源码看collective.pyray.util.collective采用后端注册机制导入时会尝试注册两个后端——NCCL对应NCCLGroup位于nccl_collective_group.py和 GLOO对应TorchGLOOGroup基于torch.distributed.gloo。若对应依赖未安装注册会被跳过可通过nccl_available()/gloo_available()/is_backend_available()查询可用性。支持矩阵来自 ray-collective.rst后端GLOOtorch.distributed.glooNCCL设备CPUGPUCPUGPUsend / recv✔✘✘✔broadcast✔✘✘✔allreduce✔✘✘✔reduce✔✘✘✔allgather✔✘✘✔reduce_scatter✔✘✘✔barrier✔✘✘✔gather / scatter / all-to-all✘✘✘✘支持三种张量类型torch.Tensor、numpy.ndarray、cupy.ndarray。ReduceOp枚举定义于 types.py支持SUM、PRODUCT、MIN、MAX。安装方面ray.util.collective已随 Ray wheel 打包只需按后端补充第三方依赖——GLOO 需要torchNCCL 需要cupy-cudaxxxx 对应环境的 CUDA 版本。创建集合通信组命令式与声明式集合原语都作用在组collective group上组内的一组进程通常是 Ray 管理的 Actor 或 Task会一起进入集合调用。使用前必须先静态声明这些进程为组。文档与 ray-collective.rst 给出了两种初始化 API。命令式在每个进程内部初始化init_collective_group(world_size, rank, backend, group_name)参数含义见 collective.pyworld_size为组内总进程数rank为当前进程排名backend为 NCCL 或 GLOOgroup_name为组名同一进程可属于多个组组名是唯一标识gloo_timeout为 GLOO 操作超时毫秒。源码还会校验rank必须满足0 rank world_size。声明式在 driver 中声明create_collective_group(actors, world_size, ranks, backend, group_name)把一组 Actor 一次性声明为组。源码collective.py会校验每个 actor 对应一个 rank、ranks 必须是0..len(ranks)-1的排列、后端必须已注册且可用随后会创建一个名为info_group_name的 detached NamedActor 来保存组信息供各进程后续查询。一个完整的 NCCL allreduce 示例精简自 nccl_allreduce_example.pyimport cupy as cp import ray import ray.util.collective as collective ray.remote(num_gpus1) class Worker: def __init__(self): self.send cp.ones((4,), dtypecp.float32) self.recv cp.zeros((4,), dtypecp.float32) def setup(self, world_size, rank): collective.init_collective_group(world_size, rank, nccl, default) return True def compute(self): collective.allreduce(self.send, default) return self.send def destroy(self): collective.destroy_group() if __name__ __main__: ray.init(num_gpus2) num_workers 2 workers, init_rets [], [] for i in range(num_workers): w Worker.remote() workers.append(w) init_rets.append(w.setup.remote(num_workers, i)) _ ray.get(init_rets) results ray.get([w.compute.remote() for w in workers]) ray.shutdown()对同一组进程可以构建多个集合组以group_name区分从而在进程的不同子集之间描述更复杂的通信拓扑。集合通信原语的行为特征ray.util.collective当前提供的集合 API 是命令式的文档明确了三个行为特征同步阻塞所有集合 API 都是阻塞调用必须全员参与每个 API 只描述通信的一部分需要组内每个参与进程都发起调用并相互 rendezvous会合后通信才真正发生并继续带外执行API 必须在参与进程actor/task的代码内部使用。常用的集合原语在 collective.py 中的签名包括原语签名要点行号allreduce(tensor, group_name, op)全组归约默认ReduceOp.SUML312barrier(group_name)组内同步屏障L352reduce(tensor, dst_rank, group_name, op)归约到指定 rankL362broadcast(tensor, src_rank, group_name)从src_rank广播L421allgather(tensor_list, tensor, group_name)全组收集L468reducescatter(tensor_list, tensor, group_name, op)归约后再分发L511send(tensor, dst_rank, group_name)/recv(tensor, src_rank, group_name)点对点收发L567 / L624点对点通信与常见反模式P2P 的send/recv与集合函数行为一致同步阻塞成对出现。示例逻辑来自 ray-collective.rst# 假设 A、B 已在同一个组中rank 分别为 0、1 ray.get([A.do_send.remote(target_rank1), B.do_recv.remote(src_rank0)]) # 正确 ray.get([A.do_send.remote(target_rank1)]) # 反模式会挂起文档明确警告只调用 send 而不实例化对面的 recv 是反模式程序会一直挂起——因为 send/recv 必须成功 rendezvous 才能推进。多 GPU 集合通信当一台机器有多个 GPU 时可利用 NVLINK 等 GPU-GPU 带宽显著提升通信性能。ray.util.collective提供多 GPU 版本 API如allreduce_multigpu、send_multigpu、recv_multigpu允许一个进程如num_gpus4的 actor管理多块 GPU。相关限制见 ray-collective.rst仅支持 NCCL 后端参与多 GPU 调用的进程必须拥有相同数量的 GPU输入通常是张量列表每个张量位于调用方进程持有的一块不同 GPU 上。多 GPU 场景比单 GPU 原语 与 GPU 数量相同的进程通常更具性能优势。仓库在 examples 下提供了nccl_allreduce_multigpu_example.py、nccl_p2p_example_multigpu.py等可直接运行参考的示例。环境变量const.py 定义了NCCL_USE_MULTISTREAM环境变量默认开启用于控制 NCCL 后端是否使用多流multistream执行可作为调优入口。场景三在 Actor 内嵌 HTTP 服务器另一种常见的带外通信是在 actor 内部启动一个 HTTP 服务器并对外暴露端点让Ray 集群之外的用户可以直接与 actor 通信而不必经过 Ray 的方法调用层。文档给出的完整示例位于 actor-http-server.pyimport ray import asyncio import requests from aiohttp import web ray.remote class Counter: async def __init__(self): self.counter 0 asyncio.get_running_loop().create_task(self.run_http_server()) async def run_http_server(self): app web.Application() app.add_routes([web.get(/, self.get)]) runner web.AppRunner(app) await runner.setup() site web.TCPSite(runner, 127.0.0.1, 25001) await site.start() async def get(self, request): return web.Response(textstr(self.counter)) async def increment(self): self.counter self.counter 1 ray.init() counter Counter.remote() [ray.get(counter.increment.remote()) for i in range(5)] r requests.get(http://127.0.0.1:25001/) assert r.text 5这个示例的精妙之处在于把两条通信路径打通并协同带内路径driver 通过counter.increment.remote()调用 actor 方法 5 次递增内部计数器带外路径actor 的__init__中通过asyncio.get_running_loop().create_task(...)启动一个 aiohttp 服务器监听127.0.0.1:25001外部客户端示例中是requests直接发起 HTTP 请求读取当前计数值最终断言返回5验证两条路径操作的是同一份状态。要点说明actor 构造函数是异步的async def __init__在初始化时就把 HTTP 服务任务注册进事件循环从而在 actor 存活期间持续服务采用asyncio.create_task而非在构造函数内await阻塞保证__init__能及时返回、actor 正常完成创建监听地址与端口可按需调整示例固定为127.0.0.1:25001若需要被集群外客户端访问应绑定到可达的 IP 并处理好防火墙/网络策略。同样的思路也适用于暴露其他协议的服务器例如在 actor 内启动gRPC 服务器文档明确提及对外提供 gRPC 端点。这意味着你可以把 actor 当作一个常驻的分布式微服务进程用 Ray 管理其生命周期与资源同时用其自带的服务器协议与外部系统互操作。限制与注意事项文档在最后明确提醒带外通信的关键限制Ray 不管理 actor 之间的带外调用。方法调用/对象存储路径上 Ray 提供的分布式引用计数distributed reference counting、内存生命周期管理等保障对带外通信不生效不要把对象引用ObjectRef通过带外通道传递。由于带外通道绕过了 Ray 的引用计数对象引用的创建与释放无法被正确追踪可能造成对象被提前回收或泄漏引发难以排查的错误。因此实践中的红线是带外通道只传普通数据张量、序列化值、业务消息等绝不传 ObjectRef需要 Ray 参与管理的对象引用一律走常规的remote()调用与ray.get()/ray.put()路径。小结带外通信是 Ray Actor 编程模型的重要补充仓库文档总结的三条路径各有定位包装库进程Horovod-on-Ray、RayDP 等把通信交给成熟第三方栈NCCL/MPI、Spark RPCRay 专注调度与生命周期Ray Collectiveray.util.collective在 Ray 内部提供 NCCL/GLOO 双后端的集合与点对点原语覆盖 allreduce、broadcast、send/recv、多 GPU 等场景是高性能分布式训练的首选内嵌服务器在 actor 内启动 HTTP/gRPC 服务把 actor 变成集群外部可直接访问的服务端点适用于对外提供接口或接入外部系统的场景。使用时要始终牢记带外通信绕过 Ray 的管理面Ray 的分布式引用计数不适用切勿在带外通道中传递对象引用。把握住这一边界带外通信就能成为你手中连接Ray 调度能力与高性能通信/对外服务的桥梁。【免费下载链接】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),仅供参考