Ray 反模式剖析:带外序列化 ray.ObjectRef 引发的对象泄漏与提前回收问题
2026/9/19 12:00:47 网站建设 项目流程

Ray 反模式剖析:带外序列化 ray.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 中,ray.ObjectRef采用分布式引用计数来管理对象生命周期:Ray 会持续"钉住"(pin)底层对象,直到系统不再持有任何引用。本文聚焦 Ray Core 官方文档中的经典反模式——带外序列化(out-of-band serialization)ray.ObjectRef——分析其背后的分布式引用计数原理、它为何会导致对象被提前回收或内存泄漏,并给出可落地的检测手段(RAY_allow_out_of_band_object_ref_serialization=0环境变量)与正确的替代写法。读完本文,你将能识别并修复代码中所有绕过 Ray 引用追踪的 ObjectRef 序列化点,避免诡异的任务挂起(hang)与磁盘溢出(spilling)问题。

一、背景:Ray 的分布式引用计数与对象"钉住"机制

Ray 的任务、Actor 与对象存储(object store)共同构成一个分布式系统,ray.ObjectRef是该系统中的分布式引用计数句柄。其核心生命周期规则是:

  • Ray 会钉住ObjectRef指向的底层对象,直到系统中所有对该引用的引用都不再被使用;
  • 当指向该对象的全部引用消失后,Ray 对其进行垃圾回收(GC),把对象从系统中清理掉。

从源码看,这种"钉住"能力由 Core Worker 提供的add_object_ref_reference实现,即向本地 Core Worker 注册一个引用计数,防止对象被驱逐。引用计数机制的底层实现在 serialization.py 中体现为"带内(in-band)"与"带外(out-of-band)"两种序列化路径的分流:

  • 带内序列化(in-band)ObjectRef作为任务参数、任务返回值,或作为另一个对象的一部分被序列化时,Ray 能感知并追踪它。源码注释明确写道:"此 ObjectRef 正被存储在一个对象中,将其 ID 加入该对象包含的 ID 列表,以便只要外层对象在作用域内,内层对象的值就保持存活"(serialization.py)。
  • 带外序列化(out-of-band):用户在业务代码里直接调用pickle.dumps(obj_ref)ray.cloudpickle.dumps(obj_ref),或让ObjectRef被远程函数闭包捕获。此时引用逃逸出 Ray 的追踪范围,Ray 无法获知该引用的存活状态。

官方文档明确指出该反模式的核心危害(out-of-band-object-ref-serialization.rst):

TLDR:避免序列化ray.ObjectRef,因为 Ray 无法判断何时对底层对象进行垃圾回收。

二、为什么序列化 ObjectRef 是反模式:两条灾难路径

当用户代码对ray.ObjectRef做带外序列化时,引用计数被绕过,可能走向两种截然不同的灾难:

1. 使用普通pickle:对象被提前回收,任务意外挂起

用标准库pickle.dumps(obj_ref)序列化时,Ray 完全不知道这个引用的存在,也不会因此增加任何引用计数。序列化得到的字节串即使被传回、反序列化并尝试ray.get(),Ray 也可能早已把底层对象当成"无引用"垃圾回收掉。最终表现为:ray.get()无限等待或抛出GetTimeoutError,任务莫名挂起,且极难排查。

2. 使用ray.cloudpickle:对象被钉住整个 Worker 生命周期,引发泄漏与磁盘溢出

与普通pickle不同,ray.cloudpickle是 Ray 自带的序列化器,能够识别ObjectRef并主动为其注册引用。但为了"安全",它采取了一种代价高昂的策略:只要检测到带外序列化,就把对象钉住,直到对应的 owner worker 死亡。源码 serialization.py 的注释写得很直白:

"如果这是带外序列化(例如直接调用 cloudpickle,或被远程函数/Actor 捕获),那么通过添加一个永远不会被移除的本地引用,将对象钉住整个 Worker 的生命周期。"

所谓"钉住",即该对象永远无法从对象存储中被驱逐(evict)。一旦此类代码被高频执行(例如在循环里反复cloudpickle.dumps引用),钉住的对象会不断累积,形成 Ray 对象泄漏,最终挤占对象存储内存并触发磁盘 spill,拖垮集群性能

三、反模式代码示例与运行验证

