Skip to content

实验 2:ObjectRef 与 wait

本实验关注两个 API:

  • ray.put:把对象放进对象存储。
  • ray.wait:等一部分结果 ready,不必等所有任务完成。

运行示例

bash
uv run python examples/ray_demos/02_objects_and_wait.py

代码核心:

python
matrix_ref = ray.put(np.ones((1000, 1000)))
pending = [consume_matrix.remote(matrix_ref, i) for i in range(6)]

while pending:
    ready, pending = ray.wait(pending, num_returns=1)
    print("ready:", ray.get(ready[0]))

为什么用 ray.put

如果把同一个大矩阵直接传给 6 个任务:

python
pending = [consume_matrix.remote(matrix, i) for i in range(6)]

运行时可能需要重复序列化和传输。使用 ray.put 后,任务参数变成对象引用:

python
matrix_ref = ray.put(matrix)

这表达了“多个任务共享同一个大对象”。

为什么用 ray.wait

ray.get(pending) 会等所有任务完成。很多工程场景里,我们希望谁先完成就先处理谁:

  • 抓取网页。
  • 并行推理。
  • 多个分片处理。
  • 批量任务流水线。

ray.wait 返回两个列表:

python
ready, remaining = ray.wait(pending, num_returns=1)
  • ready:已经完成的引用。
  • remaining:还没完成的引用。

wait 做背压

如果数据源很大,不要一次提交无限任务。可以控制 pending 数量:

python
MAX_PENDING = 32
pending = []

for item in items:
    pending.append(work.remote(item))
    if len(pending) >= MAX_PENDING:
        ready, pending = ray.wait(pending, num_returns=1)
        handle(ray.get(ready[0]))

for ref in pending:
    handle(ray.get(ref))

这会让 Driver 保持稳定内存占用,同时持续给集群喂任务。

Mini Ray 对照

bash
uv run python examples/mini_ray_runtime/demos/02_objects_and_wait.py

Mini Ray 的 wait 用轮询和 threading.Event 实现。真实 Ray 会在分布式对象状态上做更复杂的等待和通知。

小练习

  1. num_returns=1 改成 2,观察每轮输出数量。
  2. ray.wait 增加 timeout=0.1,看看超时时 ready 是否可能为空。
  3. ray.put 去掉,直接传大对象,观察性能和内存变化。

常见错误

错误说明
ObjectRef 当真实对象使用需要 ray.get(ref) 才能拿到真实值
wait 后忘记更新 pending会重复等待已经处理过的引用
一次提交过多任务Driver 内存和调度队列压力上升,需要背压

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