news 2026/9/21 3:26:59

Ray Client 架构指南:深入解析 Ray 分布式运行时的 gRPC 客户端/服务器设计与实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Ray Client 架构指南:深入解析 Ray 分布式运行时的 gRPC 客户端/服务器设计与实现

Ray Client 架构指南:深入解析 Ray 分布式运行时的 gRPC 客户端/服务器设计与实现

【免费下载链接】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

导读

本文以仓库内 python/ray/util/client/ARCHITECTURE.md 为骨架,系统讲解 Ray Client 的整体架构:它本质上是一对 gRPC 客户端与服务器,服务器端运行ray.init()扮演普通的 Ray Driver,并通过 gRPC 连接被远程客户端控制。读完本文,你将掌握 Ray Client 的代码布局、gRPC 协议(含数据通道与日志通道)的设计动机、基于 pickle/cloudpickle 的自定义序列化协议如何弥合客户端桩对象与服务端真实对象,以及 Ray Client 如何通过 client mode hook 与 Ray core 集成,并了解其两套测试模式。

总体架构:一个受 gRPC 控制的"远程 Driver"

核心模型

Ray Client 是一对 gRPC 客户端与服务器:

  • 服务器:在远程节点上运行ray.init(),行为与普通的 Ray Driver 完全一致,只是它受 gRPC 连接控制。服务器负责所有簿记工作(bookkeeping),并为其接入的客户端保持对象的作用域(scope)。
  • 客户端:用户代码运行的一端,通过 gRPC 把请求发给服务器。

代码布局上,客户端代码位于ray/util/client,服务器代码位于ray/util/client/server。查看服务器启动入口 python/ray/util/client/server/server.py 中的serve()可以看到,服务器会同时注册三个 gRPC Servicer:

task_servicer = RayletServicer(ray_connect_handler) data_servicer = DataServicer(task_servicer) logs_servicer = LogstreamServicer() ray_client_pb2_grpc.add_RayletDriverServicer_to_server(task_servicer, server) ray_client_pb2_grpc.add_RayletDataStreamerServicer_to_server(data_servicer, server) ray_client_pb2_grpc.add_RayletLogStreamerServicer_to_server(logs_servicer, server)

依赖方向的刻意设计

仓库约定:ray/util/client避免直接 importray;而服务器端本质上是另一个 Ray 应用,因此允许直接依赖ray。这一分离有两个目的:

  1. 避免依赖环:客户端不反向依赖 Ray core,避免循环导入;
  2. 未来可拆分:如果将来需要,可以将任一端单独抽成独立仓库或独立安装包(例如pip install ray_client)。

从源码看,客户端代码确实大量通过from ray.util.client import ray(即RayAPIStub)而不是直接import ray工作,而服务器端 server_pickler.py 则直接import rayimport ray.cloudpickle,印证了这一约定。

RayAPIStub:以对象代替模块的 API 表面

模块级全局变量ray(类型为RayAPIStub)在 python/ray/util/client/init.py 中定义,它充当与ray包等价的 API 表面——ray命名空间中的函数都成了RayAPIStub对象上的方法。RayAPIStub是类而非模块,这一"微妙差异"是后续 client mode hook 能够工作的前提(见下文"集成点"章节),因为类支持__getattr__,可以动态地把未覆盖的 API 转发到客户端实现。

客户端对象与服务端对象的对应关系

ray命名空间中的许多对象,在客户端都有对应物,大部分集中在 python/ray/util/client/common.py。例如:

ObjectRef <-> ClientObjectRef ActorID <-> ClientActorRef RemoteFunc <-> ClientRemoteFunc

这条对应关系有实际的调试价值:如果你在 bug 报告中看到的对象类型是ClientObjectRef,那么它必然来自客户端代码路径(由服务器返回后构造)。在 common.py 中可以看到ClientObjectRef继承自raylet.ObjectRef,内部持有Future延迟绑定 id;ClientActorRef继承自raylet.ActorID。两者的析构函数都会在对象失活时向服务器发送 release 请求,这是客户端侧引用计数的体现(见"数据通道"小节)。ClientRemoteFuncClientActorClassClientActorHandleClientRemoteMethod则分别对应远程函数、Actor 类、Actor 句柄与 Actor 方法,@ray.remote在客户端模式下由remote_decorator把它们包装为ClientStub子类(common.py)。

