Skip to content

Task 生命周期

Ray Task 的用户 API 很简单,但一次 .remote() 背后会经过提交、依赖解析、调度、执行、对象写入和结果获取。

从函数到远程任务

普通函数:

python
def square(x):
    return x * x

Ray 远程函数:

python
@ray.remote
def square(x):
    return x * x

ref = square.remote(3)

两个关键变化:

  1. square 不再是普通函数,而是 Ray 包装后的远程函数句柄(RemoteFunction)。
  2. square.remote(3) 不会立即执行函数返回 9,而是提交任务给运行时,立即返回一个 ObjectRef

任务生命周期总览

.remote() 为什么立即返回

Ray 的并行性来自提交和等待的分离。

python
refs = [square.remote(i) for i in range(100)]
results = ray.get(refs)

如果 .remote() 阻塞到函数执行结束,100 个任务会逐个执行。Ray 让 .remote() 立即返回 ObjectRef,所以 Driver 可以快速构建任务图。

参数依赖

当一个任务的参数里包含 ObjectRef 时,Ray 会把它识别成依赖。

python
@ray.remote
def load():
    return [1, 2, 3]

@ray.remote
def total(values):
    return sum(values)

values_ref = load.remote()
total_ref = total.remote(values_ref)

total 不应该在 load 完成前执行。Ray 会等待 values_ref 对应对象可用,再把真实数据交给 total

这就是 Ray 能表达 DAG 的原因。你写的是普通函数调用形状,运行时看到的是对象依赖图。

Task 的状态流转

真实 Ray 的内部状态比下面更细,但学习时可以先记住这条主线:

每个状态背后的问题不同:

状态运行时正在解决的问题
Submitted任务已经被 Driver 交给运行时
WaitingForDeps参数对象还没准备好
Schedulable依赖 ready,可以找资源
Queued资源暂时不够或等待 Worker
RunningWorker 正在执行
Finished返回对象已经 ready
Failed / Retrying失败记录、异常传播和重试

返回值与异常

Task 成功时,返回值进入对象存储:

python
ref = square.remote(4)
assert ray.get(ref) == 16

Task 失败时,异常也会通过 ObjectRef 传播:

python
@ray.remote
def fail():
    raise ValueError("bad input")

ref = fail.remote()
ray.get(ref)  # 抛出 RayTaskError,包含远程异常信息

这个设计非常重要。远程异常不会在 .remote() 时抛出,因为任务那时可能还没运行。异常会在你取结果时暴露。

任务粒度

Ray 能处理大量任务,但每个任务都有提交、调度、序列化和对象管理的开销。任务太大浪费并行性,任务太小让调度开销占主导。

经验规则:

  • 单个任务的执行时间应该明显大于调度开销(通常至少几十毫秒)。
  • 不要把微秒级函数拆成独立的 Ray Task。
  • 小任务可以合并成 batch 处理。
  • 大对象用 ray.put 放入对象存储,多个任务共享引用,不要作为普通参数反复复制。

Mini Ray 中的实现路径

Mini Ray 的任务路径在 examples/mini_ray_runtime/mini_ray/core.py

python
@ray.remote
def add(a, b):
    return a + b

ref = add.remote(1, 2)

对应调用链:

  1. remote(add) 返回 RemoteFunction
  2. RemoteFunction.remote() 调用 Runtime.submit_task()
  3. ObjectStore.reserve() 先创建未完成对象。
  4. ThreadPoolExecutor 异步运行任务。
  5. 任务线程解析 ObjectRef 参数。
  6. ResourceScheduler.acquire() 获取逻辑资源。
  7. 函数执行后 ObjectStore.set_result() 唤醒等待者。
  8. ray.get(ref) 从对象存储取值。

Mini Ray 把多进程和网络通信替换成线程,但保留了“提交任务先返回引用,结果稍后 ready”的关键语义。这也是理解 Ray 对象系统的最小闭环。

最常见的错误

在循环里提交一个任务就立刻 ray.get,等于串行执行:

python
# 慢:提交一个等一个,跟串行没区别
results = [ray.get(work.remote(i)) for i in range(100)]

# 快:先全部提交,再统一等待
refs = [work.remote(i) for i in range(100)]
results = ray.get(refs)

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