Ray 是一个面向 Python 的分布式计算框架,最初作为 UC Berkeley RISELab 的研究项目发展而来,如今已演进为覆盖数据加载、分布式训练、超参调优、在线推理等环节的完整 AI 计算生态。本文从核心抽象、总体架构、底层机制与上层生态四个层面梳理 Ray 的设计与实现。

一、 引言:Ray 解决的问题

在 Ray 出现之前,机器学习流水线通常需要拼接多套系统:用 Spark 处理数据,用 Horovod 进行分布式训练,用 Celery 执行异步任务,用 Kubernetes 部署在线推理服务。这种组合带来几个问题:

  1. 多套系统的学习与运维成本叠加。
  2. 系统之间的数据交换依赖磁盘或网络序列化,存在额外开销。
  3. Python 本身缺少细粒度的分布式原语。

Ray 提供了一套统一的编程模型:开发者只需在普通 Python 函数或类上添加 @ray.remote 装饰器,即可将单机代码扩展到多节点集群。它同时支持无状态的任务(Task)和有状态的服务(Actor),用同一套底座承载不同类型的 AI 计算负载。


二、 Ray 的核心概念与编程模型

Ray 的三个核心抽象是 Task(任务)、Actor(参与者)和 Object(对象)

1. Task(无状态任务)

Task 是 Ray 调度的基本单位。任何 Python 函数加上 @ray.remote 装饰器后即成为远程函数。调用该函数(通过 .remote())时,Ray 会在集群的某个节点上异步执行它,并立即返回一个 ObjectRef(对象引用),而非阻塞等待结果。

2. Actor(有状态参与者)

Task 是无状态的。若需要在多次操作之间保持状态(例如维护模型权重或数据库连接),则使用 Actor。Python 类加上 @ray.remote 后,实例化时会在集群中创建一个常驻进程,即 Actor。对 Actor 方法的调用同样是异步的,且在同一 Actor 内按调用顺序串行执行。

3. Object 与 ObjectRef(分布式对象)

Task 的返回值以及通过 ray.put() 显式放入集群的数据都属于 Ray 的 Object。Object 不可变(Immutable),存储在 Ray 的分布式内存对象存储中。Ray 返回一个 ObjectRef,用于跨节点安全地引用该数据,而无需主动拷贝。

概念关系连线图

下面的关系图展示了这三个概念如何协同工作的:

上图刻画了 Ray 编程模型中三类核心抽象之间的完整数据流,可以从以下三个层次来理解:

图中的粉色节点客户端/驱动程序是整个计算图的数据源,它通过两种方式向集群提交工作:一是将无状态的计算逻辑作为 Task 提交(对应 Task1Task2 两条实线箭头,标注为”提交无状态计算”);二是通过 @ray.remote 装饰器创建带内部状态的 Actor(对应”创建状态副本”实线箭头)。值得注意的是,所有提交动作都是异步的——Task.remote() 与 Actor 方法调用都会立即返回,不会阻塞客户端。

每个 Task 执行完毕后,其返回值会被写入 Ray 的分布式内存对象存储,形成不可变的 Object(图中圆柱体 Obj1Obj2Obj3)。客户端拿到的只是指向这些对象的 ObjectRef 句柄,而非数据本身。Task1 的输出成为 Obj1Task2 的输出成为 Obj2,这两个对象随后通过虚线箭头(”作为输入引用 ObjectRef”)分别被 Task2Actor 消费——这正是 Ray 数据依赖的传递方式:Task 之间无需直接通信,只需引用彼此输出的 ObjectRef 即可建立依赖关系,Ray 调度器会根据这些依赖自动决定任务的执行顺序与数据放置位置。

绿色节点 Actor 的行为与 Task 有本质区别:它在集群中作为常驻进程运行,通过 State(菱形节点”维护内部状态”)保存跨调用持久化的成员变量——例如模型权重、数据库连接等;同时它对外”处理请求并输出” Obj3。由于 Actor 方法在同一实例内按调用顺序串行执行,天然提供了有状态服务所需的并发一致性。

实线箭头表示”提交/创建/输出”这类触发型关系(客户端→Task/Actor,Task/Actor→Object);虚线箭头表示”作为输入引用”这类依赖型关系(ObjectRef→Task/Actor);菱形 State 节点则代表 Actor 特有的内部状态保持关系。理解了这幅图,就把握住了 Ray 的核心设计理念:无状态计算用 Task、有状态服务用 Actor、跨节点传数据一律走 ObjectRef


三、 Ray 的总体架构设计