协议层:gRPC 服务 + pickle 数据编码

Ray Client 存在两套协议,需要区分清楚:

  1. gRPC 服务:API 表面,用于实现远程、瘦客户端的各项功能;
  2. 数据编码协议:函数与数据的编码方式。由于两端都是 Python,这一层采用pickle(实际是cloudpickle)——客户端对象(包括客户端桩对象)通过 pickle 被透明地序列化并传输到服务器。

简言之:gRPC 服务定义了 API 形状,pickle 协议定义了 API 之上的数据编码方式

gRPC 服务定义

唯一的一份 gRPC 规范位于 src/ray/protobuf/ray_client.proto。proto 文件注释详尽,每个字段都说明了用途。它定义了三个 service:

Service用途
RayletDriver一元(unary)RPC 与传统请求/响应,如InitGetObjectPutObjectWaitObjectScheduleTerminateClusterInfo、KV 操作等
RayletDataStreamer双向流式数据通道Datapath,承载全部请求/响应模式
RayletLogStreamer双向流式日志通道Logstream
一元 RPC:get、put 与函数调用

Client 最初以一组一元 RPC 起步,功能刚好覆盖最常用的 API。作为 RPC API 的入门,它们很适合用来理解协议工作方式:

  • Get / Put / Wait:标准的对象存取与等待语义;
  • Schedule:最有意思的一个——它隐含了一个前置的 Put(把要执行的函数 put 上去),然后再执行它。

这些一元 RPC 至今仍保留,但已处于可弃用(ripe for deprecation)状态,核心问题在于它们不与持久连接绑定:

设想负载均衡器后面挂着多台 client-server,任何客户端都可能通过一元 RPC 命中任意一台服务器的任意状态。由于我们需要持有 ray ObjectRef 等句柄、保证它们不因超出作用域而被丢弃,这些句柄必须与已连接的客户端保持同步。一元 RPC 下,x = f.remote()可能打到服务器 A,而后续的ray.get(x)却打到服务器 B——而 B 上根本没有对应的 ObjectRef。

对应的消息结构(ClientTask)定义了RemoteExecType枚举:FUNCTIONACTORMETHODSTATIC_METHODNAMED_ACTOR,并携带payload_id(对已 put 函数的引用)、args/kwargs(已废弃,改由序列化的data字段承载)、client_idnamespace、任务选项optionsbaseline_options,以及大对象分块的chunk_id/total_chunks。服务器端 server.py 的Schedule会根据task.type分派到_schedule_function_schedule_actor_schedule_method_schedule_named_actor

数据通道:双向流 + ClientID + 引用计数

数据通道是客户端的双向流式连接Datapath),它封装了与一元 RPC 相同的全部请求/响应模式(见 ray_client.proto 中的DataRequest/DataResponse,其oneof type覆盖 get、put、release、init、task、terminate、acknowledge 等全部消息)。

其设计要点:

  1. ClientID 关联:连接建立之初,客户端用 UUID 生成ClientID与服务器关联。只要通道保持打开,客户端就算"在线"。通过跟踪 ClientID,服务器可以追踪它为某个特定客户端持有的全部资源,并且一旦通道断开(channel drops),即可判定客户端已断连。
  2. 引用计数:客户端跟踪自己对各种 Ray client 对象的引用数;当对象超出作用域时,客户端可主动发送ReleaseRequest给服务器做乐观清理(optimistic cleanup)。其余情况,引用计数全部在客户端完成,服务器只需知道何时清理。服务器也可以在它认为安全时,清理某客户端持有的全部引用。ReleaseRequestids是要释放的引用集合(ray_client.proto)。
  3. 未来扩展:文档指出,将来可以增加显式的ClientDisconnection消息,以区分"客户端主动完成、永远不会回来"与"客户端正遭遇连接问题"两种场景。

