ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Ray 反模式解析:一次 `ray.get()` 拉取过多对象导致 OOM 的成因与批量处理修复方案

Ray 反模式解析:一次 `ray.get()` 拉取过多对象导致 OOM 的成因与批量处理修复方案 人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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 中并行提交大量任务后最直观的做法是把所有ObjectRef收集起来一次性调用ray.get()拿回全部结果。本文基于 Ray 官方 Design Patterns Anti-patterns 文档中的ray-get-too-many-objects一节见 doc/source/ray-core/patterns/ray-get-too-many-objects.rst讲解这种写法为什么会触发堆内存溢出heap out-of-memory或对象存储空间不足object store out-of-space并给出按批拉取 按完成顺序处理的可靠修复方案。读完本文你将掌握ray.get()与ray.wait()的正确配合方式能够写出在面对成千上万个任务结果时依然内存安全的 Ray 应用。一次 ray.get() 拉取过多对象示意图 拉取过多对象会导致调用方节点内存与对象存储同时承压)反模式对全部对象一次性调用 ray.get()核心结论TLDR避免对过多对象调用ray.get()因为它会导致堆内存溢出或对象存储空间不足正确做法是每次只拉取并处理一个批次的结果。当你有大量任务需要并行执行时如果试图一次性ray.get()所有任务的结果Ray 必须把全部对象在同一时刻拉到调用方节点这极可能触发两类资源故障堆内存溢出heap out-of-memoryray.get()返回的是反序列化后的 Python 对象全部结果同时进入调用方进程的堆内存对象存储空间不足object store out-of-space对象在传输到调用方节点时先落入本机对象存储由 Plasma 管理的内存/共享内存海量对象同时落地会撑爆本地对象存储。正确策略是拉取一批、处理一批每处理完一批Ray 就会回收该批次对象占用的空间为后续批次腾出位置。下面是文档配套示例 doc/source/ray-core/doc_code/anti_pattern_ray_get_too_many_objects.py 中的反模式代码import ray import numpy as np ray.init() def process_results(results): # custom process logic pass ray.remote def return_big_object(): return np.zeros(1024 * 10) NUM_TASKS 1000 object_refs [return_big_object.remote() for _ in range(NUM_TASKS)] # This will fail with heap out-of-memory # or object store out-of-space if NUM_TASKS is large enough. results ray.get(object_refs) process_results(results)这段代码同时触犯了本文要讨论的反模式先通过列表推导式一次性派发 1000 个任务这本身没有问题任务可以并行执行随后用ray.get(object_refs)等待全部完成并一次性取回。当NUM_TASKS足够大、每个对象体积足够大时调用方节点将同时保存 1000 个np.zeros(1024 * 10)的对象进程堆和对象存储都面临被撑爆的风险。更好的做法分批拉取配合 ray.wait() 按完成顺序处理同样来自配套示例的改进代码如下BATCH_SIZE 100 while object_refs: # Process results in the finish order instead of the submission order. ready_object_refs, object_refs ray.wait(object_refs, num_returnsBATCH_SIZE) # The node only needs enough space to store # a batch of objects instead of all objects. results ray.get(ready_object_refs) process_results(results)这个循环同时解决了两类问题控制并发内存峰值ray.wait(object_refs, num_returnsBATCH_SIZE)每次只等待 100 个对象就绪ray.get(ready_object_refs)只把这 100 个对象拉回本地。节点只需要为一个批次的对象预留空间而不是为全部对象预留。BATCH_SIZE应根据单对象大小与节点内存预算来设定——对象越大、内存越紧张批次就应该越小。按完成顺序而非提交顺序处理ray.wait()返回的是已就绪的对象引用而不是按提交顺序排列的对象引用。这样先完成的任务结果先被处理不会因为某个慢任务straggler拖累整体进度。这一点在姊妹篇文档 doc/source/ray-core/patterns/ray-get-submission-order.rst 中有专门论述其配套示例 doc/source/ray-core/doc_code/anti_pattern_ray_get_submission_order.py 用 100 个随机耗时任务验证了两种顺序在总耗时上的差异。注意while object_refs循环中object_refs变量被不断替换为ray.wait()返回的未就绪列表循环终止于所有对象被消费完毕整个循环期间无需持有任何已处理对象的引用Ray 的对象存储即可回收对应空间。原理剖析为什么一次 get 会同时压垮堆与对象存储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. 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()是阻塞调用并且强制把对象传输到调用方节点的本地对象存储然后反序列化到调用方进程堆中。当对象引用列表规模巨大时每个对象的反序列化副本同时驻留在进程堆内存中堆内存峰值 所有对象体积之和传输过程中的对象同时在本地对象存储中占一份空间对象存储容量 所有对象体积之和若对象的创建者owner进程已退出或对象被回收还可能触发ObjectLostError——worker.py 中对ray.exceptions.ObjectLostError的处理分支明确提到对象丢失的原因之一就是the object store is full and objects needed to be evicted对象存储已满而必须逐出对象。对象存储侧的错误路径OutOfDisk 与 OUT_OF_DISK_ERROR当拉取对象时本地磁盘/共享内存空间不足错误会沿 C 对象管理器一路向上传播src/ray/object_manager/object_buffer_pool.cc 在创建接收缓冲区失败时检查s.IsOutOfDisk()并直接返回该状态src/ray/object_manager/object_manager.cc 收到分块写入失败时若状态为IsOutOfDisk()则调用pull_manager_-SetOutOfDisk(object_id)src/ray/object_manager/pull_manager.cc 的SetOutOfDisk会通过fail_pull_request_(object_id, rpc::ErrorType::OUT_OF_DISK_ERROR)终止该对象的拉取请求最终在 Python 侧表现为任务失败。一次性ray.get()大量对象正是最容易触发这条OUT_OF_DISK_ERROR路径的场景之一分批拉取则从根源上避免了本地对象存储的并发写入洪峰。ray.wait() 的语义非阻塞轮询式等待ray.wait()同样在 python/ray/_private/worker.py 中实现其签名与关键参数如下参数类型默认值说明ray_waitablesList[ObjectRef]必填待等待的对象引用列表必须唯一num_returnsint1至少就绪多少个对象即返回timeoutfloatNone最多等待秒数None表示一直等到num_returns个对象就绪fetch_localboolTrue是否等待对象下载到本地节点False时只要对象在集群任意节点可用即视为就绪ray.wait()返回二元组(ready_list, remaining_list)第一个列表是已就绪的对象引用第二个列表是仍未就绪的引用且两个列表都保持输入顺序。它本身不传输对象内容除非fetch_localTrue时触发下载因此可以安全地作用于规模庞大的引用列表只有拿到就绪子集后再调ray.get()才真正产生数据搬运。这就是先用ray.wait()圈定小批次、再对小批次ray.get()这一模式在底层不爆内存的原因。与其它 ray.get() 反模式的关联与边界一次性ray.get()过多对象只是ray.get()系列反模式之一doc/source/ray-core/patterns/index.rst 收录了完整的 Design Patterns Anti-patterns 目录其中与ray.get()强相关的还包括ray-get-loop在循环里调用 ray.get 破坏并行ray.get()是阻塞调用若在循环中先get再派发下一个任务就退化为完全串行执行。正确做法是先派发全部任务再统一等待结果。这与本文反模式正好互补——一个是等待太早一个是等待太多。unnecessary-ray-get不必要的 ray.get 拖慢性能只要不直接操纵对象就不要ray.get()它直接传递ObjectRef让下游任务在目标节点上隐式解引用可以省去把对象搬运到 driver 再搬运回去的额外拷贝。这提醒我们本文的批处理方案用于必须由当前进程消费结果的场景若结果只是中间产物应优先考虑不落地、直接传引用。nested-ray-get任务内部对参数 get 使执行串行化任务参数应直接传ObjectRef让 Ray 自动解析避免在任务内部阻塞等待。pipelining流水线化避免空转等待在生产者-消费者模式下用边产出边消费的流水线替代全部产出后再一次性消费与本文的批量拉取理念同源。如果任务结果需要在批次间累积例如求和、聚合可参考 anti_pattern_ray_get_submission_order.py 中按完成顺序累加、最终断言结果一致的写法确保按完成顺序处理不会改变聚合语义。实践要点小结永远不要一次性ray.get()大量对象。把ray.get()的目标规模控制在一个可承受的批次内节点只需为单个批次预留堆与对象存储空间。用ray.wait()划分批次并按完成顺序消费既能限制内存峰值又能规避慢任务拖累整体耗时批内对象处理完毕即释放引用让 Ray 及时回收空间。根据对象体积与节点规格选择BATCH_SIZE大对象场景如机器学习推理产出的超大数组应显著调小批次需要控制单次等待时间时可配合timeout参数。区分必须本地消费与可以传引用只有当前进程真正需要操作结果时才ray.get()中间产物优先直接传递ObjectRef见 unnecessary-ray-get。这套分批拉取 完成顺序消费的模式可以直接落地到任意 Ray 任务集群中是规避对象存储OUT_OF_DISK_ERROR与堆 OOM 的最实用手段。赞分享人工智能分布式训练强化学习任务调度模型推理服务【免费下载链接】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 对象系统实战指南ObjectRef、ray.get 与分布式对象传递Ray 对象系统实战指南ObjectRef、ray.get 与分布式对象传递 本篇指南以 Ray 官方文档《Objects》为核心系统讲解 Ray 分布式对人工智能分布式训练强化学习任务调度模型推理服务网络受限怎么办cloud-uploader 代理配置完整教程手把手带你解锁上传网络受限怎么办cloud uploader 代理配置完整教程手把手带你解锁上传 cloud uploader 是一款专为 Mac 用户打造的网易云音乐云盘上人工智能分布式训练强化学习任务调度模型推理服务RuboCop批量处理一次性修复整个代码库RuboCop批量处理一次性修复整个代码库 引言代码修复的痛点与解决方案 你是否曾面对一个拥有数千个Ruby文件的代码库其中充斥着数百个代码风格违规和潜在代码质量Lint格式化静态分析开发工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表