官方文档配套的可运行示例位于 anti_pattern_out_of_band_object_ref_serialization.py,它一次性演示了上述两条灾难路径与检测手段,完整代码如下:

import ray import pickle from ray._private.internal_api import memory_summary import ray.exceptions ray.init() @ray.remote def out_of_band_serialization_pickle(): obj_ref = ray.put(1) import pickle # object_ref is serialized from user code using a regular pickle. # Ray can't keep track of the reference, so the underlying object # can be GC'ed unexpectedly, which can cause unexpected hangs. return pickle.dumps(obj_ref) @ray.remote def out_of_band_serialization_ray_cloudpickle(): obj_ref = ray.put(1) from ray import cloudpickle # ray.cloudpickle can serialize only when # RAY_allow_out_of_band_object_ref_serialization=1 env var is set. # However, the object_ref is pinned for the lifetime of the worker, # which can cause Ray object leaks that can cause spilling. return cloudpickle.dumps(obj_ref) print("==== serialize object ref with pickle ====") result = ray.get(out_of_band_serialization_pickle.remote()) try: ray.get(pickle.loads(result), timeout=5) except ray.exceptions.GetTimeoutError: print("Underlying object is unexpectedly GC'ed!\n\n") print("==== serialize object ref with ray.cloudpickle ====") # By default, it's allowed to serialize ray.ObjectRef using # ray.cloudpickle. ray.get(out_of_band_serialization_ray_cloudpickle.options().remote()) # you can see objects are still pinned although it's GC'ed and not used anymore. print(memory_summary()) print( "==== serialize object ref with ray.cloudpickle with env var " "RAY_allow_out_of_band_object_ref_serialization=0 for debugging ====" ) try: ray.get( out_of_band_serialization_ray_cloudpickle.options( runtime_env={ "env_vars": { "RAY_allow_out_of_band_object_ref_serialization": "0", } } ).remote() ) except Exception as e: print(f"Exception raised from out_of_band_serialization_ray_cloudpickle {e}\n\n")

运行结果解读

依次运行上述三段验证,可以看到:

  1. pickle.dumps(obj_ref)路径:反序列化后调用ray.get(pickle.loads(result), timeout=5)抛出ray.exceptions.GetTimeoutError,并打印Underlying object is unexpectedly GC'ed!——底层对象已经被提前回收,这正是导致任务意外挂起的根因。
  2. ray.cloudpickle.dumps(obj_ref)路径:默认(环境变量未设置时)序列化被允许,但调用ray._private.internal_api.memory_summary()打印对象存储内存摘要时,可以看到该对象虽然逻辑上已不再被使用,却仍处于 pinned 状态,即"钉住"在对象存储中,构成泄漏。
  3. 开启检测路径:通过runtime_envenv_varsRAY_allow_out_of_band_object_ref_serialization设为"0"后,序列化时立即抛出异常(详见下一节)。

四、检测手段:RAY_allow_out_of_band_object_ref_serialization=0

为了帮助开发者发现代码中是否存在该反模式,Ray 提供了开关环境变量:

RAY_allow_out_of_band_object_ref_serialization=0

当该变量为0(关闭)时,只要ray.cloudpickle检测到带外序列化ray.ObjectRef,就会抛出ray.exceptions.OufOfBandObjectRefSerializationException异常(异常类定义于 exceptions.py),并附带极具诊断价值的错误消息。

异常消息中的关键信息

从 serialization.py 的抛错代码可以看出,该异常消息包含:

  • 被序列化的ObjectRef的十六进制 ID(object_ref.hex());
  • 明确的修复指引:若确实需要放行,可设置RAY_allow_out_of_band_object_ref_serialization=1,但同时警告"对象将被钉住整个 Worker 生命周期并可能导致 Ray 对象泄漏";
  • 调用点(callsite)定位:异常会打印调用点,但前提是开启RAY_record_ref_creation_sites=1来记录引用创建位置,否则会提示 "Disabled. Set RAY_record_ref_creation_sites=1"。因此排障时建议同时设置:
RAY_allow_out_of_band_object_ref_serialization=0 RAY_record_ref_creation_sites=1

重要使用前提:环境变量必须提前设置

