Windows下进程内创建线程共享JoinableQueue时程序卡住问题排查
问题分析与解决方案
我尝试在2个进程中各创建3个线程,让所有线程共享
multiprocessing.JoinableQueue类型的队列。worker_func函数用于创建线程,thread_func函数从队列中取值并打印。但程序在time.sleep或队列的get()方法处卡住了,我哪里做错了?我在Windows电脑上运行程序。
原代码:
import threading from multiprocessing import Pool, Manager, JoinableQueue import multiprocessing from threading import Thread import time def thread_func(q, disp_lock): with disp_lock: print('thread ', threading.current_thread().name, ' in process ', multiprocessing.current_process().name , ' reporting for duty') while True: time.sleep(0.1) try: val = q.get_nowait() with disp_lock: print('thread ', threading.current_thread().name, ' in process ', multiprocessing.current_process().name , ' got value: ',val) q.task_done() except: with disp_lock: print('queue is empty: ', q.qsize()) def worker_func(num_threads, q, disp_lock): threads = [] for i in range(num_threads): thread = Thread(target= thread_func, args=( q, disp_lock,)) thread.daemon = True thread.start() if __name__ == "__main__": manager = Manager() lock = manager.Lock() q1 = JoinableQueue()#manager.Queue() q1_length = 20 for i in range(q1_length): q1.put(i) processes = [] num_processes = 2 # 2 processes num_threads = 3 for _ in range(num_processes): p = multiprocessing.Process(target=worker_func, args=( num_threads, q1, lock, )) # create a new Process p.daemon = True p.start() processes.append(p) q1.join()
错误原因分析
JoinableQueue跨进程共享失败:Windows采用spawn方式启动多进程,原生JoinableQueue()无法在进程间共享内存,子进程根本访问不到主进程放入队列的元素,导致线程一直阻塞等待数据。- 子进程无线程等待逻辑:
worker_func启动线程后直接返回,子进程会立即退出;而线程被设为daemon=True,子进程退出时会直接终止所有守护线程,线程来不及处理队列任务。 - 线程无限循环无终止条件:
thread_func的while True没有退出逻辑,即使队列空了也会一直循环尝试取值,既浪费资源也导致程序无法正常结束。 - 主进程未等待子进程:主进程仅调用
q1.join(),但如果子进程提前退出,队列的task_done()可能未被正确调用,导致join()一直阻塞。
修正后的代码
import threading import multiprocessing from threading import Thread import time def thread_func(q, disp_lock): with disp_lock: print(f'thread {threading.current_thread().name} in process {multiprocessing.current_process().name} reporting for duty') while True: try: # 用带超时的get替代get_nowait,避免无意义的频繁循环 val = q.get(timeout=1) with disp_lock: print(f'thread {threading.current_thread().name} in process {multiprocessing.current_process().name} got value: {val}') q.task_done() except multiprocessing.queues.Empty: # 检查队列所有任务是否完成,确认后退出线程 if q.unfinished_tasks == 0: with disp_lock: print(f'thread {threading.current_thread().name} in process {multiprocessing.current_process().name} exiting: queue completed') break continue except Exception as e: with disp_lock: print(f'thread {threading.current_thread().name} error: {str(e)}') break def worker_func(num_threads, q, disp_lock): threads = [] for i in range(num_threads): thread = Thread(target=thread_func, args=(q, disp_lock,)) thread.start() threads.append(thread) # 等待所有线程处理完任务 for t in threads: t.join() if __name__ == "__main__": manager = multiprocessing.Manager() lock = manager.Lock() # 通过Manager创建可跨进程共享的JoinableQueue q1 = manager.JoinableQueue() q1_length = 20 for i in range(q1_length): q1.put(i) processes = [] num_processes = 2 num_threads = 3 for idx in range(num_processes): p = multiprocessing.Process(target=worker_func, args=(num_threads, q1, lock,), name=f'Process-{idx+1}') # 不设置为守护进程,让主进程等待子进程完成 p.start() processes.append(p) # 等待队列所有任务标记完成 q1.join() # 等待所有子进程正常退出 for p in processes: p.join() print("All tasks completed.")
关键修改说明
- 改用
manager.JoinableQueue():确保队列能在Windows多进程环境下实现跨进程共享。 - 子进程等待线程:
worker_func中添加thread.join(),保证子进程在所有线程处理完任务后再退出。 - 线程添加退出条件:通过
q.unfinished_tasks == 0判断队列任务是否全部完成,完成后线程主动退出循环。 - 主进程等待子进程:在
q1.join()后添加p.join(),确保所有子进程都正常结束后,主进程才终止。 - 优化队列取值逻辑:用带超时的
get()替代get_nowait(),减少无意义的循环消耗。
内容的提问来源于stack exchange,提问作者dleal
相关产品推荐
相关产品推荐

