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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 00:50:35