Skip to content

对象存储与序列化

Ray 的对象系统是理解性能的关键。很多新手以为瓶颈在 CPU 计算,其实瓶颈经常在对象复制、序列化、跨节点传输和内存压力上。

为什么需要对象存储

远程任务有两个基本需求:

  1. 把参数交给 Worker。
  2. 把返回值交给后续任务或 Driver。

如果每次任务调用都通过 RPC 直接传完整对象,会遇到几个问题:

  • 大对象被反复复制,浪费内存和带宽。
  • 多个任务需要同一份数据时,每份都要传一次。
  • 跨节点传输没有统一管理,对象存在哪里、谁负责传输都不清楚。

Ray 用 Object Store 把“数据”从“任务消息”里分离出来:对象存到存储里,任务之间传的是引用(ObjectRef),不是数据本身。

ray.put

ray.put 把本地对象放入对象存储,返回 ObjectRef

python
weights_ref = ray.put(model_weights)
refs = [score.remote(weights_ref, shard) for shard in shards]

这样做的核心价值是共享:大对象只进入对象系统一次,后续任务传的是引用。

对象是不可变的

Ray 的对象存储按不可变语义工作。一个对象被创建后(通过任务返回值或 ray.put),其他任务只能读取,不能原地修改。

这个约束让运行时可以放心地:

  • 在节点之间复制对象,不用担心数据不一致。
  • 多个任务共享同一份数据,不用担心并发写。
  • 用引用计数管理生命周期。

如果你需要可变状态,比如一个会持续更新的模型或计数器,应该用 Actor。

ObjectRef 是数据依赖

当你把 ObjectRef 传给远程函数时,Ray 不会把引用本身当普通参数交给函数,而是把它作为依赖处理。

python
@ray.remote
def load():
    return [1, 2, 3]

@ray.remote
def mean(values):
    return sum(values) / len(values)

values_ref = load.remote()
mean_ref = mean.remote(values_ref)

运行时会确保 mean 拿到的是 [1, 2, 3],而不是 ObjectRef 对象。

对象位置

在多节点集群里,对象可能存在于某个节点的对象存储中。

调度器会尽量考虑数据局部性:如果一个任务需要的大对象已经在某个节点上,把任务放在那个节点通常更划算。

Spill(溢写)

对象存储的内存有限。当对象太多放不下时,Ray 可以把不常用的对象 spill 到磁盘或外部存储,避免 OOM。但 spill 有性能代价,从磁盘读回来比从内存读慢很多。

实际开发中要关注:

  • Object Store 的内存使用量(Dashboard 里可以看到)。
  • 是否频繁发生 spill(频繁说明对象太多或太大)。
  • 大对象是否被重复生成(应该用 ray.put 共享)。
  • 是否有 Driver 或 Actor 长期持有引用,导致对象无法释放。

序列化

Ray 需要把 Python 对象序列化,才能跨进程或跨节点传输。不是所有对象都适合序列化。

常见问题:

问题处理方式
传入不可序列化对象,如打开的文件句柄在 Worker 内重新打开资源,或用 Actor 封装
大模型权重反复作为普通参数传入使用 ray.put 或 Actor 持有
返回对象过大拆分对象、流式处理、使用 Ray Data
对象里包含线程锁、连接对象避免作为 Task 参数传输

Mini Ray 的对象存储

Mini Ray 用一个字典和 threading.Event 模拟对象存储:

python
class ObjectStore:
    def reserve(self) -> ObjectRef:
        ...

    def set_result(self, ref, value):
        entry.value = value
        entry.ready.set()

    def get(self, ref):
        entry.ready.wait()
        return entry.value

它没有共享内存,也没有跨节点位置管理,但保留了两个关键语义:

  • .remote() 先生成一个未完成对象。
  • get() 等对象 ready 后再返回。

这就是理解 Ray 对象系统的最小闭环。

小练习

把下面代码改成先 ray.put 一次大对象,再让多个任务共享:

python
refs = [score.remote(model_weights, shard) for shard in shards]

参考答案:

python
weights_ref = ray.put(model_weights)
refs = [score.remote(weights_ref, shard) for shard in shards]

再思考一个问题:如果 model_weights 是一个会被频繁更新的对象,为什么 Actor 可能比 ray.put 更合适?

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