跳转到内容

Ray 2018 — 把任务和演员放进同一个分布式舞台

待复核

Ray 是一个面向新一代 AI 应用的分布式执行框架。日常类比:以前办一场大型演出,要分别找舞台调度、演员经纪、道具仓库和后台通信;Ray 想做的是同一个后台系统,既能派一次性小工,也能管理长期待命的演员。

这篇论文的核心句子是:把无状态任务有状态 actor放进同一个动态任务图里执行。任务适合很多短小、可重试的计算;actor 适合保留模型、模拟器、参数服务器这类内部状态。

论文最初瞄准的是强化学习,因为强化学习不是单纯训练一个模型,而是训练、仿真、服务三件事不断互相喂数据。Ray 后来变成现代 Python 分布式计算、LLM 训练和推理编排常见的底座之一。

不理解 Ray 这篇,很多系统设计就会断一截:

  • 为什么 Spark 这类批处理系统很强,却不适合每秒百万级小任务和动态仿真循环。
  • 为什么强化学习需要同时关心 CPU 仿真、GPU 训练、低延迟策略服务,而不是只管训练吞吐。
  • 为什么现代 AI infra 常说 task、actor、object store、scheduler、lineage,这套词在 Ray 里被捏成一套工程语言。
  • 为什么 LLM 时代的分布式推理、微调、数据处理常需要一个”能动态编排 Python 函数”的运行时。

Ray 的设计可以拆成三件事:

  1. task 是一次性跑腿。类比:你把一张便签交给跑腿员,他完成后把结果放回柜台。Ray 的 remote function 返回 future,输入输出都是不可变对象,所以失败后可以按 lineage 重跑。

  2. actor 是长期坐班的人。类比:厨师有自己的锅和调料,下一道菜会受上一道菜准备状态影响。Ray 的 actor 方法串行执行,适合模型副本、仿真环境、参数服务器这类不能轻易序列化的状态。

  3. GCS 把控制状态集中保存,执行组件保持无状态。类比:所有排班表、道具位置、任务血缘都写在共享白板上,调度员坏了可以换一个继续读白板。这样 scheduler、object store、worker 都能独立扩展和恢复。

三件事合起来,Ray 不是单纯”远程函数库”,而是一套能动态生成任务图、调度异构资源、在失败后重建结果的分布式运行时。

案例 1:把普通函数变成远程 task

Section titled “案例 1:把普通函数变成远程 task”
import ray
ray.init()
@ray.remote
def score(policy, state):
return policy(state)
future = score.remote(policy, state)
action = ray.get(future)

逐部分解释:

  • @ray.remote 把本地函数注册成远程任务。
  • score.remote(...) 不会立刻阻塞,而是返回一个 future。
  • ray.get(future) 才真正等待结果;中间系统可以把任务放到任意合适节点执行。

案例 2:用 actor 保存会变化的状态

Section titled “案例 2:用 actor 保存会变化的状态”
@ray.remote(num_gpus=1)
class Trainer:
def __init__(self):
self.step = 0
def update(self, batch):
self.step += 1
return train_one_step(batch)
trainer = Trainer.remote()
loss_id = trainer.update.remote(batch)

逐部分解释:

  • Trainer 是 actor,内部的 step 会在多次调用之间保留。
  • num_gpus=1 是资源声明,告诉调度器这个 actor 需要 GPU。
  • actor 方法也返回 future,但同一个 actor 上的方法按顺序执行,避免状态被并发写乱。

案例 3:强化学习的训练、仿真、服务放在一张图里

Section titled “案例 3:强化学习的训练、仿真、服务放在一张图里”
rollouts = [sim.rollout.remote(policy) for sim in simulators]
ready, rest = ray.wait(rollouts, num_returns=4)
policy = trainer.update.remote(ready)
more = [sim.rollout.remote(policy) for sim in simulators]

逐部分解释:

  • sim.rollout 是仿真 actor,可能有的快、有的慢。
  • ray.wait 先拿到一部分完成结果,不必等所有仿真一起到齐。
  • 训练后的新策略又能立刻喂回仿真,任务图在运行中继续长出来。
  1. 把 Ray 当 Spark 替代品:Spark 擅长大批量数据流水线,Ray 擅长动态小任务和有状态 actor,场景错了会两边都不舒服。

  2. 把 actor 用成全局变量:actor 保留状态也意味着位置较固定,所有人都打同一个 actor 会形成热点。

  3. 忽略对象不可变:Ray object store 里的对象默认不可变,这是为了简化一致性和失败恢复;想原地改共享对象会撞上模型边界。

  4. 以为 scheduler 什么都知道:Ray 的任务图是动态长出来的,调度器不能提前看完整 DAG,所以很多优化只能靠运行时估计。

适用

  • 强化学习:仿真、策略服务、训练紧密循环,且任务耗时差异很大。
  • Python AI 任务编排:很多函数需要并行跑,还要保留少量长期状态。
  • 异构资源调度:CPU 做数据和仿真,GPU 做训练或推理。
  • 需要失败后重跑:无状态 task 可以靠 lineage 重建结果。

不适用

  • 纯 SQL 分析或固定批处理:查询优化器和列式执行引擎更合适。
  • 单机就能跑完的小脚本:分布式调度开销反而是负担。
  • 强一致共享内存:Ray 倾向不可变对象 + actor 串行状态,不是分布式锁数据库。
  • 需要完整预编译静态图的深度学习内核优化:TensorFlow、JAX、编译器栈更专门。
  • 2012 年前后:Spark 用 RDD 解决迭代批处理和交互式分析,但任务模型仍偏批处理阶段。
  • 2016 年:TensorFlow 把神经网络训练画成静态数据流图,训练强,但仿真和动态控制流不自然。
  • 2017 年:Berkeley 团队先提出 Real-Time Machine Learning,指出实时 AI 缺少一套动态分布式执行底座。
  • 2018 年:Ray 在 OSDI 发表,把 task、actor、GCS、bottom-up scheduler 组合成完整系统。
  • 之后几年:Ray 上层长出 RLlib、Tune、Serve 等库,逐渐从强化学习底座扩展到通用 AI infra。
  1. 抽象合并比再造专用系统更有力量:Ray 没给训练、仿真、服务各造一个框架,而是用 task + actor 覆盖它们。
  2. 控制面和数据面分开很关键:GCS 保存 lineage 和对象位置,object store 负责搬数据,scheduler 负责选位置。
  3. 动态任务图需要低延迟调度:论文报告 Ray 在 100 个节点上超过 180 万 task/s,这不是锦上添花,而是细粒度 AI 工作负载的门槛。
  4. 工程取舍都围绕失败恢复:无状态任务可重跑,有状态 actor 需要 checkpoint,GCS 自己要复制。
  • tensorflow-osdi-2016 —— TensorFlow 解决大规模训练图,Ray 补上仿真和服务的动态部分。
  • dataflow-model-2015 —— 两者都把系统问题拆成可组合抽象,只是一个偏流处理,一个偏 AI 执行。
  • mapreduce —— MapReduce 是批处理起点,Ray 是细粒度动态任务的另一条路。
  • borg-omega-kube-2016 —— 都关心调度,但 Ray 面向毫秒级 task,集群调度器面向服务和作业。
  • a3c-2016 —— 强化学习里的异步 actor 思路,是 Ray 最初工作负载的重要来源。
  • ppo —— 论文用 PPO 说明 Ray 如何把仿真 actor 和 GPU 训练拼在一起。
  • alpa-2022 —— Alpa 自动找深度学习并行策略,Ray 更像承载上层编排的运行时。