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
相关产品推荐
相关产品推荐

