实验 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.pyMini Ray 的 wait 用轮询和 threading.Event 实现。真实 Ray 会在分布式对象状态上做更复杂的等待和通知。
小练习
- 把
num_returns=1改成2,观察每轮输出数量。 - 给
ray.wait增加timeout=0.1,看看超时时 ready 是否可能为空。 - 把
ray.put去掉,直接传大对象,观察性能和内存变化。
常见错误
| 错误 | 说明 |
|---|---|
把 ObjectRef 当真实对象使用 | 需要 ray.get(ref) 才能拿到真实值 |
wait 后忘记更新 pending | 会重复等待已经处理过的引用 |
| 一次提交过多任务 | Driver 内存和调度队列压力上升,需要背压 |