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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 04:27:43