Ray 解决什么问题
先看一段很普通的代码:你有一个 Python 函数 expensive(x),处理一条数据要 2 秒,你有 1000 条数据要处理。写成列表推导要 2000 秒,太慢了。
results = [expensive(x) for x in items] # 串行,2000 秒你可能的第一反应是用 multiprocessing。能用,但只在一台机器上有效;如果你有 4 台机器,就要自己写 socket 通信、任务分发、结果收集、失败重试。这些代码量往往比业务逻辑还大。
Ray 的做法是:把这些分布式程序里一定会遇到的问题包进一个运行时,让你用接近写普通 Python 的方式来用它。你写的还是函数和类,但运行时帮你做并行、传对象、调度资源和处理失败。
为什么不用现有的方案
如果你试过用 threading、multiprocessing、Celery、Spark 或 Dask 来解决类似问题,你可能已经踩过这些坑:
| 方案 | 擅长 | 典型限制 |
|---|---|---|
threading | I/O 并发 | Python GIL 限制 CPU 密集任务 |
multiprocessing | 单机 CPU 并行 | 跨机器、对象共享和容错要自己做 |
| Celery | 后台任务队列 | 适合“投递、消费、返回”模式,复杂依赖图和大对象共享不是它的强项 |
| Spark | 批处理和 SQL | 数据集模型很强,但细粒度 Python 函数调用和有状态 Actor 不太自然 |
| Dask | Python 数据科学任务图 | 擅长 DataFrame 和 delayed 计算,Actor 模式和在线服务不是重点 |
| Ray | Python 函数、类和对象的分布式运行时 | 需要理解什么时候提交任务、什么时候等待结果、对象放在哪里 |
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
一个常见的错误写法是在循环里立刻等结果:
for item in items:
result = ray.get(expensive.remote(item)) # 每次都等上一个完成这跟串行执行没什么区别:每个任务都要等前一个结束才开始。
正确的做法是先全部提交,再统一取结果:
refs = [expensive.remote(item) for item in items] # 快速提交 1000 个任务
results = ray.get(refs) # 等所有任务完成refs 就是 ObjectRef,也就是任务结果的占位符。.remote() 几乎瞬间返回 ObjectRef,运行时在后台调度和执行任务。等你真正需要结果时,ray.get 才会阻塞等待。
ObjectRef 还承担两个运行时职责:
- 依赖表达:一个任务可以把另一个任务的
ObjectRef当参数传入,Ray 自动建立依赖关系,确保执行顺序正确。 - 对象定位:运行时知道每个对象存在哪个节点的内存里,自动帮你传输,不用自己写网络代码。
Ray 和任务队列有什么不同
你可能会想:这跟 Celery 或 RQ 有什么区别?普通任务队列通常是“投递消息、后台消费、返回状态”的模式。Ray 在几个方面更进一步:
- 依赖传递:Celery 里你要自己管理“任务 A 完成后再提交任务 B”;Ray 用
ObjectRef自动表达依赖。 - 大对象管理:普通队列把参数序列化进消息体;Ray 有独立的对象存储,大对象只存一份,多个任务共享引用。
- 资源感知调度:调度器知道 CPU、GPU 和自定义资源,能把 GPU 任务调度到有 GPU 的节点上。
- 有状态 Actor:Celery Worker 是无状态的;Ray 的 Actor 可以长期存在,持有模型、连接池、缓存等状态。
Ray 的内部结构
写 Ray 代码时你只接触 ray.remote、ray.get、ray.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 帮不了你。