Ray 的架构采用去中心化与中心化相结合的设计。整体由三个主要部分构成:Global Control Store (GCS)、Raylet 和 Client/Worker。

1. 架构示意图

上图展示了 Ray 集群的基本拓扑。整个集群由一个 Head Node 和多个 Worker Node 组成,图中以两个 Worker Node 为例。三类节点、三类连线,下面分别说明。

先看节点。Head Node 承担管理职责,内部运行三个进程:

  • GCS:中央元数据服务,保存集群级状态(节点健康、Actor 位置、Placement Group 等)。
  • Raylet:头节点自身的 Raylet,负责本节点的资源管理与任务调度。
  • Driver:客户端驱动程序,是用户代码提交 Task、创建 Actor 的入口。

Worker Node 负责实际的计算,每个节点运行以下进程:

  • Raylet:节点内的调度中枢,同时是与 GCS 通信的桥梁。
  • Plasma:节点内的内存对象存储,存放本节点 Worker 产生的 Object。同一节点上的 Worker 通过共享内存读写该存储,无需跨进程拷贝。
  • Worker Process:执行无状态 Task 的工作进程,一个节点可以有多个(如 Worker Node 1 中的 Worker1_AWorker1_B)。
  • Actor Process:承载有状态服务的常驻进程。并非每个节点都必须有 Actor——图中 Worker Node 1 有一个 Actor,Worker Node 2 只有 Worker,没有 Actor。

再看连线,图中包含三类通信关系:

  • GCS 与所有 Raylet 之间的双向连线(GCS ↔ Raylet_Head / Raylet1 / Raylet2):这是控制面通信。每个 Raylet 启动时向 GCS 注册自身,之后持续上报心跳与资源使用情况;各节点也通过这条链路查询 GCS 保存的 Actor 位置注册表等集群级信息。
  • 节点内部 Raylet 与 Plasma、Worker、Actor 之间的双向连线:这是数据面与执行面的通信。Raylet 把任务分发给本节点的 Worker 并接收结果;Worker/Actor 读写 Object 时经由 Raylet 协调对应的 Plasma 对象存储。因此 Raylet 是节点内唯一的调度与协调入口。
  • Plasma1 与 Plasma2 之间的虚线双向连线:标注为”对象的 P2P 传输”。当一个节点的 Task 需要读取存储在另一个节点上的 Object 时,由本节点的 Object Manager 直接从远端 Plasma 拉取数据,不经过 GCS 中转。这也是跨节点数据流的主要通道。

整体来看,GCS 与 Raylet 构成星形的控制面,各节点内部的 Raylet、Plasma、Worker/Actor 构成自洽的执行面,节点之间的对象传输走 P2P 链路,三者各司其职、互不干扰。

2. 核心组件解析

2.1 Global Control Store (GCS)

GCS 是 Ray 集群的中央元数据服务。早期版本基于 Redis 实现,后续版本改为自研的独立服务。它存放的是集群级信息,包括:

  • 各节点的健康状态与资源可用量。
  • Actor 的位置注册表。
  • Job 与 Placement Group 等元数据。

注意,单个任务(Task)的运行状态并不集中存在 GCS 里,而是由创建它的 worker 通过所有权模型自行管理(见下文第四节的“对象所有权模型”)——GCS 只负责集群层面这些跨节点共享的元数据。

GCS 简化了容错实现:节点失效后,新节点通过查询 GCS 即可重建集群状态。

2.2 Raylet(节点管理器)

每个计算节点上运行一个 Raylet 守护进程,负责节点内的资源管理与任务调度,包含两个子组件:

  • Node Manager(节点管理器):接收来自本地 Worker 或其他 Raylet 的任务,根据本地资源情况分配给本地 Worker 进程执行;本地资源不足时,将任务转发(Spillback)给其他节点。
  • Object Manager(对象管理器):管理当前节点的对象存储,并在需要时从其他节点拉取对象数据(P2P 传输)。

2.3 内存对象存储

Ray 通过共享内存实现对象存储。同一节点上的不同 Worker 读取同一份数据时,往往可以避免数据在进程间复制——尤其是 numpy 这类基于缓冲区的对象,做的是零拷贝(Zero-copy)访问;一般 Python 对象仍需反序列化。这一机制对海量数据与模型预加载场景尤为重要。

说明:早期版本中对象存储由独立的 Plasma 进程提供;Ray 2.x 起 Plasma 被移除,对象存储改为内嵌于 Raylet 进程内实现,但共享内存与零拷贝的语义保持不变。上图中的 Plasma 对应的是这一对象存储组件。


