Python多进程含Queue时进程终止异常及队列容量Bug问询
问题场景
用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的元素还没被主进程读取:
p4处理完所有数据后,会往queue4里放一个None标记结束,但此时queue4已经被填满,put(None)直接阻塞p4因为阻塞无法退出,导致p4.join()一直卡着- 元素少的时候,
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

