对象存储与序列化
Ray 的对象系统是理解性能的关键。很多新手以为瓶颈在 CPU 计算,其实瓶颈经常在对象复制、序列化、跨节点传输和内存压力上。
为什么需要对象存储
远程任务有两个基本需求:
- 把参数交给 Worker。
- 把返回值交给后续任务或 Driver。
如果每次任务调用都通过 RPC 直接传完整对象,会遇到几个问题:
- 大对象被反复复制,浪费内存和带宽。
- 多个任务需要同一份数据时,每份都要传一次。
- 跨节点传输没有统一管理,对象存在哪里、谁负责传输都不清楚。
Ray 用 Object Store 把“数据”从“任务消息”里分离出来:对象存到存储里,任务之间传的是引用(ObjectRef),不是数据本身。
ray.put
ray.put 把本地对象放入对象存储,返回 ObjectRef。
weights_ref = ray.put(model_weights)
refs = [score.remote(weights_ref, shard) for shard in shards]这样做的核心价值是共享:大对象只进入对象系统一次,后续任务传的是引用。
对象是不可变的
Ray 的对象存储按不可变语义工作。一个对象被创建后(通过任务返回值或 ray.put),其他任务只能读取,不能原地修改。
这个约束让运行时可以放心地:
- 在节点之间复制对象,不用担心数据不一致。
- 多个任务共享同一份数据,不用担心并发写。
- 用引用计数管理生命周期。
如果你需要可变状态,比如一个会持续更新的模型或计数器,应该用 Actor。
ObjectRef 是数据依赖
当你把 ObjectRef 传给远程函数时,Ray 不会把引用本身当普通参数交给函数,而是把它作为依赖处理。
@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 模拟对象存储:
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 一次大对象,再让多个任务共享:
refs = [score.remote(model_weights, shard) for shard in shards]参考答案:
weights_ref = ray.put(model_weights)
refs = [score.remote(weights_ref, shard) for shard in shards]再思考一个问题:如果 model_weights 是一个会被频繁更新的对象,为什么 Actor 可能比 ray.put 更合适?