FEATURED · 精选文章

Ray 序列化完全指南:Plasma 对象存储、零拷贝传输与自定义序列化实践

发布时间 / 2026/9/20 1:30:00
来源 / 创域科博编辑部
栏目 / 资讯中心
Ray 序列化完全指南:Plasma 对象存储、零拷贝传输与自定义序列化实践 Ray 序列化完全指南Plasma 对象存储、零拷贝传输与自定义序列化实践【免费下载链接】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 分布式运行时中的数据序列化机制从 Plasma 共享内存对象存储与 Pickle 协议 5 的技术选型到 NumPy 数组零拷贝读取、assignment destination is read-only 的经典坑再到可选的 PyTorch 只读张量零拷贝序列化以及三类自定义序列化方案与inspect_serializability故障排查工具。读完本文你将理解 Ray 在跨进程、跨节点传输对象时的底层原理并能针对不可序列化对象、性能瓶颈和只读数组报错给出可直接落地的解决方案。Ray 序列化概述由于 Ray 的多个进程并不共享内存空间worker 之间、节点之间传输数据时对象必须经历**序列化serialization与反序列化deserialization**两个阶段。Ray 借助 Plasma 对象存储高效地在不同进程、不同节点之间搬运对象同节点上的 worker 之间共享对象存储中的 NumPy 数组时走的是零拷贝反序列化路径。从技术选型上看Ray 采用定制化的Pickle 协议版本 5PEP 574回移植实现来替代早期的 PyArrow 序列化器消除了此前的一系列限制例如无法序列化递归对象。当前 Ray 完全兼容 Pickle 协议 5同时借助 cloudpickle 支持更广泛的对象类型例如lambda 与嵌套函数、动态类等原生 pickle 难以处理的对象。协议版本的影响值得留意Python 大多数发行版默认使用的 pickle 协议是协议 3而协议 4 与协议 5 对体积较大的对象更高效。Ray 内部统一以协议 5 进行序列化见 python/ray/_private/serialization.py 中_serialize_to_pickle5的pickle.dumps(value, protocol5, buffer_callback...)因此大对象跨 worker 传输的性能优于默认 pickle。Plasma 对象存储Plasma 是一个内存对象存储最初作为 Apache Arrow 的一部分开发。在 Ray 1.0.0 版本之前Ray 将 Arrow 的 Plasma 代码 fork 进自身代码库以便贴合 Ray 的架构与性能需求进行解耦和持续演进。Plasma 的核心设计约束有两点对象不可变immutable所有放入 Plasma 对象存储的对象都不可修改并被保存在**共享内存shared memory**中。这样设计是为了让同一节点上的多个 worker 都能高效访问同一份数据而无需各自拷贝。数据默认不广播每个节点拥有自己独立的对象存储。数据写入对象存储后并不会自动广播到其他节点数据会一直留在写入方节点本地直到另一个节点上的任务或 Actor 显式请求该对象即触发跨节点拉取。因此Plasma 承担的是进程间零拷贝共享 按需跨节点传输的职责同节点多 worker 靠共享内存省去拷贝跨节点则按需传输。NumPy 数组零拷贝读取与只读陷阱Ray 针对 NumPy 数组做了专门优化使用Pickle 协议 5 的 out-of-band 数据带外缓冲传输。具体而言NumPy 数组在对象存储中以只读对象形式存放同一节点上的所有 Ray worker 可以直接读取对象存储中的数组无需复制zero-copy readsworker 进程中的每个 NumPy 数组对象只是持有一个指向共享内存中对应数组的指针。正因如此任何针对只读对象的写入操作都必须先把数组复制到本地进程内存。如果忽略这一点把ray.put得到、或作为参数传入远程函数的 NumPy 数组直接原地修改就会触发如下报错import ray import numpy as np ray.remote def f(arr): arr[0][0] 1 # 报错: assignment destination is read-only ray.init() a np.zeros((10, 10)) ray.get(f.remote(a))运行时会抛出ValueError: assignment destination is read-only。解决方式是在目标端手动复制一份再修改ray.remote def f(arr): arr arr.copy() # 复制到本地进程内存后即可修改 arr[0][0] 1 return arr需要说明的是arr arr.copy()等价于关闭了 Ray 提供的零拷贝反序列化能力——你以一次显式复制换取了可写性这在需要就地更新数据的场景中是必要代价。序列化注意事项以下要点是 Ray 官方序列化指南反复强调的实践原则也均有对应源码与测试支撑协议版本Ray 使用 Pickle 协议 5。大多数 Python 发行版默认的 pickle 协议是 3对于大对象协议 4 与 5 比协议 3 更高效。Ray 在_serialize_to_pickle5中固定以protocol5结合buffer_callback完成序列化见 python/ray/_private/serialization.py。引用共享语义得到保留对于非原生类型对象即使同一对象在目标对象中被多次引用Ray 也始终只保留一份拷贝。例如import ray import numpy as np obj [np.zeros(42)] * 99 l ray.get(ray.put(obj)) assert l[0] is l[1] # 没问题引用关系被保留这得益于 pickle5/cloudpickle 对对象图引用结构的忠实还原对依赖共享可变状态的业务逻辑尤为重要。性能建议只要条件允许尽量使用 NumPy 数组或 NumPy 数组的 Python 集合list/dict 等以获得最大传输性能。锁对象基本不可序列化复制一把锁没有意义还可能引发严重的并发问题。如果你的对象包含锁必须另寻方案——例如把锁改为不参与序列化的字段或在自定义序列化器中重建锁。通用规避思路只使用原生类型如 NumPy 数组、由 NumPy 数组与基本类型构成的 list/dict或者用 Actor 持有无法序列化的对象Actor 的对象不随每次调用序列化传输而是常驻在所属 worker 中往往能直接绕开序列化问题。零拷贝序列化只读 PyTorch 张量除了 NumPyRay 还为只读 PyTorch 张量提供了可选的零拷贝序列化能力默认关闭。工作原理与适用条件启用后Ray 在序列化torch.Tensor时不再走默认的整块复制路径而是将张量转换为 NumPy 数组uint8视图再借助 pickle5 的零拷贝缓冲共享完成传输从而避免复制底层张量数据详见 python/ray/_private/tensor_serialization_utils.py 中的zero_copy_tensors_reducer与_zero_copy_tensors_deserializer实现序列化时执行detach()→.cpu()→contiguous()→ 扁平化为 1D →view(torch.uint8)→ 转 NumPy反序列化时用torch.from_numpyview还原 dtype reshape还原形状 to(device)还原设备。这一特性在跨任务、跨 Actor 传递大张量时可显著降低延迟但必须谨慎使用PyTorch 原生不支持只读张量。启用后 Ray 不会拷贝数据也允许对共享内存的写入——这意味着同节点上两个进程若共置运行一个进程在ray.get()之后修改张量可能被另一个进程看到造成数据污染。该特性在以下条件下效果最佳张量requires_grad False即已从 autograd 计算图分离张量内存连续tensor.is_contiguous()为 True张量位于 CPU 内存时收益更大GPU 张量或非连续张量虽仍受支持但 Ray 会先自动搬运到 CPU / 先做连续化产生初始拷贝见 python/ray/_private/ray_constants.py 的注释说明你未使用 Ray Direct Transport直接传输路径。从源码看若张量不满足CPU、已分离、连续三个条件reducer 会以一次或两次完整拷贝为代价补足条件并发出ZeroCopyTensorsWarning警告见 python/ray/_private/tensor_serialization_utils.py此时零拷贝收益会被抵消。如何启用该特性默认关闭RAY_ENABLE_ZERO_COPY_TORCH_TENSORS环境变量的默认值为False定义于 python/ray/_private/ray_constants.py。你需要在运行脚本之前于驱动进程中设置环境变量export RAY_ENABLE_ZERO_COPY_TORCH_TENSORS1也可以像下面示例一样在ray.init时通过runtime_env注入环境变量从而对集群内所有进程driver 与 worker生效。端到端示例求 1GiB 张量的和下面示例利用ray.get()计算一个 1GiB 张量的总和全程走零拷贝序列化import ray import torch import time ray.init(runtime_env{env_vars: {RAY_ENABLE_ZERO_COPY_TORCH_TENSORS: 1}}) ray.remote def process(tensor): return tensor.sum() x torch.ones(1024, 1024, 256) start_time time.perf_counter() result ray.get(process.remote(x)) elapsed_time time.perf_counter() - start_time print(fElapsed time: {elapsed_time}s) assert result x.sum()官方文档记录的实测对比同一脚本在启用开关前后的端到端耗时表明启用零拷贝序列化可将端到端延迟降低约66.3%# 未启用零拷贝序列化 Elapsed time: 23.53883756196592s # 启用零拷贝序列化 Elapsed time: 7.933729998010676s注意上表数据为文档所记录环境下的单次实测参考实际收益取决于张量尺寸、设备与集群拓扑。仓库中对应的回归测试覆盖了标量、高维张量、非连续stride张量、以及列表/字典/NamedTuple/dataclass 嵌套张量等丰富场景见 python/ray/tests/test_cpu_tensor_zero_copy_serialization.py 与 python/ray/tests/test_gpu_tensor_zero_copy_serialization.py测试夹具ray_start_cluster_with_zero_copy_tensors通过monkeypatch设置环境变量并注入runtime_env见 python/ray/tests/conftest.py。自定义序列化当 Ray 默认的序列化器pickle5 cloudpickle不满足需求时——例如某些对象序列化失败、或默认序列化对特定对象太慢——你有至少 3 种自定义方案。文档与源码python/ray/util/serialization.py共同构成了完整的自定义体系。方式一在类内部定义__reduce__如果你能访问类的源码可以在类内定义__reduce__方法这是大多数 Python 库普遍采用的做法。__reduce__返回一个(callable, args)二元组反序列化时 Ray 会调用callable(*args)重建对象import ray import sqlite3 class DBConnection: def __init__(self, path): self.path path self.conn sqlite3.connect(path) # 没有 __reduce__ 时该实例无法序列化 def __reduce__(self): deserializer DBConnection serialized_data (self.path,) return deserializer, serialized_data original DBConnection(/tmp/db) print(original.conn) copied ray.get(ray.put(original)) print(copied.conn)输出两个连接对象在反序列化后均被成功重建sqlite3.Connection object at ... sqlite3.Connection object at ...这里的技巧是连接对象本身不可 pickle但重建它所需的**构造参数路径**是可序列化的因此序列化构造参数而非连接对象即可绕过限制。方式二通过ray.util.register_serializer注册类级序列化器当你无法访问或修改类源码时可以为该类型注册自定义的序列化/反序列化函数。以包含threading.Lock的类为例——锁不可序列化默认直接ray.put必然失败import ray import threading class A: def __init__(self, x): self.x x self.lock threading.Lock() # 无法被序列化 try: ray.get(ray.put(A(1))) # 失败 except TypeError: pass def custom_serializer(a): return a.x def custom_deserializer(b): return A(b) # 为类 A 注册序列化器与反序列化器 ray.util.register_serializer( A, serializercustom_serializer, deserializercustom_deserializer) ray.get(ray.put(A(1))) # 成功 # 可以随时注销该序列化器 ray.util.deregister_serializer(A) try: ray.get(ray.put(A(1))) # 失败 except TypeError: pass # 注销一个不存在的序列化器不会产生任何影响 ray.util.deregister_serializer(A)从 python/ray/util/serialization.py 的源码可以确认register_serializer与deregister_serializer均为 PublicAPI内部实现为把(custom_deserializer, (custom_serializer(obj),))构造成一个 cloudpickle reducer 并写入pickle.CloudPickler.dispatch[cls]反序列化时即调用custom_deserializer(custom_serializer(obj))注销时从 dispatch 表中弹出对应条目。仓库中的test_custom_serializer测试python/ray/tests/test_serialization.py完整复现了注册→成功→注销→失败→重复注销无副作用的流程并验证了在ray.init()之前注册、之后依然生效的行为。使用该 API 时需牢记三点序列化器按 worker 本地管理每个 Ray worker 都有独立的序列化上下文因此每个 worker 内都需要注册注销同样只对当前 worker 生效。注册即替换为同一类注册新序列化器会立即替换旧序列化器。幂等性重复注册相同序列化器不会产生副作用。方式三为特定对象定制序列化辅助类如果只想定制某一个对象实例而非整个类型的序列化行为可以包装一个辅助类并定义其__reduce__import ray import threading class A: def __init__(self, x): self.x x self.lock threading.Lock() # 无法序列化 try: ray.get(ray.put(A(1))) # 失败 except TypeError: pass class SerializationHelperForA: 用于序列化的辅助类。 def __init__(self, a): self.a a def __reduce__(self): return A, (self.a.x,) ray.get(ray.put(SerializationHelperForA(A(1)))) # 成功 # 该序列化器只对特定对象生效而非所有 A 实例 try: ray.get(ray.put(A(1))) # 仍然失败 except TypeError: pass自定义异常序列化器默认 pickle 机制无法序列化的异常同样可以借助ray.util.register_serializer处理——注意必须在 driver 和所有 worker 中都注册。下面示例演示一个携带不可序列化锁字段的自定义异常如何跨 worker 边界正常传播import ray import threading class CustomError(Exception): def __init__(self, message, data): self.message message self.data data self.lock threading.Lock() # 无法被序列化 def custom_serializer(exc): return {message: exc.message, data: str(exc.data)} def custom_deserializer(state): return CustomError(state[message], state[data]) # 在 driver 中注册 ray.util.register_serializer( CustomError, serializercustom_serializer, deserializercustom_deserializer ) ray.remote def task_that_registers_serializer_and_raises(): # 在 worker 中注册自定义序列化器 ray.util.register_serializer( CustomError, serializercustom_serializer, deserializercustom_deserializer ) # 抛出自定义异常 raise CustomError(Something went wrong, {complex: data}) # 自定义异常将跨 worker 边界被正确序列化 try: ray.get(task_that_registers_serializer_and_raises.remote()) except ray.exceptions.RayTaskError as e: print(fCaught exception: {e.cause}) # 这里就是 CustomError自定义异常在远程任务中被抛出时Ray 的处理流程是使用你注册的自定义序列化器序列化该异常将其包装进ray.exceptions.RayTaskError反序列化后的原始异常可通过ray_task_error.cause属性访问如e.cause。此外当序列化本身失败时Ray 会抛出ray.exceptions.UnserializableException其中包含原始堆栈的字符串表示便于定位问题。故障排查定位不可序列化对象使用inspect_serializabilityray.util.inspect_serializability是定位棘手 pickle 问题的首选工具。它能够追踪任意 Python 对象函数、类或对象实例内部潜在的不可序列化部分返回(是否可序列化, 不可序列化对象集合)其签名与返回类型见 python/ray/util/check_serialize.py。以下面这个闭包引用了threading.Lock的函数为例from ray.util import inspect_serializability import threading lock threading.Lock() def test(): print(lock) inspect_serializability(test, nametest)输出 Checking Serializability of function test at 0x7ff130697e50 !!! FAIL serialization: cannot pickle _thread.lock object Detected 1 global variables. Checking serializability... Serializing lock unlocked _thread.lock object at 0x7ff1306a9f30... !!! FAIL serialization: cannot pickle _thread.lock object WARNING: Did not find non-serializable object in unlocked _thread.lock object at 0x7ff1306a9f30. This may be an oversight. Variable: FailTuple(lock [objunlocked _thread.lock object at 0x7ff1306a9f30, parentfunction test at 0x7ff130697e50]) was found to be non-serializable. There may be multiple other undetected variables that were non-serializable. Consider either removing the instantiation/imports of these variables or moving the instantiation into the scope of the function/class. 输出会明确指出哪个全局变量lock、挂在哪个父对象函数test之下、为何失败cannot pickle _thread.lock object并给出建议——移除这些变量的实例化/导入或将其实例化移入函数/类的作用域内部。仓库中对应的警告可操作性测试见 python/ray/tests/test_serialization.py 的test_inspect_serializability_warning_message_is_actionable。开启 verbose 调试如果inspect_serializability仍不足以定位问题可以在导入 Ray 之前设置环境变量export RAY_PICKLE_VERBOSE_DEBUG2这会切换为基于 Python 的序列化后端而非 C-Pickle从而允许你在序列化中途用 Python 调试器深入调试代码。代价是序列化速度显著变慢仅适合排查阶段使用。已知问题与规避Python 3.8 与 3.9 的某些版本存在 pickle 模块的内存泄漏 bugbugs.python.org/issue39492可能导致用户遇到内存持续增长。该问题在Python 3.8.2rc1、Python 3.9.0 alpha 4 及之后的版本中已被修复因此遇到此类内存异常时优先核对并升级 Python 小版本即可。小结Ray 的序列化体系可以概括为三层底层是 Plasma 共享内存对象存储不可变对象、按需跨节点传输中间是 pickle5 cloudpickle 的协议 5 序列化器支持递归对象、lambda、动态类与 out-of-band 缓冲上层则是面向用户的零拷贝与自定义扩展能力NumPy 只读数组、可选 PyTorch 张量零拷贝、register_serializer/__reduce__自定义、inspect_serializability诊断。掌握这些机制与对应源码位置python/ray/_private/serialization.py、python/ray/_private/tensor_serialization_utils.py、python/ray/util/serialization.py、python/ray/util/check_serialize.py即可在实践中自如地规避只读数组报错、压降大对象传输开销并为任何第三方类型定制高效的序列化方案。【免费下载链接】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),仅供参考
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