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

multiprocessing.Queue队列过大导致进程挂起问题咨询

multiprocessing.Process与Queue队列过大导致进程挂起的问题解析

问题描述

使用Python的multiprocessing模块的Process与Queue时,当队列元素过多会出现进程挂起现象,目前未确定触发问题的队列临界大小。

最初推测挂起是因为队列未清空且未关闭导致,但队列仅含少量元素时,即使不清空关闭也不会触发问题;而清空并关闭队列则可避免该问题,对此存在困惑,需阐明进程挂起的原因。


演示

  1. 使用p = Process(target=some_function_that_does_not_break)时,控制台输出:
Function started
Function ended
queue size is=100
J1
J2
main ended
  1. 使用p = Process(target=some_function_that_also_does_not_break)时,控制台输出:
Function started
Function ended. Time filling queue: 1.999461717 seconds
Emptying queue
queue size is=3005706
J1
Queue empty! Time emptying queue: 26.81551289 seconds
Queue closed
J2
main ended

(不清楚为何清空队列的耗时远大于填充队列)

  1. 使用p = Process(target=some_function_that_breaks)时,控制台输出:
Function started
Function ended
queue size is=3152815
J1
(execution hanging here)

代码

#!/usr/bin/env python
import time
from multiprocessing import Process, Queue, Value

q = Queue()
stop = Value("b", False)

# [EDIT]: 新增try...except...是为了避免关于队列已满的无意义讨论,
# 因为该队列没有设置最大容量,理论上不会满。
def some_function_that_breaks():
    print("Function started")
    try:
        while not stop.value:
            q.put("Item")
    except Exception as e:
        print(f"发生异常: {e}")

    print("Function ended")


def some_function_that_does_not_break(queue_size=100):
    print("Function started")
    for _ in range(queue_size):
        q.put("Item")
    print("Function ended")

def some_function_that_also_does_not_break():
    print("Function started")
    time_start = time.perf_counter_ns()
    while not stop.value:
        q.put("Item")

    time_filling_s = (time.perf_counter_ns() - time_start) / 1e9
    print(f"Function ended. 填充队列耗时: {time_filling_s} seconds")

    print("清空队列")
    time_start = time.perf_counter_ns()
    while not q.empty():
        q.get()

    time_emptying_s = (time.perf_counter_ns() - time_start) / 1e9

    print(f"队列已空! 清空队列耗时: {time_emptying_s} seconds")
    q.close()
    print("队列已关闭")


p = Process(target=some_function_that_does_not_break)

p.start()
time.sleep(2)
stop.value = True
time.sleep(1)

print(f"队列大小={q.qsize()}")

print("J1")
p.join()
print("J2")
p.join()

print("main ended")

问题原因分析

1. 进程挂起的核心原因

multiprocessing.Queue 底层基于**管道(pipe)**和后台线程实现:

  • 子进程调用put()时,数据先写入管道,由Queue的后台feeder线程将管道数据转存到队列的内存缓冲区。
  • 子进程结束时,会自动触发Queue的close()操作,该操作需要等待feeder线程完成所有未完成的写入任务,同时确保管道内的所有数据都被主进程端的Queue接收。
  • 如果队列累积了大量未被消费的元素,主进程未调用get()读取这些数据,子进程的close()操作会一直等待管道数据被处理,进而导致p.join()挂起——因为join()需要等待子进程完全释放所有资源后才能返回。

而队列元素较少时,管道内的数据能快速处理完毕,因此即使不消费也不会触发挂起。

2. 清空队列耗时远大于填充的原因

  • put()的批量优化:put()操作时,后台feeder线程会批量将数据写入管道,减少了进程间通信的系统调用次数,速度更快。
  • get()的逐个读取:get()是逐个从队列中取出元素,每次操作都涉及进程间的状态同步,系统调用开销更大。加上代码中用while not q.empty()循环(q.empty()并非线程安全,但此处为单主进程读取),每次循环都要检查队列状态,进一步增加了耗时。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 15:23:10