跨进程非序列化数据传递 (Passing Values That Cannot Be Serialized)
在复杂的 AI 工作流中,许多重型运行时对象根本无法转换为纯文本或 JSON 数据。例如加载到 GPU 现存中的 Diffusers Pipeline 管道、Latent 潜空间张量、已打开的文件指针句柄、C++ 底层驱动连接等。
如果你的算子库运行在物理隔离的独立 Worker 子进程中(参见 基于 Worker 子进程的算子物理隔离),主编排进程与 Worker 之间默认通过 JSON 进行通信。要实现重型内存对象在节点间高效流转,必须借助引擎的内存对象持久化与跨进程 Key 代理机制。
核心极简使用准则:
在产出该对象的输出端口上显式声明
serializable=False并赋给原生对象;消费它的下游节点无需任何特殊声明,像往常一样正常读取即可。
class LoadPipeline(ControlNode):
def __init__(self, **kwargs) -> None:
super().__init__(**kwargs)
# 1. 在生产端声明 serializable=False
self.add_parameter(
Parameter(
name="pipeline",
output_type="Pipeline",
tooltip="已装载进 VRAM 的 Diffusers 管道实例",
serializable=False,
allowed_modes={ParameterMode.OUTPUT},
)
)
def process(self) -> None:
# 直接赋给原生复杂的 Python 管道对象
self.parameter_output_values["pipeline"] = load_pipeline(...)
class Generate(ControlNode):
def __init__(self, **kwargs) -> None:
super().__init__(**kwargs)
# 2. 消费端无需任何标记,正常声明普通端口即可
self.add_parameter(Parameter(name="pipeline", input_types=["Pipeline"], tooltip="待执行的管道"))
def process(self) -> None:
# 3. 提取到的直接就是原本真实的 Python 原生管道对象
pipeline = self.get_parameter_value("pipeline")
...
背后底层发生了什么?(What actually happens)
该原生重量级对象自始至终被牢牢锁死在生产它的那个 Python 物理进程内存中,绝对不会跨越网络或进行任何二进制序列化复制。
当该参数准备穿越跨进程边界时,引擎底层会自动将其置换为一个短小的不透明 Key(访问凭据) 并传输该 Key。当下游位于同一 Worker 内部的节点读取该端口时,底层又瞬间通过该 Key 还原为原本的原生对象内存指针。
必须内化的三大推论:
1. 当前节点的字典内永远持有真身:写入 parameter_output_values["pipeline"] 后紧接着原地读取,拿到的直接是真身,内部绝对不会被偷换为 Key;
2. 在主进程 Shared 模式下完全退化为原生传引:若未开启 Worker 隔离,数据直接遵循 Python 原生内存引用传递,无任何额外性能损耗;
3. 唯有生产者需要显式声明:消费端完全不用关心上游是否非序列化,直接声明普通 Parameter。
serializable=False 的确切语义与数据路由表现
serializable=False 既能阻止该字段被写入磁盘工作流文件,又指导引擎拦截并跨进程代理非数据对象:
赋予 serializable=False 的参数值类型 |
跨进程跨越时的实际形态 | 底层架构原因 |
|---|---|---|
| Pipeline 管道、Tensor 张量、驱动句柄 | 生成全局 Key 代理;真身留在当前进程 | 无法进行任何数据形式的序列化 |
| 普通的 API 密钥字符串 | 原样字符串传输直通 | 本身就是轻量数据,若生成 Key 反而在远端无法解析 |
数字字典 dict、文本列表 list |
原样 JSON 字典/列表传输直通 | 已具备标准数据结构 |
ImageUrlArtifact 等资产类 |
生成全局 Key 代理;真身留在当前进程 | 见下文特别说明 |
⚠️ 特别提示:若希望包含 URL 的图像等资产正常跨进程漫游,切勿在输出端口添加 serializable=False!未声明该标记的普通 ImageUrlArtifact 会自动以轻量 URL 字符串形式正常跨进程流转。
跨多次运行的昂贵显存模型缓存 (local_objects)
对于加载一次可能耗时 30 秒的重型模型,若希望在多次工作流执行之间保持长驻复用,使用节点专属的 local_objects 缓存 API:
def process(self) -> None:
# 1. 基于配置哈希生成专属 Key
key = self.local_objects.key_for(self._config_hash())
# 2. 尝试从本地进程内存中直接提取已就绪的模型
pipeline = self.local_objects.get(key)
if pipeline is None:
# 3. 缓存未命中:执行物理耗时加载并登记显存释放回调
pipeline = build_pipeline(...)
self.local_objects.put(pipeline, key=self._config_hash(), on_drop=release_vram)
# 执行下游推理...
| 核心 API 调用 | 功能语义与职责 |
|---|---|
local_objects.put(value, *, key, on_drop=None) |
在指定 Key 下持有该对象,可附加释放清理回调 |
local_objects.get(key) |
提取该进程内持有的真实对象;未命中返回 None |
local_objects.key_for(suffix) |
根据后缀派生带有库命名空间前缀的全局安全 Key |
local_objects.drop(key) |
手动主动释放并销毁指定的一个驻留对象 |
local_objects.drop_all() |
释放并清空本库在当前进程内驻留的全部显存对象(“清理缓存”节点专用) |
显存与 GPU 资源的主动释放机制 (on_local_object_drop)
在 Python 中单纯丢弃变量引用并不能立即触发 PyTorch 清理 CUDA VRAM 显存。
在声明参数时传入 on_local_object_drop(或在 put 时传入 on_drop),引擎在清理销毁该对象时会自动回调它:
Parameter(
name="pipeline",
output_type="Pipeline",
tooltip="已加载的管道",
serializable=False,
allowed_modes={ParameterMode.OUTPUT},
# 当缓存决定丢弃该管道时,自动将其转移至 CPU 内存或清空显存缓存
on_local_object_drop=lambda pipeline: pipeline.to("cpu"),
)
何时会触发显存释放回调? - 算子再次运行,为该输出端口换入了全新的不同对象; - 该算子节点在画布中被用户主动删除,且无其他连线依赖; - 算子库被卸载或整机热重载; - 当前画布工作流被关闭或清空。
严禁踩入的绝对禁区 (What you cannot do)
- 容器参数 (
ParameterList/ParameterDictionary) 无法承载持有状态: 容器参数的数值完全由其动态子成员拼装而来,不存在单一的宿主实体挂载释放回调。若需输出一组 Tensor 批次,请使用单一的普通Parameter承载list[Tensor]; - 对象绝对无法跨出创建它的物理子进程:
若试图在主编排进程(或另一个不同的 Worker 子进程)中去反查由该 Worker 产出的 Key,会直接报出明确错误:
Attempted to read the value for parameter 'pipeline' on node 'Generate'. Failed due to: it is held in another process...不同算子库之间若需传递资产,必须通过落地为磁盘文件或公网 URL 进行交互; - 输入端口绝不参与缓存:
Key 的发放只发生在输出端。若一个对象未在输出端标记
serializable=False就直接灌入下游,跨进程时会被直接打回原形。