Skip to content

Ray 解决什么问题

先看一段很普通的代码:你有一个 Python 函数 expensive(x),处理一条数据要 2 秒,你有 1000 条数据要处理。写成列表推导要 2000 秒,太慢了。

python
results = [expensive(x) for x in items]  # 串行,2000 秒

你可能的第一反应是用 multiprocessing。能用,但只在一台机器上有效;如果你有 4 台机器,就要自己写 socket 通信、任务分发、结果收集、失败重试。这些代码量往往比业务逻辑还大。

Ray 的做法是:把这些分布式程序里一定会遇到的问题包进一个运行时,让你用接近写普通 Python 的方式来用它。你写的还是函数和类,但运行时帮你做并行、传对象、调度资源和处理失败。

为什么不用现有的方案

如果你试过用 threadingmultiprocessing、Celery、Spark 或 Dask 来解决类似问题,你可能已经踩过这些坑:

方案擅长典型限制
threadingI/O 并发Python GIL 限制 CPU 密集任务
multiprocessing单机 CPU 并行跨机器、对象共享和容错要自己做
Celery后台任务队列适合“投递、消费、返回”模式,复杂依赖图和大对象共享不是它的强项
Spark批处理和 SQL数据集模型很强,但细粒度 Python 函数调用和有状态 Actor 不太自然
DaskPython 数据科学任务图擅长 DataFrame 和 delayed 计算,Actor 模式和在线服务不是重点
RayPython 函数、类和对象的分布式运行时需要理解什么时候提交任务、什么时候等待结果、对象放在哪里

Ray 的设计目标不是“让 for 循环快一点”。它更想解决的是:当计算要跨进程、跨机器、跨 CPU/GPU 时,业务代码还能不能保持清楚。Task、ObjectRef 和 Actor 就是 Ray 给出的三种基本工具。

Ray Core 的三个基本抽象

Ray Core 提供三个核心概念,后面所有章节都是在展开它们:

Task 是远程函数调用。你用 @ray.remote 装饰一个函数,调用 .remote() 时 Ray 不会立即执行它,而是把调用提交给运行时,运行时找一个有空闲资源的 Worker 去执行。

ObjectRef 是结果的占位符。.remote() 返回的不是计算结果,而是一个引用。你什么时候需要结果,什么时候用 ray.get(ref) 去取。这个设计让 Ray 能并行提交大量任务,不用等前一个完成再提交下一个。

Actor 是有状态的远程对象。Task 执行完就结束了,不保留任何状态。但有些场景(比如推理服务、参数服务器、游戏环境)需要一个长期存活、持续更新状态的对象,这就是 Actor。

这三个概念组合起来,可以覆盖很多场景:批量数据处理、模型训练、超参搜索、在线推理服务、强化学习环境、Agent 工作流。

为什么 Ray 要引入 ObjectRef

一个常见的错误写法是在循环里立刻等结果:

python
for item in items:
    result = ray.get(expensive.remote(item))  # 每次都等上一个完成

这跟串行执行没什么区别:每个任务都要等前一个结束才开始。

正确的做法是先全部提交,再统一取结果:

python
refs = [expensive.remote(item) for item in items]  # 快速提交 1000 个任务
results = ray.get(refs)                              # 等所有任务完成

refs 就是 ObjectRef,也就是任务结果的占位符。.remote() 几乎瞬间返回 ObjectRef,运行时在后台调度和执行任务。等你真正需要结果时,ray.get 才会阻塞等待。

ObjectRef 还承担两个运行时职责:

  1. 依赖表达:一个任务可以把另一个任务的 ObjectRef 当参数传入,Ray 自动建立依赖关系,确保执行顺序正确。
  2. 对象定位:运行时知道每个对象存在哪个节点的内存里,自动帮你传输,不用自己写网络代码。

Ray 和任务队列有什么不同

你可能会想:这跟 Celery 或 RQ 有什么区别?普通任务队列通常是“投递消息、后台消费、返回状态”的模式。Ray 在几个方面更进一步:

  • 依赖传递:Celery 里你要自己管理“任务 A 完成后再提交任务 B”;Ray 用 ObjectRef 自动表达依赖。
  • 大对象管理:普通队列把参数序列化进消息体;Ray 有独立的对象存储,大对象只存一份,多个任务共享引用。
  • 资源感知调度:调度器知道 CPU、GPU 和自定义资源,能把 GPU 任务调度到有 GPU 的节点上。
  • 有状态 Actor:Celery Worker 是无状态的;Ray 的 Actor 可以长期存在,持有模型、连接池、缓存等状态。

Ray 的内部结构

写 Ray 代码时你只接触 ray.remoteray.getray.put 这几个 API,但它们背后有好几层组件配合工作:

  • 用户 API:你日常写的代码在这一层。
  • Core Worker:把你的 remote() 调用翻译成运行时任务,管理 ObjectRef 和依赖关系。
  • Raylet + Object Store:每个节点一个 Raylet 进程,负责本地资源管理、Worker 进程池和对象存储。
  • GCS:集群级元数据服务,管理节点列表、资源视图和 Actor 位置。

后面的源码导读章节会逐层展开这些组件。现在你只需要知道:API 很薄,真正的逻辑在下面。

本教程关注的边界

Ray 生态很大:Ray Data、Ray Train、Ray Tune、Ray Serve、RLlib、Jobs、Dashboard、KubeRay……本教程只讲 Ray Core。

原因很简单:Core 是地基。Task、Actor、Object Store 和 Scheduler 是所有上层库的共同基础。你先理解 Core,再去看 Data 或 Serve,会清楚它们在 Core 之上做了什么封装。

常见误解

Ray 不会自动让任意 Python 程序变快。你需要自己把计算拆成合理粒度的任务,避免在循环里频繁 ray.get,也要理解大对象传输和资源声明的实际成本。如果一个问题本身没有并行性,Ray 帮不了你。

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