多进程Queue存储列表列表时进程无法终止问题求助
问题:多进程Queue导致进程无法终止,回溯中的锁是什么?
我尝试在proc进程函数里把嵌套列表放到multiprocessing.Queue里,调用event.set()后希望进程能自动终止。从打印信息看proc函数已经执行完了,但进程还在运行。如果减小每次put的列表数量(batchperq变量)或者嵌套列表的大小,程序就能正常运行。按下键盘中断后得到了如下回溯信息,想问问它尝试获取的“锁”是什么?是不是和Queue有关?
代码示例
import multiprocessing as mp import queue import numpy as np import time def main(): trainbatch_q = mp.Queue(10) batchperq = 50 event = mp.Event() tl1 = mp.Process(target=proc, args=( trainbatch_q, 20, batchperq, event)) tl1.start() time.sleep(3) event.set() tl1.join() print("Never printed..") def proc(batch_q, batch_size, batchperentry, the_event): nrow = 100000 i0 = 0 to_q = [] while i0 < nrow: rowend = min(i0 + batch_size,nrow) somerows = np.random.randint(0,5,(rowend-i0,2)) to_q.append(somerows.tolist()) if len(to_q) == batchperentry: print("adding..", i0, len(to_q)) while not the_event.is_set(): try: batch_q.put(to_q, block=False) to_q = [] break except queue.Full: time.sleep(1) i0 += batch_size print("proc finishes")
回溯信息
Traceback (most recent call last): File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/multiprocessing/process.py", line 252, in _bootstrap util._exit_function() File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/multiprocessing/util.py", line 322, in _exit_function _run_finalizers() File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/multiprocessing/util.py", line 262, in _run_finalizers finalizer() File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/multiprocessing/util.py", line 186, in __call__ res = self._callback(*self._args, **self._kwargs) File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/multiprocessing/queues.py", line 198, in _finalize_join thread.join() File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/threading.py", line 1056, in join self._wait_for_tstate_lock() File "/software/Anaconda3-5.0.0.1-el7-x86_64/lib/python3.6/threading.py", line 1072, in _wait_for_tstate_lock elif lock.acquire(block, timeout): KeyboardInterrupt
问题分析与解决方案
首先明确:回溯里的锁确实和multiprocessing.Queue直接相关,而且这就是进程无法终止的核心原因。
为什么会卡住?
multiprocessing.Queue内部有个后台线程,专门负责把数据从当前进程的内存复制到跨进程共享的存储区域(比如管道或者共享内存)。这个线程在运行时会持有一个线程状态锁,用来保证操作的线程安全。
你的问题出在:
proc函数执行完了,已经打印了proc finishes,但你往Queue里放的数据从来没被主进程消费过,队列里还积压着大量数据。- 当子进程要退出时,
Queue的终结器(_finalize_join)会被触发,它会尝试等待后台线程完成所有数据的发送工作,也就是调用thread.join()。 - 但因为主进程一直没取数据,后台线程会一直阻塞在发送数据的环节,永远停不下来。
join()内部需要获取线程的状态锁来确认线程是否结束,结果就卡在了锁的获取上,导致子进程无法正常退出。
至于为什么减小batchperq或者嵌套列表大小就正常?因为这时候队列的容量(你设的是10)能装下这些数据,后台线程很快就能把所有数据处理完,子进程退出时后台线程已经结束了,自然不会卡住。
怎么解决?
给你几个实用的方案:
方案1:主进程主动消费队列数据
在子进程join()之前,把队列里的所有数据都取出来,让后台线程能完成工作:
def main(): trainbatch_q = mp.Queue(10) batchperq = 50 event = mp.Event() tl1 = mp.Process(target=proc, args=( trainbatch_q, 20, batchperq, event)) tl1.start() time.sleep(3) event.set() # 新增:消费队列所有数据 while True: try: trainbatch_q.get(block=False) except queue.Empty: break tl1.join() print("Now this will be printed!")
方案2:显式关闭队列并等待后台线程
在确认子进程已经完成数据写入后,调用Queue.close()关闭队列,然后用join_thread()等待后台线程结束:
def main(): trainbatch_q = mp.Queue(10) batchperq = 50 event = mp.Event() tl1 = mp.Process(target=proc, args=( trainbatch_q, 20, batchperq, event)) tl1.start() time.sleep(3) event.set() # 等待子进程完成数据写入(也可以用另一个Event做更可靠的同步) time.sleep(1) trainbatch_q.close() trainbatch_q.join_thread() tl1.join() print("Now this will be printed!")
方案3:调整队列容量(仅适合数据量固定场景)
如果你的数据总量是固定的,可以把队列容量设得足够大,能装下所有要放入的数据,这样后台线程能一次性处理完:
# 计算总批次:100000行 / (20行/批次 * 50批次/队列条目) = 100个条目 trainbatch_q = mp.Queue(100)
内容的提问来源于stack exchange,提问作者phdscm
相关产品推荐
相关产品推荐