在客户端实现中,python/ray/util/client/dataclient.py 的DataClient维护一个专用线程ray_client_streaming_rpc运行双向流,用req_id(int32 递增计数器,溢出回绕)把异步请求与响应配对;outstanding_requests记录未完成请求以便断线后重放,ready_data存放阻塞式响应,asyncio_waiting_data存放异步回调。服务器端 python/ray/util/client/server/dataservicer.py 的DataServicer同样处理这些请求,并用OrderedResponseCache配合客户端的AcknowledgeRequest(每收到 32 个响应发送一次 ACK,见 dataclient.py)做去重与缓存清理。

大对象分块传输

为了支持跨网络传输大对象,Ray Client 将对象切成5 MiB 的块OBJECT_TRANSFER_CHUNK_SIZE,见 common.py),put 与 task 的请求体在发送端被惰性分块(dataclient.py 中的chunk_put/chunk_task),get 的响应在服务器端分块、客户端用ChunkCollector按序重组并支持从start_chunk_id断点续传。gRPC 单条消息上限被放宽到 2GiB(GRPC_MAX_MESSAGE_SIZE),超过 2GiB 的对象会触发用户警告,建议改用 S3 等远程 URI 传输。此外,common.py 还设置了 30 秒的 keepalive ping 与 600 秒的超时,以适配 ELB 等负载均衡器的 60 秒空闲超时。

日志通道:独立于数据通道的双向流

与数据通道类似,还有一个关联的日志通道,把日志回传给客户端:

  • 它是独立的通道,因为日志是"附属品"——即使因断连丢失一些日志,通常也可以接受;
  • 独立通道还让日志聚合器不必实现完整 API 即可接入;
  • 它是双向流:客户端发送LogSettingsRequest控制日志的开关(enabled)与级别(loglevel),服务器把产生的日志以LogDatamsg+level+name)流式回传(ray_client.proto)。协议注释约定:level > 0遵循 Python logging 的级别;level == -1表示 stdout;level == -2表示 stderr。

关于 CloudPickle:自定义 Pickler 解决桩对象混用问题

如本节开头所述,pickle/cloudpickle 是"把数据编码为可被 Python 执行的数据"以走 gRPC 通用传输的方式。Ray Client 在 python/ray/util/client/client_pickler.py 与 python/ray/util/client/server/server_pickler.py 中提供了自己的 pickle/unpickle 子类,目的是解决Client*桩对象混用的问题。

用一个例子说明问题的由来:假设RemoteFunc f()调用了另一个RemoteFunc g()f()需要持有对g()的引用(比如调用g.remote()),于是f()在序列化时,序列化数据里会包含一个RemoteFunc类的对象,供 worker 端反序列化。在 Ray Client 中,f()g()都是ClientRemoteFunc,行为同理。在早期版本里,一个ClientRemoteFunc必须知道自己是在服务器端还是客户端,并据此决定"像普通 RemoteFunc 一样运行"还是"在客户端发起调用"。这导致了一些棘手的 bug——尤其是把客户端桩对象传出去或(更糟)返回回来时(设想f()调用g(),而g()构造并返回了一个新的闭包h())。

自定义 pickler 的解决方式:只要客户端桩对象被序列化,就用一个结构体(写作时是一个元组PickleStub)替代它;反序列化时,服务器"填入"对应的非桩对象。反之亦然——如果服务器要编码一个返回/响应中的ObjectRef,就把元组放到线上,客户端反序列化器再把它还原成ClientObjectRef

PickleStub是一个命名元组,字段为typeclient_idref_idnamebaseline_options(client_pickler.py)。客户端ClientPickler.persistent_idRayAPIStubClientObjectRefClientActorHandleClientRemoteFuncClientActorClassClientRemoteMethod分别生成对应 stub(如"Ray""Object""Actor""RemoteFunc""RemoteActor""RemoteMethod");而服务器端ClientUnpickler.persistent_load按 stub 类型回填:"Object"self.server.object_refs[pid.client_id][pid.ref_id]"Actor"self.server.actor_refs[pid.ref_id]"RemoteFunc"/"RemoteActor"lookup_or_register_func/lookup_or_register_actor"RemoteMethod"→ 从 actor 句柄取方法(server_pickler.py)。反向路径上,服务器ServerPickler.persistent_id遇到ray.ObjectRef/ray.actor.ActorHandle时把它们登记到对应客户端的对象/actor 表中并生成 stub,ServerUnpickler.persistent_load(客户端侧)再把"Object"/"Actor"stub 还原为ClientObjectRef/ClientActorHandle(client_pickler.py)。

