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

Python搭配pebble使用multiprocessing queue时join()永久阻塞问题

问题原因

你代码里存在两个核心错误导致q.join()永久阻塞:

  • worker函数的终止条件调用错误:将判断事件是否触发的done.is_set()误写为设置事件触发的done.set()。set()方法返回值永远为None,所以if done.set()的判断永远不成立,worker不会主动退出;且每次循环都会主动把done事件设为触发状态,打乱全局状态。
  • 进程池容量不足导致调度异常:默认n_checkers=1时,pebble进程池只有1个工作进程,你先提交了5个生产者任务,再提交消费者任务,消费者会被排到所有生产者任务之后才会执行,如果你是要实现生产消费并行的场景,进程池容量至少要大于等于生产者+消费者的最小并发数。
修复方案

修改worker函数的终止条件,根据需要调整进程池容量即可正常运行,修复后代码如下:

from functools import partial
import multiprocessing as mp
import pebble
import queue
import time    

def add_to_queue(num, q):
    # 往队列q中添加元素
    time.sleep(2) # 模拟业务耗时
    print("putting on queue")
    q.put(num)
    print("put on queue done")
    return num    

def worker(q, output, done):
    # 持续从队列拉取元素,直到done事件被触发
    while True:
        # 修复:改为判断事件是否触发,而非主动设置事件
        if done.is_set():
            return
        try:
            print("Getting from queue")
            num = q.get(block=True, timeout=10)
            print("Got from queue")
        except queue.Empty:
            print("EMPTY QUEUE")
            continue
        time.sleep(num)
        output.append(num)
        # 标记元素处理完成
        q.task_done()
        print("task done")    

def main(n_checkers=2): # 修复:进程池容量至少留1个给消费者
    mgr = mp.Manager()
    q = mgr.Queue()
    output = mgr.list()
    done = mgr.Event()
    workers = []    

    add_partial = partial(add_to_queue, q=q)
    with pebble.ProcessPool(n_checkers) as pool:
        nums = [1, 2, 3, 4, 5]
        map_future = pool.map(add_partial, nums)
        for i in range(1): # 启动1个消费者
            print("SCHEDULING WORKER", i)
            ftr = pool.schedule(worker, args=(q, output, done))
            workers.append(ftr)
    
        for r in map_future.result():
            print(r)
    
        print("Joining Queue")
        q.join()
        done.set()
        for w in workers:
            w.result()
        
    print(output)
    
        
if __name__ == "__main__":
    main()

如果仍然存在阻塞问题,可以将Manager.Queue替换为普通multiprocessing.Queue,避免跨进程管理器的同步bug。

内容的提问来源于stack exchange,提问作者RSHAP

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 23:36:01