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

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}')
异常情况

异常发生时输出如下:
异常输出截图

问题原因与修复方案

核心问题

  1. 强制终止子进程导致结果丢失:调用self.qtasks.join()后立刻执行terminate(),此时子进程可能还在执行q_out.put(res)操作。JoinableQueue.join()仅保证所有任务被标记为task_done,不确保子进程已将结果完全写入输出队列,强制终止会中断未完成的put,导致队列状态异常。
  2. 依赖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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 09:10:48