Skip to content

核心概念

这一章把 Ray Core 用到的关键概念逐个讲清楚。后面的架构、生命周期和源码章节都会反复用到这些术语。

Driver

Driver 是启动 Ray 程序的 Python 进程。你在 Driver 里调用 ray.init(),定义远程函数和 Actor,然后提交任务。

Driver 不是“控制一切的主节点”。它更像一个客户端:负责发起任务、持有 ObjectRef、调用 ray.get 等待结果。真正执行任务的是 Worker 进程,调度和对象管理由 Raylet 等运行时组件完成。

python
import ray

ray.init()

@ray.remote
def f(x):
    return x + 1

ref = f.remote(41)
print(ray.get(ref))

Task

Task 是一次远程函数调用。@ray.remote 把普通函数包装成远程函数,.remote() 不会直接运行函数,而是把调用提交给 Ray。

python
@ray.remote
def parse(path):
    return load_and_clean(path)

refs = [parse.remote(path) for path in paths]

Task 的典型特点:

  • 默认无状态,适合可并行拆分的计算。
  • 返回 ObjectRef,结果进入 Ray 管理的对象系统。
  • 可以声明资源,例如 @ray.remote(num_cpus=2, num_gpus=1)
  • 失败后可以按配置重试。

ObjectRef

ObjectRef 是远程对象的引用,不是对象本身。

python
ref = parse.remote("data.csv")
print(ref)          # ObjectRef(...)
result = ray.get(ref)

它有点像 Future,但比本地 Future 更重要:Ray 可以用它表达跨进程、跨节点的数据依赖。

python
@ray.remote
def train(batch):
    ...

batch_ref = load_batch.remote("s3://...")
model_ref = train.remote(batch_ref)

train 依赖 load_batch 的输出。你不需要自己写回调或消息传递,Ray 会根据 ObjectRef 建立依赖关系。

Object Store

Object Store 保存任务返回值和 ray.put 放进去的对象。大对象不会简单地塞进每次 RPC 参数里,而是由运行时管理位置和生命周期。

python
weights_ref = ray.put(model_weights)
refs = [score.remote(weights_ref, shard) for shard in shards]

这个写法的意义是:大对象只放一次,多个任务共享同一个引用。真实 Ray 使用共享内存对象存储、序列化协议、引用计数和对象溢写机制来降低数据移动成本。

Actor

Actor 是有状态的远程对象。你把类标记为 @ray.remote,然后创建远程实例。

python
@ray.remote
class Counter:
    def __init__(self):
        self.value = 0

    def incr(self):
        self.value += 1
        return self.value

counter = Counter.remote()
refs = [counter.incr.remote() for _ in range(3)]
print(ray.get(refs))

Actor 的关键是状态和顺序。默认情况下,同一个 Actor 的方法调用会按顺序执行,因此它适合:

  • 参数服务器。
  • 在线推理模型实例。
  • 长连接客户端。
  • 缓存、计数器、协调器。

Job、Node、Worker

概念含义
Job一次 Ray 应用运行,通常由一个 Driver 发起
NodeRay 集群中的一台机器或容器
Worker实际执行 Task 或 Actor 方法的 Python 进程
Raylet每个节点上的核心守护进程,负责本地资源、Worker 管理和对象传输
GCSGlobal Control Store,保存集群级元数据

资源

Ray 的资源是逻辑资源,不一定等同于操作系统实时用量。

python
@ray.remote(num_cpus=2, num_gpus=1)
def train_model(config):
    ...

这里的 num_cpus=2 表示调度时需要占用两个 CPU 资源配额。Ray 不会自动限制你的函数只能使用两个 OS 线程;你仍然需要设置底层库的线程数。

工程经验

资源声明主要服务于调度。它告诉 Ray“这个任务需要什么位置和容量”,不是强制的性能隔离机制。

Placement Group

Placement Group 用来把一组资源请求作为一个整体调度。比如分布式训练中,你希望 4 个 GPU worker 同时拿到资源,而不是只启动一半。

python
from ray.util.placement_group import placement_group

pg = placement_group([{"GPU": 1}, {"GPU": 1}, {"GPU": 1}, {"GPU": 1}])
ray.get(pg.ready())

后续任务或 Actor 可以绑定到这个 Placement Group 上,保证资源布局符合预期。

概念之间的关系

一句话总结

可以把 Ray Core 理解成一组配合工作的机制:remote 把本地调用变成运行时任务,ObjectRef 把结果变成可传递的依赖,Object Store 管理对象数据,Scheduler 决定在哪里执行,Worker 真正运行 Python 代码,Actor 把状态变成可调度的远程服务。

这些概念之间的关系,在后续章节会通过生命周期、架构图和源码逐步展开。

面向学习目的的 Ray Core 中文导读与 Mini Ray 机制预览。