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

Python多进程含Queue时进程终止异常及队列容量Bug问询

Python多进程Queue导致join阻塞的问题分析与解决

问题场景

用multiprocessing.Queue做进程间通信时,调用process.join()会出现阻塞:只有当队列被清空时进程才会终止。而且这个问题只在队列元素超过一定数量(比如533个以上)时触发,元素少(如100个)时完全正常。示例代码如下:

import random
import multiprocessing
import time


def process1(data, queue1):
    for item in data:
        queue1.put(item)

    queue1.put(None)


def process2(queue1, queue2):
    while True:
        item = queue1.get()
        if item is None:
            queue2.put(None)
            break
        queue2.put(item * 2)


def process3(queue2, queue3):
    while True:
        item = queue2.get()
        if item is None:
            queue3.put(None)
            break
        queue3.put(item * 2)


def process4(queue3, queue4):
    while True:
        item = queue3.get()
        if item is None:
            queue4.put(None)
            print("None received")
            break
        queue4.put(item * 2)
    print("End process4")


if __name__ == "__main__":
    start_time = time.time()

    data = [random.randint(1, 10000) for _ in range(10000)]

    queue1 = multiprocessing.Queue()
    queue2 = multiprocessing.Queue()
    queue3 = multiprocessing.Queue()
    queue4 = multiprocessing.Queue()

    p1 = multiprocessing.Process(target=process1, args=(data, queue1))
    p2 = multiprocessing.Process(target=process2, args=(queue1, queue2))
    p3 = multiprocessing.Process(target=process3, args=(queue2, queue3))
    p4 = multiprocessing.Process(target=process4, args=(queue3, queue4))

    p1.start()
    p2.start()
    p3.start()
    p4.start()
    print("Proc started")

    p1.join()
    print("p1 join")
    p2.join()
    print("p2 join")
    p3.join()
    print("p3 join")
    p4.join()
    print("p4 join")

    results = []
    print("While loop")
    while True:
        item = queue4.get()
        if item is None:
            break
        results.append(item)

    end_time = time.time()

    print(end_time - start_time)

原因拆解

multiprocessing.Queue内部有默认的缓冲区阈值(不同环境下数值不同,这里碰到的是~533),当队列元素达到这个阈值时,put()操作会阻塞,直到有元素被get()取出、腾出空间。

你的代码逻辑里,主进程启动所有子进程后,立刻按顺序调用p1.join()→p2.join()→p3.join()→p4.join(),但此时queue4的元素还没被主进程读取:

  1. p4处理完所有数据后,会往queue4里放一个None标记结束,但此时queue4已经被填满,put(None)直接阻塞
  2. p4因为阻塞无法退出,导致p4.join()一直卡着
  3. 元素少的时候,queue4缓冲区没被占满,put(None)能顺利执行,所以p4可以正常退出,join()也能完成

设计逻辑

multiprocessing.Queue的阻塞设计是为了防止内存溢出,属于生产者-消费者模型的同步机制:当队列满时,生产者必须等待消费者消费数据,避免无限制写入导致内存被耗尽。

解决办法

方法1:调整等待顺序(最直接,符合需求)

不要先调用p4.join(),而是先读取queue4的结果,等所有元素消费完再join。这样queue4的空间会被及时释放,p4的put()不会阻塞:

修改主进程的代码顺序:

if __name__ == "__main__":
    start_time = time.time()

    data = [random.randint(1, 10000) for _ in range(10000)]

    queue1 = multiprocessing.Queue()
    queue2 = multiprocessing.Queue()
    queue3 = multiprocessing.Queue()
    queue4 = multiprocessing.Queue()

    p1 = multiprocessing.Process(target=process1, args=(data, queue1))
    p2 = multiprocessing.Process(target=process2, args=(queue1, queue2))
    p3 = multiprocessing.Process(target=process3, args=(queue2, queue3))
    p4 = multiprocessing.Process(target=process4, args=(queue3, queue4))

    p1.start()
    p2.start()
    p3.start()
    p4.start()
    print("Proc started")

    p1.join()
    print("p1 join")
    p2.join()
    print("p2 join")
    p3.join()
    print("p3 join")

    # 先读取queue4的结果,再join p4
    results = []
    print("While loop")
    while True:
        item = queue4.get()
        if item is None:
            break
        results.append(item)
    
    p4.join()
    print("p4 join")

    end_time = time.time()

    print(end_time - start_time)

方法2:使用JoinableQueue(适合复杂场景)

JoinableQueue支持任务确认机制,通过task_done()和join()来同步,不依赖队列缓冲区状态。需要修改所有子进程的逻辑:

# 替换所有Queue为JoinableQueue
queue1 = multiprocessing.JoinableQueue()
queue2 = multiprocessing.JoinableQueue()
queue3 = multiprocessing.JoinableQueue()
queue4 = multiprocessing.JoinableQueue()

# 修改子进程函数,每个get后调用task_done()
def process1(data, queue1):
    for item in data:
        queue1.put(item)
    queue1.put(None)
    queue1.join()  # 等待所有元素被处理

def process2(queue1, queue2):
    while True:
        item = queue1.get()
        queue1.task_done()  # 标记当前元素已处理
        if item is None:
            queue2.put(None)
            queue2.join()
            break
        queue2.put(item * 2)

# process3、process4同理,每个get操作后添加task_done()调用

方法3:扩大队列缓冲区(不推荐)

创建Queue时通过maxsize设置更大的缓冲区,比如queue4 = multiprocessing.Queue(maxsize=10000)。但这种方法只是延缓问题,当数据量超过设置值时,阻塞依然会出现,还可能占用过多内存。

内容的提问来源于stack exchange,提问作者The Gum Machine

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 17:13:18