这一设计带来的收益:

  • 客户端侧尽可能只与桩对象打交道,服务器侧永远看不到桩对象,两端界限干净;
  • 服务器侧出现了ClientObjectRef属于错误情形而非需要特殊处理的场景;
  • 服务器侧可以像普通 Ray 一样工作、只处理普通 Ray 对象,并在发送时把它们透明地编码成客户端对象;
  • 因为"永不相交"(never the twain shall meet),建模与调试都容易得多。

另外,client_pickler.py 还处理了递归/自引用:函数(或 Actor 类)在 put 自己过程中被编码时,_refInProgressSentinel,会生成"RemoteFuncSelfReference"/"RemoteActorSelfReference"stub,服务器端对应构造ClientReferenceFunction/ClientReferenceActor(server_pickler.py),从而支持"函数的参数里包含它自己"这类递归闭包场景。

与 Ray core 的集成点:client mode hook

为了提供与 Ray core 无缝的客户端体验,需要包装一部分核心 Ray 函数(例如ray.get())。Python 的动态特性在此帮了大忙:如前所述,RayAPIStub是类而非模块,因此可以在 API 层实现__getattr__把调用重定向到任何想去的地方(模块做不到这一点,除非等 Python 3.6 废弃后使用 PEP 562 的__module_getattr__)。如果raycore 本身是对象而非命名空间里的函数,就不需要包装、直接替换实现即可,但项目必须保持向后兼容。

所有与 Ray core 的集成点都位于 python/ray/_private/client_mode_hook.py:

  • client_mode_hook装饰器:用来包装 Ray core 函数。当client_mode_should_convert()根据环境变量返回True时,装饰器生效,把调用转发给ray/util/client对象(getattr(ray, func.__name__)(*args, **kwargs));否则照常调用原函数(client_mode_hook.py)。
  • 模式开关is_client_mode_enabled默认关闭,在ray.client(...).connect()或测试中开启;RAY_CLIENT_MODE=1环境变量可让整个进程默认开启 client mode(主要用于测试,见 client_mode_hook.py)。
  • disable_client_hook()上下文管理器:线程本地地把 hook 状态置为 False,服务器端所有真正执行 Ray 操作的地方都用它包住(例如 server.py 中InitSchedulePutObjectWaitObject等实现都写在with disable_client_hook():内),保证服务器端调用的是"真实的" Ray 函数而不是客户端实现,避免死循环。
  • client_mode_convert_function/client_mode_convert_actor:用于把"预先注册的 RemoteFunction / ActorClass"透明转换为ClientRemoteFunc/ClientActorClass,典型场景是函数在库加载早期、尚未进入 client mode 时就被@ray.remote装饰(client_mode_hook.py)。
  • client_mode_wrap:用于实现不属于主ray.*API、却需要服务器端执行的功能,例如 Placement Group 的创建——客户端模式下会把函数包装成ray.remote(num_cpus=0)任务去服务器端执行(client_mode_hook.py)。

对应地,python/ray/util/client/init.py 的_ClientContext.connect()会调用_explicitly_enable_client_mode()强制开启 client mode,并在断连时由disconnect()恢复。

连接生命周期与服务器部署

虽然 ARCHITECTURE.md 以代码布局为重点,但结合 python/ray/util/client/init.py 与 python/ray/util/client/server/server.py,可以完整还原一条连接的生命周期:

  1. 启动服务器python -m ray.util.client.server --host <host> --port <port> --mode <proxy|legacy|specific-server> [--address <ray集群地址>]。默认 mode 为proxy(即 Ray Client Proxy,把请求转发给 Ray 集群),legacy/specific-server则直接serve()起一个内嵌 ray.init 的服务器(server.py)。init_and_serve则是进程内启动并自动连接的便捷入口,默认监听本地 50051 端口(init.py)。
  2. 客户端连接ray.util.client.ray.connect("host:port", ...)创建Worker,通过数据通道发送InitRequest(内含 pickle 化的job_configray_init_kwargs);服务器端RayletServicer.Init反序列化后调用ray_connect_handler执行ray.init(),并校验客户端/服务器版本(check_version_info,可被ignore_versionRAY_IGNORE_VERSION_MISMATCH覆盖)。
  3. 使用与清理:客户端通过Schedule提交任务、GetObject/PutObject存取对象;连接关闭或通道断开时,服务器通过 ClientID 追踪到的引用表(object_refsactor_refsactor_owners)被release_all清空(server.py)。
  4. 断线重连DataClient在 RPC 错误后可恢复的情况下,先 ping 服务器判断通道是否仍可用,必要时重建 gRPC channel,并把outstanding_requests中尚未确认的请求全部重放(dataclient.py);服务器端的ResponseCache/OrderedResponseCache则保证重放的请求不会在服务器上被执行两次(幂等语义,见 common.py)。