四、 关键技术细节剖析

在相对简单的 API 背后,Ray 在底层做了若干针对性优化。以下分析几个重点机制。

1. 分布式的自下而上调度策略(Bottom-up Scheduling)

与 Hadoop/Spark 这类由中心调度器统一分配任务的“自上而下”模式不同,Ray 采用分布式调度,并配合自下而上的 Spillback 机制,以避免中心调度器在高并发下成为瓶颈:

  1. Driver 提交 Task 时,先提交给本地节点的 Raylet(Node Manager)。
  2. 本地 Raylet 优先尝试在本地执行该 Task(检查 CPU/GPU 资源与依赖的 Object 是否就绪)。
  3. 若满足条件,直接调起本地 Worker 执行,无需跨节点网络通信。
  4. 若本地无法执行(资源耗尽,或依赖数据位于远端且体积较大),本地 Raylet 将 Task 溢出(Spillback)给有剩余资源的远端 Raylet。

该策略降低了任务调度的延迟:单机场景下本地 Task 的提交开销很小(亚毫秒级),吞吐量也高于中心化调度。

2. 对象所有权模型(Ownership Model)与分布式 GC

早期 Ray 将对象元数据集中维护在 GCS 中,带来较大的通信开销。后续版本引入所有权模型:调用创建任务 API 的 Worker 即“拥有”该任务返回值的 ObjectRef

Owner 进程负责:

  • 记录 Object 的位置、状态与生命周期。
  • 在 Task 失败时决定重试逻辑(而非由全局中心决定)。
  • 管理垃圾回收(GC):当作用域内的 ObjectRef 在 Python 层的引用计数归零后,Owner 通过内部 RPC 通知所在节点的 Raylet,从对象存储中删除该 Object。

当对象总量超过内存容量时,Ray 的对象存储支持 Object Spilling(对象溢出):按 LRU 策略将最冷的数据异步序列化并持久化到本地磁盘或云存储(如 AWS S3),对上层应用透明。

3. Actor 的生命周期与容错恢复(Fault Tolerance)

Ray 的容错分为数据容错与计算容错:

  • Lineage Reconstruction(血统重建):对于无状态 Task,若节点宕机导致 Object 丢失,只要 Owner 仍在,Owner 知道该 Object 的生成方式,可重新提交 Task 重新计算。
  • Actor 重启机制:当包含 Actor 的节点宕机,GCS 通过心跳检测到失效后,会在其他存活节点重新调度并启动该 Actor 进程。但这只重建了 Actor 本身,不保证其运行时的内部状态(如成员变量),业务级容错需要开发者结合 Checkpoint 机制自行实现。

五、 Ray 的上层生态系统

Ray 的价值不仅在于底层的 Core,还在于其官方构建的覆盖机器学习各环节的组件。

  1. Ray Data:负责数据加载与预处理。它不取代 Spark 处理复杂 SQL 查询,而是将清洗后的特征数据或大模型语料以 Pipeline 形式高效喂入 GPU 进行训练。
  2. Ray Train:对 PyTorch 或 TensorFlow 的分布式训练进行封装,屏蔽节点变更导致的 RANK 更新等细节,并集成 Ray 的容错机制。
  3. Ray Serve:面向复杂模型流的在线推理框架。实际推理可能涉及多个模型(如文生图、文生文、审核模型)组成的 DAG。Ray Serve 可为 DAG 上不同节点分配不同资源(如文本节点用 CPU、图像节点用 GPU),实现资源差异化分配。
  4. RLlib:面向强化学习的分布式训练库,内置 PPO 等常用算法,用于大规模 RL 训练场景。

六、 生产环境实践:KubeRay 的结合

Ray 自带资源管理能力,但企业级数据中心普遍以 Kubernetes 作为资源底座。当前主流的部署方式是 KubeRay

KubeRay 是一个 Kubernetes Operator,通过 RayCluster 这个 CRD(Custom Resource Definition)管理 Ray 集群:

  • 开发者提交 YAML 声明集群规模,例如”1 个 Head、10 个带 GPU 的 Worker”。
  • KubeRay 据此拉起对应的 Pods。
  • 自动扩缩容(Autoscaling):这里的扩缩容分几个层级。当 Ray 调度器检测到大量 Task 因 GPU 不足而排队时,Ray Autoscaler(运行在 head Pod 内的 sidecar 容器)会请求新增 worker Pod——具体做法是递增 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 云原生基础设施中的常见组成部分。