Python多进程Queue意外抛出Empty异常问题排查
问题描述
在Ubuntu系统的Python 3.7环境中,我将multiprocessing.Queue封装为类,用于通过输入、输出队列批量执行任务函数。但从输出队列读取结果时,偶尔(每10-15次调用出现一次)会抛出Empty异常,尽管队列中仍存在元素。我曾尝试过带短超时的阻塞get方法,但问题依旧。此前未封装成类时代码运行正常,因此推测问题出在封装方式上。
封装类代码
import multiprocessing import os class ProcessParallel(): def __init__(self): self.qtasks = multiprocessing.JoinableQueue() def process_queue(self, target_func, q_out): pid = os.getpid() #print(f'[pid {pid}] started processing queue') while True: task = self.qtasks.get() if task is None: #print(f'[pid {pid}] got stop signal from queue') self.qtasks.task_done() break res = target_func(task) q_out.put(res) self.qtasks.task_done() def process(self, num_jobs, target_func, task_list, q_out, verbose = True): plist = [] for k in range(num_jobs): plist.append(multiprocessing.Process(target = self.process_queue, args = (target_func, q_out))) for p in plist: p.start() #--- populate the tasks queue, inc. a stop signal for each process for task in task_list: self.qtasks.put(task) for _ in range(num_jobs): self.qtasks.put(None) if verbose: print('waiting for the tasks queue to join') self.qtasks.join() if verbose: print('tasks queue joined') print(f'terminating {len(plist)} process') for p in plist: p.terminate() if verbose: print('done')
测试代码
def myfun(x): return 3 * x + 1 q_out = multiprocessing.Queue() ppar = ProcessParallel() ppar.process(4, myfun, [1,2,3,4], q_out, True) print(f'{q_out.qsize()} results in output queue') for _ in range(q_out.qsize()): #r = q_out.get_nowait() # non-blocking call also raises Empty exception occasionally r = q_out.get(timeout = 0.01) print(f'got item from queue: {r}')
异常情况
异常发生时输出如下:
问题原因与修复方案
核心问题
- 强制终止子进程导致结果丢失:调用
self.qtasks.join()后立刻执行terminate(),此时子进程可能还在执行q_out.put(res)操作。JoinableQueue.join()仅保证所有任务被标记为task_done,不确保子进程已将结果完全写入输出队列,强制终止会中断未完成的put,导致队列状态异常。 - 依赖
qsize()不可靠:qsize()返回的是队列元素数量的瞬时快照,多进程环境下调用后队列实际数量可能已变化,以此作为循环次数会引发读取错误。
修复代码
修改后的封装类
import multiprocessing import os class ProcessParallel(): def __init__(self): self.qtasks = multiprocessing.JoinableQueue() def process_queue(self, target_func, q_out): pid = os.getpid() #print(f'[pid {pid}] started processing queue') while True: task = self.qtasks.get() if task is None: #print(f'[pid {pid}] got stop signal from queue') self.qtasks.task_done() break res = target_func(task) q_out.put(res) self.qtasks.task_done() def process(self, num_jobs, target_func, task_list, q_out, verbose = True): plist = [] for k in range(num_jobs): plist.append(multiprocessing.Process(target = self.process_queue, args = (target_func, q_out))) for p in plist: p.start() #--- populate the tasks queue, inc. a stop signal for each process for task in task_list: self.qtasks.put(task) for _ in range(num_jobs): self.qtasks.put(None) if verbose: print('waiting for the tasks queue to join') self.qtasks.join() if verbose: print('tasks queue joined') print(f'waiting for {len(plist)} processes to exit') # 等待子进程正常退出,而非强制终止 for p in plist: p.join() if verbose: print('done')
修改后的测试代码
def myfun(x): return 3 * x + 1 q_out = multiprocessing.Queue() task_list = [1,2,3,4] ppar = ProcessParallel() ppar.process(4, myfun, task_list, q_out, True) # 根据任务数量确定读取次数,不依赖qsize() for _ in range(len(task_list)): r = q_out.get() print(f'got item from queue: {r}')
额外说明
- Unix系统(包括Ubuntu)中
multiprocessing.Queue.qsize()的结果仅作参考,不能精确依赖其判断队列元素数量。 - 避免使用
terminate()强制终止子进程,让子进程通过处理None信号正常退出,可确保所有I/O操作完成,避免队列状态异常。
内容的提问来源于stack exchange,提问作者Itamar Katz
相关产品推荐
相关产品推荐

