You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

传递至工作进程时multiprocessing Queue引用丢失问题求助

问题:multiprocessing Manager Queue传递至工作进程池时出现引用丢失错误

主进程创建全局multiprocessing.Manager,为每个待处理任务生成该Manager下的Queue对象并传递给工作进程池(每个任务对应一个Queue),访问Queue时出现引用丢失,报错如下:

Traceback (most recent call last):
  File "/net/home/cmosig/miniconda3/envs/env_conda/lib/python3.12/multiprocessing/managers.py", line 209, in _handle_request
    result = func(c, *args, **kwds)
             ^^^^^^^^^^^^^^^^^^^^^^
  File "/net/home/cmosig/miniconda3/envs/env_conda/lib/python3.12/multiprocessing/managers.py", line 438, in incref
    raise ke
  File "/net/home/cmosig/miniconda3/envs/env_conda/lib/python3.12/multiprocessing/managers.py", line 426, in incref
    self.id_to_refcount[ident] += 1
    ~~~~~~~~~~~~~~~~~~~^^^^^^^
KeyError: '712cf1996780'

最小复现代码

import multiprocessing as mp
from joblib import Parallel, delayed

GLOBAL_QUEUE_MANAGER = mp.Manager()


def process(queue):
    queue.put("test")
    queue.close()
    return

def job_generator():
    for i in range(10):
        queue = GLOBAL_QUEUE_MANAGER.Queue()
        yield {"queue": queue}


Parallel(n_jobs=3, batch_size=1,
         backend="multiprocessing")(delayed(process)(**p)
                                    for p in job_generator())

约束条件

  • 避免提前创建所有任务的Queue对象(任务可能多达数千个,担心内存开销过大)
  • 工作进程为守护进程,无法自行创建Queue

问题原因

错误源于joblib的multiprocessing后端处理生成器任务时,Manager创建的Queue代理对象的引用计数管理异常:生成器中创建的Queue对象在被传递到工作进程前,可能已被主进程的垃圾回收机制回收,导致Manager端丢失了该对象的引用记录。

解决方法

1. 手动维护Queue引用,防止提前回收

在主进程中保留所有已创建Queue的引用,直到任务全部完成,避免垃圾回收提前清理这些对象:

import multiprocessing as mp
from joblib import Parallel, delayed

GLOBAL_QUEUE_MANAGER = mp.Manager()
active_queues = []  # 保存Queue引用,阻止提前回收

def process(queue):
    queue.put("test")
    queue.close()
    return

def job_generator():
    for i in range(10):
        queue = GLOBAL_QUEUE_MANAGER.Queue()
        active_queues.append(queue)
        yield {"queue": queue}

# 执行任务
Parallel(n_jobs=3, batch_size=1, backend="multiprocessing")(
    delayed(process)(**p) for p in job_generator()
)

# 任务完成后释放引用
active_queues.clear()

2. 改用标准库multiprocessing.Pool替代joblib

multiprocessing.Pool对Manager代理对象的引用管理更稳定,适合此类场景:

import multiprocessing as mp

GLOBAL_QUEUE_MANAGER = mp.Manager()

def process(queue):
    queue.put("test")
    queue.close()
    return

def job_generator():
    for i in range(10):
        queue = GLOBAL_QUEUE_MANAGER.Queue()
        yield (queue,)

if __name__ == "__main__":
    with mp.Pool(3) as pool:
        pool.starmap(process, job_generator())

3. 分批次生成并处理任务

通过分批次生成任务参数,既不会一次性创建数千个Queue,也能保证当前批次的Queue引用被保留,直到该批次任务完成:

import multiprocessing as mp
from joblib import Parallel, delayed

GLOBAL_QUEUE_MANAGER = mp.Manager()

def process(queue):
    queue.put("test")
    queue.close()
    return

def batch_generator(batch_size=50):
    total_tasks = 10  # 替换为实际任务总数
    generated = 0
    while generated < total_tasks:
        batch = []
        batch_count = min(batch_size, total_tasks - generated)
        for _ in range(batch_count):
            queue = GLOBAL_QUEUE_MANAGER.Queue()
            batch.append({"queue": queue})
            generated += 1
        yield batch

# 分批次执行
for batch in batch_generator():
    Parallel(n_jobs=3, batch_size=1, backend="multiprocessing")(
        delayed(process)(**p) for p in batch
    )

内容的提问来源于stack exchange,提问作者cmosig

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 02:20:10