测试策略:两种互补的模式

测试 client 代码有两种主要方式:

方式一:以"客户端/服务器应用"的视角测试

明确知道自己正在测试一个客户端/服务器应用:同时运行连接的两端,调用(已知是客户端侧的)客户端,观察服务器端的效果,反之亦然。形如test_client*.py的测试集合采用该方式,仓库中包括:

  • python/ray/tests/test_client.py:核心客户端行为;
  • python/ray/tests/test_client_reconnect.py:断线重连;
  • python/ray/tests/test_client_references.py:引用计数与释放;
  • python/ray/tests/test_client_proxy.py:Proxy 模式;
  • python/ray/tests/test_client_builder.py:ray.client(...)builder API;
  • python/ray/tests/client_test_utils.py:公共测试工具(fixture)。

一般原则:如果实现的是"让 client 工作"的特性,应该放进test_client系列——因为你同时控制两端、可以直接测客户端代码。

方式二:把 Ray 的既有测试作为 fixture 跑在 client 模式上

另一种方式是把 Ray 自己的测试当作 fixture,像普通 Ray 一样测试 Ray 的 API。这是一个非常强大的模式:

  • 某测试在此模式下通过,意味着用户有信心"一切工作方式与以往一致",无论是否使用 client;
  • 实现难度更高:因为接入/关闭一个客户端/服务器对与单节点 ray 实例的 setup/shutdown 不同,尤其是这些测试可能对 fixture 的搭建方式做了假设;
  • 它是唯一能测试集成点的方式(即验证 client mode hook 对 Ray core 函数的透明转发)。

因此:如果修的是用户侧 API bug 或与 Ray core 的集成问题,通常的做法是把一个既有单元测试改造/纳入 Ray Client 的测试集

小结:从架构到调试

回顾全文,Ray Client 的架构可以归纳为三个层次:

  1. API 层RayAPIStub以对象形式提供与ray包等价的 API 表面,客户端桩对象(ClientObjectRefClientActorRefClientRemoteFunc等)在 common.py 中定义;
  2. 传输层:gRPC 的三个 service(ray_client.proto)分别承载一元 RPC、双向数据通道与日志通道;ClientID 关联连接、ReleaseRequest 实现引用计数、分块与缓存机制支撑大对象与重连;
  3. 编码层:client/server 两侧的自定义 pickler(client_pickler.py 与 server/server_pickler.py)用PickleStub元组在"客户端桩对象"与"服务端真实对象"之间做无损翻译。

对开发者而言,理解这套分层最大的收益在于定位问题:看到ClientObjectRef就知道它来自客户端路径;看到PickleStub就知道正处于序列化转换边界;而disable_client_hook()包裹的代码段则代表服务器端真实执行的 Ray 调用。结合 client_mode_hook.py 的转发逻辑与test_client*系列测试,你可以快速判断一个 bug 是出在 API 转发、gRPC 传输还是序列化编码层。

【免费下载链接】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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/21 3:22:23

低功耗Bandgap设计实战:结构、启动电路与验证方法

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/21 3:19:20

SAP MM工厂间调拨:301与303移动类型选型指南与实战避坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/21 3:14:35

Modbus RTU现场通信故障排查:地电位、终端电阻与字节序的坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/21 3:13:23

public-apis实战指南:从API选型到生产级治理

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/21 3:03:33

反激变压器设计全流程:12V/1A宽压输入算例详解

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华