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,由它单方面决定何时释放。

这带来两个生产级问题:

  1. 引用计数成为热点:短生命周期的小对象成千上万时,计数 RPC 把单进程 Store 打满;
  2. 语义不安全:没有所有权归属,一旦持有者崩溃,内存是泄漏还是回收全靠超时猜测;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 说明的是数据粒度太粗,应该改成流式分块处理。

七、生产上真正会踩的坑

  1. 小对象风暴:把百万级小对象塞进对象存储,管理开销远大于数据本身。解决方式是批处理(一个 task 处理一批,ray.get() 一次拿一批),或者让下游以 ObjectRef 形式消费,避免数据回到 driver。
  2. 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))
  1. driver 成瓶颈:所有结果都回 driver 就退化成单机。让 task 之间直接传 ObjectRef,形成分布式 DAG,driver 只做编排。
  2. Actor 串行化:Actor 方法默认单线程串行执行。CPU 密集场景要么用 max_concurrency 配 asyncio,要么拆成多个 Actor 分片。
  3. 序列化意外:Ray 用 cloudpickle 捕获闭包。捕获了大对象或不可序列化的句柄(数据库连接、锁)会在提交时静默变慢甚至失败。养成"显式传参,不靠闭包"的习惯。

八、工程观点

Ray 最值得借鉴的不是某个组件,而是它对控制面与数据面的严格分界:GCS 只管元数据和租约,调度下放到每个节点,对象存储用共享内存 + Arrow 把"传输"变成"映射"。这三条让它在 AI workloads(训练、超参搜索、批量推理)这种"任务巨多、单任务不大、数据是大张量"的负载上碾压传统中心化调度器。

反过来说,如果你的负载是稳定的长驻服务、强事务语义、需要秒级精确恢复,Ray 并不是答案——血缘重建的不确定性和 Actor 的弱容错保证会让你痛苦。技术选型上认清这一点,比学会任何一个 API 都重要。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部