Ray 分布式计算引擎深度实战:从 Plasma 共享内存对象存储、Actor 所有权模型到 GCS 与分布式调度的工程全解
摘要:Ray 把"一台机器上的多进程编程"无缝延伸到上千节点,靠的不是某种魔法调度器,而是三块互相咬合的基石:以 Arrow 为内存布局的分布式对象存储、以所有权(ownership)为核心的分布式引用计数与血缘重建、以及一个刻意保持"薄"的全局控制服务(GCS)。本文从这三条主线切入,配合可运行代码,讲清 Ray 的数据怎么传、对象什么时候被释放、故障如何恢复、任务落在哪台机器上,以及生产上真正会踩的坑。
一、Ray 要解决的真实问题
Python 生态里从来不缺并行:multiprocessing 能起进程,concurrent.futures 能提交任务,Dask 能分数据。但它们大多停在"单机多核"这一层,跨机就要自己写 gRPC、自己做服务发现、自己解决失败重试。
Ray 的设计哲学很朴素:让分布式变成函数调用的语法糖。
import ray
ray.init(address="auto") # 连到已有集群,或 ray.init() 起本地集群
@ray.remote
def parse(raw: bytes) -> dict:
return {"size": len(raw), "lines": raw.count(b"\n")}
refs = [parse.remote(open(f"part-{i}.bin", "rb").read()) for i in range(1000)]
results = ray.get(refs) # 1000 个任务自动分发到整个集群
注意这段代码里发生的三件事:函数被序列化后通过 GCS 广播到所有节点的 worker;返回值写进每个节点本地的共享内存对象存储;refs 只是 ObjectRef 句柄,真正的数据搬移被推迟到 ray.get()。理解这三点,就理解了 Ray 的一半。
二、分布式对象存储:从 Plasma 到所有权模型
2.1 Plasma 时代的设计
Ray 早期的对象存储叫 Plasma,是一个每节点一个进程的 /dev/shm 共享内存管理器。它的价值在于同机零拷贝:
import numpy as np, ray
@ray.remote
def make():
return np.zeros((1024, 1024), dtype=np.float32) # 4MB
@ray.remote(num_cpus=1)
def consume(arr): # 参数通过共享内存传递,不经过网络栈
return float(arr.sum())
ref = make.remote()
print(ray.get(consume.remote(ref)))
make() 在节点 A 上把 ndarray 按 Apache Arrow 布局写进 Plasma;如果 consume() 恰好也被调度到节点 A,它通过 mmap 直接拿到同一块物理页,连一次 memcpy 都没有。这是因为 Arrow 的列式内存布局天然具备"偏移量 + 长度"的自描述性,Ray 只需传递一个文件描述符和 offset。
跨机时,Ray 会把 Arrow 的 RecordBatch 切片(chunk)走 gRPC 流式推送到目标节点的 Plasma,再物化成本地共享内存对象——同一份数据在不同节点上各有一份副本,而不是中心化存储。
2.2 Plasma 的致命缺陷
Plasma 采用集中式引用计数:所有 ObjectRef 的增减都要通知创建该对象的那个 Plasma Store,由它单方面决定何时释放。
这带来两个生产级问题:
- 引用计数成为热点:短生命周期的小对象成千上万时,计数 RPC 把单进程 Store 打满;
- 语义不安全:没有所有权归属,一旦持有者崩溃,内存是泄漏还是回收全靠超时猜测;
ray.get()卡死在"对象永不出现"上成了经典故障。
2.3 Ownership + 分布式引用计数
Ray 1.x 之后改为所有权模型:每个 ObjectRef 都有一个明确的 owner(创建该对象的 worker),owner 维护该对象的引用计数表和位置信息。引用变化以增量 RPC 批量上报给 owner,owner 归零则通知各节点 Plasma 释放。
@ray.remote
def a(x):
return x + 1
@ray.remote
def borrow(r): # 把 ObjectRef 传进另一个 task,形成借用
return ray.get(r)
r0 = a.remote(1)
r1 = a.remote(r0) # r0 被 task 图借用,owner 会记 borrow_id
r2 = borrow.remote(r0)
借用(borrowing)是有超时监控的:如果借用方崩溃,owner 在租约(lease)超时后回收引用。这套机制把"分布式内存管理"收敛成了"单点 owner 记账 + 增量同步",复杂度大幅下降。
三、容错:血缘重建而不是全局快照
Ray 对无状态 task 采用 lineage reconstruction(血缘重建):对象丢失时,沿任务依赖图向上找到还能访问的输入,重新执行。
@ray.remote
def stage1(x): return x * 2
@ray.remote
def stage2(x): return x + 1
r = stage2.remote(stage1.remote(10))
ray.get(r) # 若 stage1 的结果所在节点宕机,Ray 自动重跑 stage1 → stage2
好处是零运行时开销——不需要写 checkpoint、不需要像 Flink 那样做分布式快照。代价是:
- 非确定性任务(读时钟、随机数、外部副作用)重建结果会漂移,务必注入种子;
- 长链条任务的最坏恢复时间等于整条链的重算时间。
对于有状态 Actor,Ray 不做自动重建(除非显式开 max_restarts),因为重放消息流的代价不可控:
@ray.remote(num_cpus=2, max_restarts=3, max_task_retries=2)
class Shard:
def __init__(self, shard_id):
self.shard_id = shard_id
self.state = {}
def put(self, k, v):
self.state[k] = v
return len(self.state)
shard = Shard.options(name="shard-0", lifetime="detached").remote(0)
detached + name 让 Actor 脱离 driver 生命周期,可被 ray.get_actor("shard-0") 重新获取——这是生产上做常驻服务的标准姿势。
四、GCS:刻意做"薄"的控制面
GCS(Global Control Service)保存集群的元数据:节点列表、资源容量、Actor 注册信息、已提交的任务规格。它不参与数据传输,也不做细粒度调度决策。
GCS 的可用性靠租约保证:节点每几百毫秒续约一次,超时未续约的节点被标记为 dead,其上 Actor 按策略重启。为了避免每次任务提交都打 GCS,Ray 引入了 gcs-aware 客户端缓存 + 批提交;GCS 早期是单点,后续版本支持基于 Redis/内存存储的 HA 部署。
工程含义很直接:GCS 挂了,正在跑的任务不会立刻失败(数据和执行路径不经过它),但新任务无法提交、新 Actor 无法注册。所以 GCS 的可用性要求低于数据面,但仍必须上 HA。
五、调度:自下而上的分布式调度器
Ray 的调度分两类:
- Task/Actor 调度:默认自下而上(bottom-up)。每个节点本地有调度器,优先本地排队;本地资源不足时按策略 spillback 到其它节点。这与 Kubernetes 中心化 scheduler 截然不同,换来的是十万级任务/秒的提交吞吐。
- Placement Group(放置组):用于需要原子性资源预留的场景(如分布式训练的 8 卡 worker)。
from ray.util.placement_group import placement_group
pg = placement_group([
{"CPU": 4, "GPU": 1}, # bundle 0: trainer
{"CPU": 8}, # bundle 1: data loader
], strategy="PACK")
ray.get(pg.ready()) # 原子预留,要么全满足要么一直等
@ray.remote(num_cpus=4, num_gpus=1)
def train(): return "ok"
ray.get(train.options(placement_group=pg, placement_group_bundle_index=0).remote())
自定义资源可以做软性拓扑约束:
@ray.remote(num_cpus=1, resources={"zone:a": 0.01}) # 0.01 的分数资源表达"偏好"
def zone_aware(x): return x
Ray 2.x 之后引入了基于 bundle 的新调度器,统一了 task/actor/ placement group 的路径,并支持 SPREAD、STRICT_PACK 等策略。经验法则:小任务靠 bottom-up 的随机性自然均衡;大块资源必须用 placement group 显式预留,否则会出现"两个 job 各占半张卡"的死锁式饥饿。
六、对象溢出与内存红线
对象存储默认占节点可用内存的 30%。超出后 Ray 会把对象 spill 到本地磁盘,从磁盘读回时再自动 restore:
ray start --head \
--object-store-memory=$((40*1024*1024*1024)) \
--system-config='{"object_spilling_config":"{\"type\":\"filesystem\",\"params\":{\"directory_path\":\"/data/ray/spill\"}}"}'
生产建议:
- spill 目录必须放在独立的 NVMe 上,与系统盘隔离,否则节点 IO 打满会引发连锁超时;
- 不要指望 spill 当银弹。频繁 spill 说明的是数据粒度太粗,应该改成流式分块处理。
七、生产上真正会踩的坑
- 小对象风暴:把百万级小对象塞进对象存储,管理开销远大于数据本身。解决方式是批处理(一个 task 处理一批,
ray.get()一次拿一批),或者让下游以ObjectRef形式消费,避免数据回到 driver。 ray.get()位置错误:在循环里逐个ray.get()会串行化 pipeline。正确姿势是ray.wait(refs, num_returns=k, fetch_local=False)做流式消费:
pending = list(refs)
while pending:
done, pending = ray.wait(pending, num_returns=4, timeout=30)
for r in done:
process(ray.get(r))
- driver 成瓶颈:所有结果都回 driver 就退化成单机。让 task 之间直接传
ObjectRef,形成分布式 DAG,driver 只做编排。 - Actor 串行化:Actor 方法默认单线程串行执行。CPU 密集场景要么用
max_concurrency配 asyncio,要么拆成多个 Actor 分片。 - 序列化意外:Ray 用 cloudpickle 捕获闭包。捕获了大对象或不可序列化的句柄(数据库连接、锁)会在提交时静默变慢甚至失败。养成"显式传参,不靠闭包"的习惯。
八、工程观点
Ray 最值得借鉴的不是某个组件,而是它对控制面与数据面的严格分界:GCS 只管元数据和租约,调度下放到每个节点,对象存储用共享内存 + Arrow 把"传输"变成"映射"。这三条让它在 AI workloads(训练、超参搜索、批量推理)这种"任务巨多、单任务不大、数据是大张量"的负载上碾压传统中心化调度器。
反过来说,如果你的负载是稳定的长驻服务、强事务语义、需要秒级精确恢复,Ray 并不是答案——血缘重建的不确定性和 Actor 的弱容错保证会让你痛苦。技术选型上认清这一点,比学会任何一个 API 都重要。

发表评论 取消回复