Ray 是一个面向 Python 的分布式计算框架,最初作为 UC Berkeley RISELab 的研究项目发展而来,如今已演进为覆盖数据加载、分布式训练、超参调优、在线推理等环节的完整 AI 计算生态。本文从核心抽象、总体架构、底层机制与上层生态四个层面梳理 Ray 的设计与实现。
在 Ray 出现之前,机器学习流水线通常需要拼接多套系统:用 Spark 处理数据,用 Horovod 进行分布式训练,用 Celery 执行异步任务,用 Kubernetes 部署在线推理服务。这种组合带来几个问题:
Ray 提供了一套统一的编程模型:开发者只需在普通 Python 函数或类上添加 @ray.remote 装饰器,即可将单机代码扩展到多节点集群。它同时支持无状态的任务(Task)和有状态的服务(Actor),用同一套底座承载不同类型的 AI 计算负载。
Ray 的三个核心抽象是 Task(任务)、Actor(参与者)和 Object(对象)。
Task 是 Ray 调度的基本单位。任何 Python 函数加上 @ray.remote 装饰器后即成为远程函数。调用该函数(通过 .remote())时,Ray 会在集群的某个节点上异步执行它,并立即返回一个 ObjectRef(对象引用),而非阻塞等待结果。
Task 是无状态的。若需要在多次操作之间保持状态(例如维护模型权重或数据库连接),则使用 Actor。Python 类加上 @ray.remote 后,实例化时会在集群中创建一个常驻进程,即 Actor。对 Actor 方法的调用同样是异步的,且在同一 Actor 内按调用顺序串行执行。
Task 的返回值以及通过 ray.put() 显式放入集群的数据都属于 Ray 的 Object。Object 不可变(Immutable),存储在 Ray 的分布式内存对象存储中。Ray 返回一个 ObjectRef,用于跨节点安全地引用该数据,而无需主动拷贝。
下面的关系图展示了这三个概念如何协同工作的:
graph TD
UserCode[客户端/驱动程序] -->|提交无状态计算| Task1(Ray Task)
UserCode -->|提交无状态计算| Task2(Ray Task)
UserCode -->|创建状态副本| Actor[Ray Actor]
Task1 -->|输出| Obj1[(Object 1)]
Task2 -->|输出| Obj2[(Object 2)]
Actor -->|维护内部状态| State{State}
Actor -->|处理请求并输出| Obj3[(Object 3)]
Obj1 -.->|作为输入引用 ObjectRef| Task2
Obj2 -.->|作为输入引用 ObjectRef| Actor
style UserCode fill:#f9f,stroke:#333,stroke-width:2px
style Task1 fill:#bbf,stroke:#333
style Task2 fill:#bbf,stroke:#333
style Actor fill:#bfb,stroke:#333上图刻画了 Ray 编程模型中三类核心抽象之间的完整数据流,可以从以下三个层次来理解:
图中的粉色节点客户端/驱动程序是整个计算图的数据源,它通过两种方式向集群提交工作:一是将无状态的计算逻辑作为 Task 提交(对应 Task1、Task2 两条实线箭头,标注为”提交无状态计算”);二是通过 @ray.remote 装饰器创建带内部状态的 Actor(对应”创建状态副本”实线箭头)。值得注意的是,所有提交动作都是异步的——Task.remote() 与 Actor 方法调用都会立即返回,不会阻塞客户端。
每个 Task 执行完毕后,其返回值会被写入 Ray 的分布式内存对象存储,形成不可变的 Object(图中圆柱体 Obj1、Obj2、Obj3)。客户端拿到的只是指向这些对象的 ObjectRef 句柄,而非数据本身。Task1 的输出成为 Obj1,Task2 的输出成为 Obj2,这两个对象随后通过虚线箭头(”作为输入引用 ObjectRef”)分别被 Task2 和 Actor 消费——这正是 Ray 数据依赖的传递方式:Task 之间无需直接通信,只需引用彼此输出的 ObjectRef 即可建立依赖关系,Ray 调度器会根据这些依赖自动决定任务的执行顺序与数据放置位置。
绿色节点 Actor 的行为与 Task 有本质区别:它在集群中作为常驻进程运行,通过 State(菱形节点”维护内部状态”)保存跨调用持久化的成员变量——例如模型权重、数据库连接等;同时它对外”处理请求并输出” Obj3。由于 Actor 方法在同一实例内按调用顺序串行执行,天然提供了有状态服务所需的并发一致性。
实线箭头表示”提交/创建/输出”这类触发型关系(客户端→Task/Actor,Task/Actor→Object);虚线箭头表示”作为输入引用”这类依赖型关系(ObjectRef→Task/Actor);菱形 State 节点则代表 Actor 特有的内部状态保持关系。理解了这幅图,就把握住了 Ray 的核心设计理念:无状态计算用 Task、有状态服务用 Actor、跨节点传数据一律走 ObjectRef。
Ray 的架构采用去中心化与中心化相结合的设计。整体由三个主要部分构成:Global Control Store (GCS)、Raylet 和 Client/Worker。
graph TB
subgraph HeadNode["Head Node (头节点)"]
GCS[(Global Control Store)]
Raylet_Head[Raylet]
Driver[Driver Process]
end
subgraph WorkerNode1["Worker Node 1"]
Raylet1[Raylet]
Plasma1[(Plasma Object Store)]
Worker1_A[Worker Process]
Worker1_B[Worker Process]
Actor1[Actor Process]
Raylet1 <--> Plasma1
Worker1_A <--> Raylet1
Worker1_B <--> Raylet1
Actor1 <--> Raylet1
end
subgraph WorkerNode2["Worker Node 2"]
Raylet2[Raylet]
Plasma2[(Plasma Object Store)]
Worker2_A[Worker Process]
Raylet2 <--> Plasma2
Worker2_A <--> Raylet2
end
GCS <--> Raylet_Head
GCS <--> Raylet1
GCS <--> Raylet2
Plasma1 <..->|对象的 P2P 传输| Plasma2
style HeadNode fill:#EDE7F6,stroke:#7E57C2,stroke-width:2px,color:#4A148C
style WorkerNode1 fill:#E3F2FD,stroke:#42A5F5,stroke-width:2px,color:#0D47A1
style WorkerNode2 fill:#E0F2F1,stroke:#26A69A,stroke-width:2px,color:#004D40
style GCS fill:#7E57C2,stroke:#4527A0,color:#FFFFFF,stroke-width:2px
style Raylet_Head fill:#B39DDB,stroke:#5E35B1,color:#311B92
style Raylet1 fill:#B39DDB,stroke:#5E35B1,color:#311B92
style Raylet2 fill:#B39DDB,stroke:#5E35B1,color:#311B92
style Driver fill:#FFF176,stroke:#F9A825,color:#F57F17
style Worker1_A fill:#90CAF9,stroke:#1976D2,color:#0D47A1
style Worker1_B fill:#90CAF9,stroke:#1976D2,color:#0D47A1
style Worker2_A fill:#90CAF9,stroke:#1976D2,color:#0D47A1
style Plasma1 fill:#A5D6A7,stroke:#43A047,color:#1B5E20
style Plasma2 fill:#A5D6A7,stroke:#43A047,color:#1B5E20
style Actor1 fill:#FFAB91,stroke:#E64A19,color:#BF360C上图展示了 Ray 集群的基本拓扑。整个集群由一个 Head Node 和多个 Worker Node 组成,图中以两个 Worker Node 为例。三类节点、三类连线,下面分别说明。
先看节点。Head Node 承担管理职责,内部运行三个进程:
Worker Node 负责实际的计算,每个节点运行以下进程:
Worker1_A、Worker1_B)。再看连线,图中包含三类通信关系:
整体来看,GCS 与 Raylet 构成星形的控制面,各节点内部的 Raylet、Plasma、Worker/Actor 构成自洽的执行面,节点之间的对象传输走 P2P 链路,三者各司其职、互不干扰。
GCS 是 Ray 集群的中央元数据服务。早期版本基于 Redis 实现,后续版本改为自研的独立服务。它存放的是集群级信息,包括:
注意,单个任务(Task)的运行状态并不集中存在 GCS 里,而是由创建它的 worker 通过所有权模型自行管理(见下文第四节的“对象所有权模型”)——GCS 只负责集群层面这些跨节点共享的元数据。
GCS 简化了容错实现:节点失效后,新节点通过查询 GCS 即可重建集群状态。
每个计算节点上运行一个 Raylet 守护进程,负责节点内的资源管理与任务调度,包含两个子组件:
Ray 通过共享内存实现对象存储。同一节点上的不同 Worker 读取同一份数据时,往往可以避免数据在进程间复制——尤其是 numpy 这类基于缓冲区的对象,做的是零拷贝(Zero-copy)访问;一般 Python 对象仍需反序列化。这一机制对海量数据与模型预加载场景尤为重要。
说明:早期版本中对象存储由独立的 Plasma 进程提供;Ray 2.x 起 Plasma 被移除,对象存储改为内嵌于 Raylet 进程内实现,但共享内存与零拷贝的语义保持不变。上图中的 Plasma 对应的是这一对象存储组件。
在相对简单的 API 背后,Ray 在底层做了若干针对性优化。以下分析几个重点机制。
与 Hadoop/Spark 这类由中心调度器统一分配任务的“自上而下”模式不同,Ray 采用分布式调度,并配合自下而上的 Spillback 机制,以避免中心调度器在高并发下成为瓶颈:
该策略降低了任务调度的延迟:单机场景下本地 Task 的提交开销很小(亚毫秒级),吞吐量也高于中心化调度。
早期 Ray 将对象元数据集中维护在 GCS 中,带来较大的通信开销。后续版本引入所有权模型:调用创建任务 API 的 Worker 即“拥有”该任务返回值的 ObjectRef。
Owner 进程负责:
ObjectRef 在 Python 层的引用计数归零后,Owner 通过内部 RPC 通知所在节点的 Raylet,从对象存储中删除该 Object。当对象总量超过内存容量时,Ray 的对象存储支持 Object Spilling(对象溢出):按 LRU 策略将最冷的数据异步序列化并持久化到本地磁盘或云存储(如 AWS S3),对上层应用透明。
Ray 的容错分为数据容错与计算容错:
Ray 的价值不仅在于底层的 Core,还在于其官方构建的覆盖机器学习各环节的组件。
mindmap
root((Ray Ecosystem))
Data Processing
[Ray Data]
分布式数据读取
流式Map/Filter处理
Model Training
[Ray Train]
桥接PyTorch DDP
桥接DeepSpeed/Horovod
容错训练
Tuning
[Ray Tune]
超参数搜索
ASHA/PBT等高级统筹算法
Serving
[Ray Serve]
模型多路复用
请求批处理 (Batching)
复杂的推理DAG图 (Ensemble)
RL
[RLlib]
强化学习/PPO算法
分布式RL训练Ray 自带资源管理能力,但企业级数据中心普遍以 Kubernetes 作为资源底座。当前主流的部署方式是 KubeRay。
KubeRay 是一个 Kubernetes Operator,通过 RayCluster 这个 CRD(Custom Resource Definition)管理 Ray 集群:
RayCluster CR 里的 replicas 字段;随后 KubeRay operator 按新的 replicas 创建对应的 Ray worker Pod。如果 K8s 集群本身的节点(虚拟机)也不够,则由用户自行配置的 Kubernetes Cluster Autoscaler 去云厂商购买新的底层节点。这样从 Task 排队一直到云底座 IaaS,形成了一条完整的弹性伸缩链路。Ray 通过分布式调度、共享内存对象存储与统一的 Task/Actor 抽象,将机器学习流水线中的多套系统收敛到单一框架。在 LLM 时代,Ray 被广泛用于大规模分布式训练、RLHF 与多模态模型的并行推理部署。
理解 Ray 的底层机制,有助于编写高性能的并行 Python 代码,也有助于构建高可用、高吞吐的 AI 基础设施。随着算力协同需求的增长,Ray 已成为 AI 云原生基础设施中的常见组成部分。