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

为什么multiprocessing队列关闭后调用Queue.get()仍会阻塞?

核心原因

你对multiprocessing.Queue的close()方法存在误解,该方法仅关闭调用方进程持有的队列写入端,不会主动广播队列关闭的状态到其他关联进程。只有满足以下两个条件时,get()调用才会触发队列关闭对应的异常:

  • 所有关联进程持有的队列写入端都已被关闭
  • 队列中已经没有未读取的剩余元素

你的代码中子进程也持有队列的完整引用,其对应的写入端没有被关闭,系统会判定仍有可能有进程向队列写入数据,因此q.get()会一直阻塞,不会抛出异常,最终导致程序挂起。

推荐解决方案

工业界最常用、兼容性最好的方案是哨兵值通知法,不需要依赖特定Python版本的队列关闭异常逻辑,稳定性更高:

from multiprocessing import Process, Queue

# 哨兵值,用于标识任务结束
SENTINEL = None

def worker_main(q):
    while True:
        message = q.get()
        # 收到哨兵值就退出循环
        if message is SENTINEL:
            break
        # 正常处理任务的逻辑写在下方
        # do_something(message)

if __name__ == '__main__':
    q = Queue()
    worker = Process(target=worker_main, args=(q,))
    worker.start()
    
    # 这里写生产任务、放入队列的逻辑,示例:
    # q.put("task1")
    # q.put("task2")
    
    # 所有任务生产完成后,放入和子进程数量相等的哨兵值
    q.put(SENTINEL)
    
    q.close()
    q.join_thread()
    worker.join()

如果你一定要用关闭队列触发异常的方案

需要确保所有进程的队列写入端都被关闭,且队列已空:

  1. 子进程启动后先关闭自身持有的队列写入端(因为子进程只需要读,不需要写队列)
  2. 主进程完成任务写入后关闭队列写入端,等待队列后台线程完成数据传输
    修改后的代码如下:
from multiprocessing import Process, Queue

def worker_main(q):
    # 子进程不需要写队列,先关闭自身的写入端
    q.close()
    while True:
        try:
            message = q.get()
            # 正常处理任务的逻辑写在下方
            # do_something(message)
        except ValueError:
            # 队列已关闭且为空,退出
            break

if __name__ == '__main__':
    q = Queue()
    worker = Process(target=worker_main, args=(q,))
    worker.start()
    
    # 生产任务逻辑写在这里
    # q.put("task1")
    
    q.close()
    q.join_thread()
    worker.join()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 04:06:01