Python Multiprocessing Queue取最后一项阻塞的原因与解决方法
Multiprocessing Queue 读取最后一条数据阻塞的原因与解决方案
阻塞原因
- 子进程未正常退出,管道写入端未关闭
multiprocessing.Queue基于操作系统管道实现,父进程调用get()时,若子进程仍存活,管道的写入端处于开放状态,父进程会认为可能还有数据写入,因此会一直阻塞等待,哪怕最后一条数据已经放入队列。 - 未明确数据结束标记
父进程没有判断子进程是否已经完成所有数据的写入,仅通过get()无限等待,没有终止条件,导致在最后一条数据读取后(或读取前)陷入阻塞。 - 缺少超时机制
代码中queue.get()未设置超时参数,一旦队列中暂时无数据(或子进程未完成写入),父进程会无限阻塞。
解决方案
1. 等待子进程完成后再读取
在父进程中,先调用子进程的join()方法,等待所有子进程执行完毕、正常退出后,再从队列中读取数据。此时管道写入端已关闭,父进程能一次性读取所有数据,不会阻塞。
示例代码:
from multiprocessing import Process, Queue def worker(q): # 写入4条数据 for i in range(4): q.put(f"data_{i}") if __name__ == "__main__": q = Queue() p = Process(target=worker, args=(q,)) p.start() p.join() # 等待子进程完全退出 # 读取所有数据 while not q.empty(): print(q.get())
2. 使用哨兵标记数据结束
在子进程写入完所有数据后,放入一个约定的特殊值(如None)作为“结束哨兵”。父进程读取到该值时,停止读取操作,避免无限等待。
示例代码:
from multiprocessing import Process, Queue def worker(q): for i in range(4): q.put(f"data_{i}") q.put(None) # 写入哨兵 if __name__ == "__main__": q = Queue() p = Process(target=worker, args=(q,)) p.start() # 读取直到遇到哨兵 while True: item = q.get() if item is None: break print(item) p.join()
3. 为读取操作设置超时
在queue.get()中添加timeout参数,避免无限阻塞。同时捕获queue.Empty异常,处理队列无数据的情况。
示例代码:
from multiprocessing import Process, Queue from queue import Empty def worker(q): for i in range(4): q.put(f"data_{i}") if __name__ == "__main__": q = Queue() p = Process(target=worker, args=(q,)) p.start() try: while True: item = q.get(timeout=3) # 设置3秒超时 print(item) except Empty: print("队列已无数据,读取结束") p.join()
注意事项
- 若有多个子进程,需确保每个子进程都写入哨兵,或统计哨兵数量(比如N个子进程就等待N个哨兵),避免提前终止读取。
- 不要在子进程中保留队列的引用,确保子进程退出时能正确释放相关资源。
内容的提问来源于stack exchange,提问作者DrM
相关产品推荐
相关产品推荐

