传递至工作进程时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
相关产品推荐
相关产品推荐

