Ray 反模式:循环内调用 ray.get 破坏并行性——正确聚合 ObjectRef 的实践指南
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
导读
Ray 是面向 AI 与分布式计算场景的运行时,ray.get()是其最常用的结果获取接口,但它是一个阻塞调用。本篇文章聚焦 Ray Core 官方设计模式文档中记录的经典反模式——"在循环中调用ray.get()",剖析它为何会彻底抹杀掉远程任务的并行度,并给出"先批量提交、再一次性聚合"的修复方案。读完本文,你将掌握ray.get()的阻塞语义、ObjectRef 的异步提交机制、ray.get(list)批量聚合的用法,以及与之相关的多个ray.get系列反模式(嵌套调用、无用调用、提交顺序陷阱、对象过多)的规避思路。
反模式速览:一句话结论
官方文档的 TLDR 非常直接:
避免在循环中调用
ray.get(),因为它是阻塞调用;只在最终结果上使用ray.get()。
ray.get()用于取回远程函数的执行结果,但它会阻塞等待,直到请求的结果可用为止。如果在循环里调用ray.get(),循环体必须等到该次调用解析完成才会进入下一轮迭代;如果远程任务恰好也在同一个循环内提交,那么"提交任务 → 等待结果 → 再提交下一个任务"的串行节奏会让并行度归零——上一轮任务还没跑完,下一轮任务根本不会被调度出去。
阻塞语义的底层依据
为什么ray.get()必然阻塞?在 Ray 的 Python 源码 python/ray/_private/worker.py 中,ray.get()的 docstring 明确写道:
"This method blocks until the object corresponding to the object ref is available in the local object store. If this object is not in the local object store, it will be shipped from an object store that has it (once the object has been created)."
即:ray.get()会一直阻塞,直到对象引用(ObjectRef)对应的对象在本地对象存储中可用;如果本地没有,它还会等待对象从其他节点(持有该对象的对象存储)传输过来。这就意味着一次ray.get()等待的不只是任务计算完成,还包含对象物化、跨节点传输的完整链路。将其放入循环体,等于把每一次任务的"完成事件"强制转成了循环的"进度闸门"。
此外,源码还展示了ray.get()的几个关键签名特性(python/ray/_private/worker.py):
- 支持传入单个 ObjectRef或ObjectRef 列表两种形式;
- 传列表时,返回结果的顺序与输入列表顺序严格一致("Ordering for an input list of object refs is preserved");
- 支持
timeout参数与异步上下文(async 环境下应改用await object_ref或asyncio.gather(*object_refs),直接调用会触发警告)。
其中"列表顺序保持"这一点,正是修复循环反模式后仍能按提交顺序拿到结果的基础——ray.get(refs)返回的列表与refs一一对应,语义与循环内逐个ray.get()完全相同,但并行行为天差地别。
官方代码示例:反模式 vs 正确写法
官方文档对应的完整可运行示例位于 doc/source/ray-core/doc_code/anti_pattern_ray_get_loop.py,其反模式与修复对比代码如下:
import ray ray.init() @ray.remote def f(i): return i # Anti-pattern: no parallelism due to calling ray.get inside of the loop. sequential_returns = [] for i in range(100): sequential_returns.append(ray.get(f.remote(i))) # Better approach: parallelism because the tasks are executed in parallel. refs = [] for i in range(100): refs.append(f.remote(i)) parallel_returns = ray.get(refs)示例文件末尾还带有断言校验,确保两种写法的结果完全一致:
assert sequential_returns == parallel_returns逐行解读反模式代码
@ray.remote def f(i): return i定义一个远程任务:它会被调度到 Ray 的 worker 进程上执行,返回的i会写入对象存储;f.remote(i)调用并不会真正执行函数体,而是提交任务并立即返回一个 ObjectRef(分布式引用句柄);- 问题出在
ray.get(f.remote(i))的组合:f.remote(i)刚提交完任务,ray.get()立刻阻塞等待该任务完成并返回结果。于是第i+1轮迭代必须等到第i轮任务的结果物化之后才会发起; - 最终 100 个任务被串行执行,与本地 for 循环逐次计算没有本质区别,分布式集群的全部算力都被闲置。
正确写法的两个关键动作
- 先批量提交:第一个循环只调用
f.remote(i)收集refs,100 个任务在同一时刻全部进入调度队列,Ray 的调度器(GCS / 核心调度组件)会把它们分发到集群空闲资源上并行执行; - 后一次性聚合:
ray.get(refs)接收整个 ObjectRef 列表,一次性等待全部任务完成并取回结果。得益于源码中"结果顺序与输入顺序一致"的保证,parallel_returns与sequential_returns内容完全相同。
修复后,假设集群有N个可用 CPU 槽位,100 个任务可同时铺开执行,总耗时从"100 个任务串行之和"下降到"约 100/N 个任务串行之和",这正是 Ray 分布式并行的核心收益。
图解:阻塞循环与并行调度
官方文档为这一反模式配了流程示意图 doc/source/ray-core/images/ray-get-loop.svg:
图示直观对比了两条路径:**上方(反模式)**在每次调度远程任务后立即调用ray.get(),循环被阻塞,任务逐个串行完成;**下方(正确)**先调度全部远程调用使其并行处理,再一次性请求所有结果。换句话说:调度(submit)与等待(block)必须分离,所有远程任务先全部提交、再统一取回,才能让后台并行真正发生。
实战扩展:批量取回的边界与技巧
掌握了"分离提交与等待"的原则后,还有几个与批量ray.get()直接相关的实战细节值得注意:
1. 列表 vs 逐个调用的语义差异
ray.get(refs)一次性等待全部任务完成,等价于"全部完成后再返回"的屏障(barrier)语义;而循环内逐个ray.get()是"完成一个拿一个"的串行消费语义。当任务的耗时差异很大时,前者会等待最慢的任务(straggler),这正是 ray-get-submission-order 反模式 讨论的问题——如果需要按完成顺序处理结果,应改用ray.wait()而非ray.get()。
2. 对象数量过多时的分批策略
如果一次性ray.get()数千上万个对象,Ray 需要同时把这些对象拉到调用方节点,可能引发堆内存溢出(heap OOM)或对象存储空间不足。官方建议按批次获取与处理,处理完一批后 Ray 会淘汰该批对象为后续批次腾出空间,详见 ray-get-too-many-objects 反模式。
3. 更进一步的思路:尽量不调用 ray.get
对于中间结果,最优做法其实是"能不 get 就不 get":直接把 ObjectRef 作为参数传给下一个远程任务,让 Ray 在任务依赖图中自动解析,避免对象先传到驱动(driver)再二次传输的额外开销,详见 unnecessary-ray-get 反模式。而在嵌套任务中,如果任务内部再对参数执行ray.get(),还可能引发资源抢占与潜在死锁,Ray 虽会临时释放调用方 CPU 资源以避免死锁,但会损害性能与稳定性,相关讨论见 nested-ray-get 反模式。
ray.get 系列反模式全景
ray.get()是 Ray 编程中最容易误用的 API 之一。官方 Design Patterns & Anti-patterns 目录 中,除本文主题外,还收录了以下直接相关的反模式条目,建议组合阅读:
| 反模式文档 | 核心问题 | 修复方向 |
|---|---|---|
| nested-ray-get | 在任务内部对参数调用ray.get(),阻塞并可能引发死锁 | 将 ObjectRef 作为直接参数传递,由 Ray 隐式解析 |
| unnecessary-ray-get | 对中间结果过早ray.get(),对象被迫传输到 driver 造成额外拷贝 | 只对最终结果调用ray.get() |
| ray-get-submission-order | 按提交顺序处理结果,被慢任务(straggler)拖累 | 用ray.wait()按完成顺序处理 |
| ray-get-too-many-objects | 一次取回过多对象导致堆内存 / 对象存储溢出 | 分批获取并处理 |
这些条目共同勾勒出一条统一的编程准则:ray.get()是"终结操作",应尽可能晚、尽可能少、尽可能批量地调用,把 ObjectRef 的传递与解析交给 Ray 的任务依赖机制去完成。
总结
循环内调用ray.get()是最常见也最容易忽视的 Ray 性能陷阱之一。它的本质是把阻塞语义与提交语义耦合在了一起:任务提交(f.remote(i))本是异步的,ray.get()却将其强行拉回同步,导致分布式集群退化为串行执行。正确的姿势是:
- 先用循环批量提交全部任务,收集 ObjectRef 列表;
- 循环结束后统一调用
ray.get(refs)一次性取回结果(列表顺序与提交顺序一致); - 依据业务需要,结合
ray.wait()按完成顺序消费、分批取回控制内存,并在中间环节尽量传递 ObjectRef 而非调用ray.get()。
实践这一原则,你的 Ray 任务才能真正利用集群的多节点并行能力。完整的可运行示例可参考 anti_pattern_ray_get_loop.py,更多模式与反模式请查阅 Ray Core 设计模式目录。
【免费下载链接】rayRay is an AI compute engine. Ray consists of a core distributed runtime and a set of AI Libraries for accelerating ML workloads.项目地址: https://gitcode.com/gh_mirrors/ra/ray
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考