有两个细节需要特别注意:

  1. 默认值为开启:该环境变量在 serialization.py 中通过ray_constants.env_bool("RAY_allow_out_of_band_object_ref_serialization", True)解析,默认值是True,即默认允许带外序列化——这正是该反模式容易悄悄潜入代码的原因。
  2. 不可动态修改:该变量在模块导入时(ray.init()之前)就被读取并缓存,无法在运行时动态改变。测试用例 test_serialization.py 的注释专门指出这一点:"Use ray.remote as a workaround because RAY_allow_out_of_band_object_ref_serialization cannot be set dynamically"。因此:
    • 若想全局开启检测,需在启动 Ray 进程前设置环境变量,例如RAY_allow_out_of_band_object_ref_serialization=0 ray start --head或启动 Python 前 export;
    • 若只想让某个远程任务受控,可以采用示例代码中的方式,通过runtime_envenv_vars注入(这会为对应任务创建独立 Worker 环境)。

测试用例佐证

仓库中的单元测试完整验证了开关两端的语义(test_serialization.py):

  • test_cannot_out_of_band_serialize_object_ref:在RAY_allow_out_of_band_object_ref_serialization=0下,无论是把ref捕获进另一个远程函数闭包,还是直接cloudpickle.dumps(ray.put(1)),都会抛出OufOfBandObjectRefSerializationException
  • test_can_out_of_band_serialize_object_ref_with_env_var:在=1下,同样的操作可以正常完成。

这说明闭包捕获显式cloudpickle.dumps是触发该反模式的两大典型来源,检测开关对两者均有效。

五、正确做法:让 ObjectRef 始终留在 Ray 的引用追踪范围内

修复该反模式的原则只有一条:不要让ObjectRef以字节形式逃逸出 Ray 的追踪体系。具体建议如下:

1. 需要跨进程传递引用时,交给 Ray 的"带内"通道

  • 作为任务/ Actor 方法的参数或返回值传递:Ray 内部会自动做带内序列化并维护引用计数,这是最常用的正确姿势;
  • ObjectRef放入另一个对象再ray.put:如第一节所述,Ray 会把内层引用登记到外层对象的依赖列表中,保证外层对象存活期间内层值不被回收;
  • 在远程函数内部解引用:直接在任务体里对ObjectRef调用ray.get()取回值,再把"值"传出去,而不是把"引用"传出去。

2. 禁止手动pickle/cloudpickle序列化引用

不要为了缓存、落盘或发送到外部系统而把ray.ObjectRefpickle.dumps/cloudpickle.dumps序列化。即便序列化的是"值"而非引用,也建议确认对象图中不嵌套任何ObjectRef

3. 如果确实需要跨集群/外部系统传递

ray.get(obj_ref)取出真实值(或把大对象落盘/上传到外部存储),在外部系统间传递值本身(或值的稳定标识),到达目的端后再重新ray.put生成新的ObjectRef。切勿直接搬运引用字节。

4. 定期用memory_summary体检对象存储

ray._private.internal_api.memory_summary()(示例代码中已使用)可打印当前对象存储中各对象的大小、引用状态(PINNED_IN_MEMORY等)与 owner 信息。若发现大量本应被回收却长期 pinned 的对象,往往是带外序列化泄漏的信号,应结合RAY_record_ref_creation_sites=1回看创建调用点。

六、小结:反模式自查清单

检查项说明
代码中是否存在pickle.dumps(obj_ref)会导致底层对象被提前 GC,引发GetTimeoutError与任务挂起
代码中是否存在ray.cloudpickle.dumps(obj_ref)默认被允许,但对象会被钉住整个 Worker 生命周期,造成泄漏与磁盘 spill
远程函数闭包是否意外捕获了外层ObjectRef同样属于带外序列化,需显式避免或改造
是否设置了检测开关调试阶段设置RAY_allow_out_of_band_object_ref_serialization=0并配合RAY_record_ref_creation_sites=1
对象存储中是否有异常 pinned 对象memory_summary()体检,定位泄漏源

该反模式是 Ray Core Patterns 系列文档中的一篇(完整清单见 patterns/index.rst),与之相关的还有 ray-get-loop.rst、nested-ray-get.rst、return-ray-put.rst 等对象引用与任务编排反模式,它们共同构成一套系统的 Ray 最佳实践。记住最核心的一句话:ray.ObjectRef是 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

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询