Python多进程Queue异常:无法从shared_queue获取数据求排查建议
排查思路与修复方案
先直接点出最核心的问题:你的dowork2方法定义有误!作为Test类的成员方法,它缺少了第一个必须的self参数,这直接导致池子里的进程调用该方法时参数完全错位,函数根本无法正常执行,自然不会往shared_queue里写入数据。
下面一步步拆解排查和修复步骤:
1. 修复类方法的参数问题
你现在的dowork2定义是:
def dowork2(queue,shared_queue):
但它属于Test类的成员方法,必须把self作为第一个参数,修正后应该是:
def dowork2(self, queue, shared_queue):
修正后,pool.apply_async中传递的(queue, shared_queue)参数才能正确对应方法里的形参,子进程才能正常执行逻辑。
2. 捕获子进程异常,定位隐藏错误
你用了apply_async提交异步任务,但没有获取任务的返回或异常信息,导致子进程出错时完全没有提示。可以在添加任务后,尝试获取结果来捕获异常:
results.append(pool.apply_async(self.dowork2,(queue,shared_queue))) # 新增:检查子进程是否抛出异常 try: results[0].get(timeout=1) except Exception as e: print(f"子进程执行出错: {e}")
这样就能快速定位到参数不匹配这类隐藏错误。
3. 验证跨进程队列的可用性
虽然这不是当前的核心问题,但要确保Manager创建的队列能在进程间正确共享。你在run方法里创建manager = mp.Manager(),再把队列传给池子里的进程,这个逻辑是没问题的,不过要注意不要提前关闭Manager进程。
修复后的完整代码
import time import multiprocessing as mp class Test(mp.Process): def __init__(self, **kwargs): mp.Process.__init__(self) self.daemon = False print('dosomething') def run(self): manager = mp.Manager() queue = manager.Queue() shared_queue = manager.Queue() pool = mp.Pool() # 提交修复后的异步任务 results = [] results.append(pool.apply_async(self.dowork2, (queue, shared_queue))) # 可选:开启这段代码检查子进程是否报错 # try: # results[0].get(timeout=2) # except Exception as e: # print(f"子进程错误: {e}") while True: time.sleep(0.2) t = time.time() queue.put('abc') queue.put('def') l = '' while not shared_queue.empty(): l = l + shared_queue.get() print(f"获取到的数据: {l}") print(f"耗时: %.4f" % (time.time()-t)) # 注意:while True是死循环,下面的close和join永远不会执行,后续优雅退出需处理该逻辑 pool.close() pool.join() # 修复:添加self参数 def dowork2(self, queue, shared_queue): while True: path = queue.get() shared_queue.put(path[-1:]) if __name__ == '__main__': t = Test() t.start()
额外提醒
你的run方法里有while True死循环,后面的pool.close()和pool.join()永远不会执行,这不会影响当前功能,但如果后续需要优雅退出程序,要补充对应的退出逻辑。
内容的提问来源于stack exchange,提问作者Relax ZeroC
相关产品推荐
相关产品推荐

