
人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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 中批量提交一批彼此独立的远程任务后开发者最常见的做法是按提交顺序逐个ray.get()拉取结果。然而任务的执行时长各不相同按提交顺序等待会让较慢的掉队任务straggler阻塞后续处理白白浪费已经完成任务的空闲时间。本文以 Ray 官方 Patterns 文档中的反模式讲解为核心深入剖析ray.get()按提交顺序消费结果的性能陷阱并结合仓库源码与示例代码给出基于ray.wait()按完成顺序消费结果的正确写法帮助你最大化分布式任务的并行收益。反模式概述为什么按提交顺序处理结果会拖慢整体运行TLDR不要使用ray.get()按提交顺序处理相互独立的任务结果因为结果就绪的顺序往往与提交顺序不一致。文档原文ray-get-submission-order.rst给出的核心结论是Avoid processing independent results in submission order usingray.get()since results may be ready in a different order than the submission order.正确做法是使用ray.wait()按完成顺序completion order逐个处理已就绪的结果从而缩短整体完成时间total time to completion。问题本质阻塞等待与乱序完成当一批任务被批量提交后每个任务执行耗时不同有些任务很快返回有些任务却执行很久。如果严格按提交顺序处理结果就可能出现如下场景先提交的慢任务迟迟未完成ray.get()一直阻塞后提交的快任务其实早已完成结果已经就绪却只能在队列里排队等待整批任务的完成时间被最慢的掉队任务支配且快任务的等待时间被完全浪费。从源码看ray.get()是一个阻塞调用。在 python/ray/_private/worker.py 中其实现明确指出This method blocks until the object corresponding to the object ref is available in the local object store.也就是说ray.get(ref)会一直等到该 ObjectRef 对应的对象就绪才返回。逐个调用时循环只能在前一个调用解决后继续快任务的结果即使已就绪也无法被及时消费。下图直观对比了按提交顺序处理与按完成顺序处理两种策略的差异反模式与正确模式的代码对比文档引用的完整示例位于 anti_pattern_ray_get_submission_order.py以下代码可直接运行复现两种写法的差异。import random import time import ray ray.init() ray.remote def f(i): time.sleep(random.random()) return i # Anti-pattern: process results in the submission order. sum_in_submission_order 0 refs [f.remote(i) for i in range(100)] for ref in refs: # Blocks until this ObjectRef is ready. result ray.get(ref) # process result sum_in_submission_order sum_in_submission_order result # Better approach: process results in the completion order. sum_in_completion_order 0 refs [f.remote(i) for i in range(100)] unfinished refs while unfinished: # Returns the first ObjectRef that is ready. finished, unfinished ray.wait(unfinished, num_returns1) result ray.get(finished[0]) # process result sum_in_completion_order sum_in_completion_order result反模式部分先通过列表推导式一次性提交 100 个任务f.remote(i)然后for ref in refs按提交顺序逐个ray.get(ref)。由于f内部time.sleep(random.random())每个任务耗时在 01 秒之间随机先提交的任务未必先完成。当遇到先提交的慢任务时循环被阻塞而后面已完成的任务结果只能闲置。正确模式部分同样先批量提交任务但用一个while unfinished循环维护未完成集合每次通过ray.wait(unfinished, num_returns1)取出最先就绪的那一个ObjectRef 进行处理再更新未完成集合直至全部处理完毕。这样任何任务一完成就立刻被消费整体耗时只取决于最慢任务的耗时而不是最慢任务 其余任务的总等待量。示例文件末尾还附带断言验证两种方式结果一致assert sum_in_submission_order sum_in_completion_order即按完成顺序处理不会改变计算语义只是减少了等待时间可以在实际项目中放心替换。ray.wait 的核心机制与参数说明ray.wait()是替代ray.get()循环的关键 API其定义位于 python/ray/_private/worker.py。其签名如下PublicAPI def wait( ray_waitables: List[Union[ObjectRef, ObjectRefGenerator]], *, num_returns: int 1, timeout: Optional[float] None, fetch_local: bool True, ) - Tuple[ List[Union[ObjectRef, ObjectRefGenerator]], List[Union[ObjectRef, ObjectRefGenerator]], ]:返回值ray.wait()返回两个列表ready第一个列表已经就绪的 ObjectRef或 ObjectRefGenerator即其下一个引用对应对象已在对象存储中可用remaining第二个列表尚未就绪的剩余对象引用。关键参数参数默认值说明ray_waitables无待等待的ObjectRef或ObjectRefGenerator列表必须去重。源码中若len(ray_waitables) ! len(set(ray_waitables))会抛出ValueError见 worker.pynum_returns1需要返回的就绪对象个数。源码要求其大于 0 且不能超过传入列表长度否则抛ValueError见 worker.pytimeoutNone最大等待秒数。为None时一直等到指定数量的对象就绪才返回超时则提前返回当前已就绪的对象。源码中timeout为负会抛ValueError见 worker.pyfetch_localTrue为True时等待对象下载到本地节点后才标记为就绪为False时只要对象在集群任意位置可用就立即返回见 worker.py顺序保持保证ray.wait()保持输入列表的顺序如果 A 在输入列表中排在 B 前面且二者都进入 ready 列表则 A 仍排在 B 前面对 remaining 列表同理见 worker.py。因此可以安全地把unfinished列表在循环中反复传给ray.wait()。阻塞与异步上下文提醒与ray.get()一样ray.wait()是阻塞调用。源码中特别指出如果在 async 上下文中使用阻塞式ray.wait()会阻塞事件循环并发出警告此时应改用await asyncio.wait(ray_waitables)见 worker.py。若确实在 asyncio 环境下等待需要采用异步替代方案。使用场景与边界适用场景ray.wait()按完成顺序消费结果的写法特别适合以下情况结果需要逐个处理且处理逻辑与任务的提交索引无关任务执行时长差异明显如包含网络 IO、随机数据量、条件分支的负载希望尽快开始处理已就绪的结果缩短端到端耗时需要以流式/增量方式消费一批结果如分批聚合、逐步落盘、增量上报。不适用场景如果结果必须按固定顺序处理例如第 i 个结果依赖第 i-1 个结果则按完成顺序处理没有意义应按索引ray.get()如果所有任务执行时长几乎一致两种写法耗时差异不大按提交顺序处理更简单直观若一次性拿到所有结果即可、无需逐个提前消费更推荐直接ray.get(list_of_refs)它会一次性返回全部结果且保持输入顺序见 worker.py。与其他 ray.get 相关反模式的关联本文所述反模式属于 Ray 官方 Patterns 系列中ray.get()相关反模式之一同系列还包括ray-get-loop在循环中调用ray.get()会破坏并行性。若在同一个循环里既提交任务又ray.get()则上一轮任务未完成就无法提交下一轮最终完全串行化。正确做法是先把所有远程调用提交出去再统一获取结果ray.wait()也可用于此场景中的增量消费。unnecessary-ray-get不必要的ray.get()会导致对象被传输到调用节点产生额外拷贝开销。最佳实践是尽量推迟ray.get()甚至全程直接传递 ObjectRef 给下游任务。ray-get-too-many-objects一次性ray.get()过多对象可能导致堆内存溢出heap out-of-memory或对象存储空间不足object store out-of-space应分批获取处理。该文档也提到可以同时使用ray.wait()按完成顺序消费以降低运行时间。这四个反模式覆盖了ray.get()使用中的四类典型问题使用位置循环内、使用时机过早/不必要、使用方式按提交顺序与使用规模一次性获取过多。它们在 Design Patterns Anti-patterns 索引页 中被统一编排并在 worker.py 中作为ray.get()的关联文档被交叉引用。总结按提交顺序消费独立任务的结果会让先提交的慢任务阻塞整个处理流程导致快任务的结果白白等待整体耗时被掉队任务放大。将ray.get()循环替换为ray.wait(unfinished, num_returns1) 维护未完成集合的写法即可按完成顺序增量处理结果在保持计算语义不变的前提下显著缩短端到端运行时间。这是 Ray 官方 Patterns 系列中最实用、最容易被忽视的性能优化点之一值得在批处理、增量聚合与流式消费场景中优先采用。赞分享人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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点击查看免费下载相关推荐ComfyUI帧插值终极指南轻松实现4倍流畅视频效果ComfyUI帧插值终极指南轻松实现4倍流畅视频效果 你是否曾经观看过卡顿的视频希望画面能更加流畅自然或者作为视频创作者想要将30帧的视频提升到120帧人工智能计算机视觉视频处理Ant Design Select 搜索模式结果排序实战用 filterSort 控制过滤项的展示顺序Ant Design Select 搜索模式结果排序实战用 filterSort 控制过滤项的展示顺序 在 Ant Design Select 中开启搜索前端UI组件设计系统Micro框架GraphQL订阅消息顺序保证FIFO与因果顺序Micro框架GraphQL订阅消息顺序保证FIFO与因果顺序 在实时应用开发中消息传递的顺序一致性直接影响用户体验和系统可靠性。当多个客户端同时订阅Gra后端微服务Web框架创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考