Ray对象存储与数据流:快速理解分布式对象的内存管理精髓
【免费下载链接】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 对象存储(Object Store)是 Ray 分布式计算引擎的"共享数据中枢"。它让每个集群节点拥有一块共享内存区域,任务与 Actor 产生的分布式对象可以被高速读写、零拷贝共享,并通过自动溢写(Spilling)与分布式引用计数实现内存的精巧管理。本文用一篇图解指南,帮你快速理解 Ray 数据流的运作方式与分布式对象的内存管理精髓。
什么是 Ray 对象存储(Object Store)?🧊
在 Ray 中,任务和 Actor 都会创建和消费数据对象,这些对象被称为远程对象(Remote Object)——它们可以存放在集群中的任意节点上,并通过 **对象引用(ObjectRef,类似指针/唯一 ID)**来间接访问。
每个节点运行一个对象存储进程,基于共享内存实现,默认占用节点可用内存的 30%(可通过启动参数调整)。它有几个关键特性:
- 不可变性:对象创建后不可修改,因此可以安全地在多节点复制而无需同步;
- 零拷贝读取:对 numpy 数组等数据类型,
ray.get()直接返回指向共享内存的视图,无需序列化拷贝; - 分布式引用计数:只要集群中还有任意
ObjectRef引用某个对象,它就"钉"在内存中;所有引用消失后自动释放。
如上图所示,在 Ray Data 这类数据流场景中,数据块在各阶段(下载、GPU 推理、写回)之间流经对象存储的共享内存,避免落盘和重复拷贝,这正是"对象存储 + 数据流"组合的威力。
相关概念详见 doc/source/ray-core/objects.rst 与 doc/source/ray-core/scheduling/memory-management.rst。
对象如何在集群中流动:put、get 与传参 🌊
对象的"生产"与"消费"只有两个入口:
ray.put():手动把对象放入对象存储,返回一个 ObjectRef;- 远程调用返回值:
task.remote()完成后,结果自动写入对象存储并返回 ObjectRef。
取回数据时用ray.get()——如果当前节点没有该对象,Ray 会自动从持有它的节点下载。
值传递 vs 引用传递
Ray 对传参方式有一个巧妙的约定:
| 传参方式 | 行为 | 适用场景 |
|---|---|---|
| 对象作为顶层参数传入任务 | Ray 先解引用,等待数据就绪才执行任务(按值) | 任务确实要用数据 |
| 对象嵌在列表/字典中传入 | 不解引用,只传"指针"(按引用) | 只转发引用、不消费数据 |
这意味着:如果你只是想把数据"转交"给下游任务而不查看它,用嵌套引用传递可以让数据完全不需要移动到当前机器上,既省内存又省网络带宽。
分布式引用计数:自动内存回收的精髓 🔢
Ray 实现了分布式引用计数,任何处于"作用域内"的ObjectRef都会把对象钉在对象存储中。共有 5 类引用来源(可通过ray memory命令逐一看到):
- LOCAL_REFERENCE:本地 Python 变量持有的引用;
- PINNED_IN_MEMORY:
ray.get()反序列化后,数据直接指向共享内存,变量未删除前对象无法驱逐; - USED_BY_PENDING_TASK:有未完成任务依赖该对象;
- CAPTURED_IN_OBJECT:引用被序列化进了另一个对象内部;
- Actor Handles:Actor 持有引用。
💡 新手最常见的内存坑:闭包捕获(closure capture)会把对象钉住直到作业结束;以及
ray.get()后忘记删除变量,导致"已用完"的数据仍占着对象存储。
详细的 5 类引用示例与ray memory调试方法见 doc/source/ray-core/scheduling/memory-management.rst。
内存不够了怎么办:对象溢写(Spilling)💾
当对象存储被填满时,Ray 不会让你等到任务失败——它会触发对象溢写:
- 自动把部分"已钉住"的对象从共享内存溢写到本地磁盘(默认临时目录)或 S3;
- 当该对象再次被
ray.get()访问时,透明地恢复回内存; - 整个过程对应用代码完全无感知,相当于用磁盘"扩容"了对象存储的容量(代价是 I/O 延迟)。
其内部分为三层,保证核心执行路径不被拖慢:
- 检测层(Plasma 存储线程):内存分配失败时立即回调;
- 编排层(Raylet 主线程):基于 LRU 与钉住状态决定"溢写谁",并批量执行;
- 执行层(独立 Python IO Worker 进程):通过 gRPC 完成磁盘/S3 读写,I/O 再慢也不阻塞调度。
架构与状态机图解见 doc/source/ray-core/internals/object-spilling.rst,用户级配置(如自定义溢写目录--object-spilling-directory)见 doc/source/ray-core/objects/object-spilling.rst。
上图是 Ray Dashboard 的节点硬件利用率面板,右下角的Node Memory (heap + object store)曲线能直接看到堆内存与对象存储内存的叠加变化——排查内存问题时第一眼就该看它。
实战排查:用 ray memory 定位内存大头 🔍
当遇到ObjectStoreFullError(对象存储满)或内存持续增长时,在应用运行期间执行ray memory,它会输出:
- 每个节点的对象内存用量、引用数量汇总;
- 每个
ObjectRef的持有者(Driver/Worker)、引用类型、对象大小、创建它的代码行; - 全集群的溢写统计(已溢写/已恢复数据量与吞吐)。
配合排序参数(如按对象大小排序、按堆栈分组),可以快速锁定哪一行代码造成了内存泄漏。
关键路径速查清单 📌
| 你想了解 | 去哪里看 |
|---|---|
| 远程对象、ObjectRef、序列化 | doc/source/ray-core/objects.rst、doc/source/ray-core/objects/serialization.rst |
内存模型与ray memory调试 | doc/source/ray-core/scheduling/memory-management.rst |
| 对象溢写原理(含架构图) | doc/source/ray-core/internals/object-spilling.rst |
| 溢写目录配置与统计 | doc/source/ray-core/objects/object-spilling.rst |
| 数据流式处理的内部机制 | doc/source/data/data-internals.rst |
| 大对象传参的反模式与优化 | doc/source/ray-core/patterns/pass-large-arg-by-value.rst |
小结
Ray 对象存储的核心思想可以浓缩为三句话:共享内存实现零拷贝数据流、分布式引用计数实现自动内存回收、对象溢写实现"内存+磁盘"的弹性扩容。理解这三点,你就掌握了 Ray 分布式对象内存管理的精髓,也能在遇到ObjectStoreFullError时从容排查 🚀
【免费下载链接】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),仅供参考