Python multiprocessing.Queue数据丢失问题求助
解决multiprocessing队列数据丢失与进程阻塞问题
问题原因分析
你的问题根源有两个:
- 错误使用
queue.cancel_join_thread():这个方法会让子进程退出时直接终止队列内部的写入线程,导致未完成写入的数据被丢弃,这就是子进程结果丢失的核心原因。 - 进程阻塞死锁:原生
multiprocessing.Queue有缓冲区限制,当子进程调用put写入大量数据时,缓冲区满后子进程会阻塞,直到主进程调用get腾出空间。你之前在所有进程结束后才读取队列,导致子进程无法完成put操作,进而无法正常join,引发卡顿。
解决方案
- 彻底移除
queue.cancel_join_thread(),让子进程正常等待队列写入线程完成数据传输。 - 调整队列读取时机:在子进程运行期间就开始读取队列数据,避免缓冲区满导致子进程阻塞。推荐用守护线程持续读取队列,既不影响主进程管理子进程,又能及时释放队列缓冲区。
修正后的完整代码
from multiprocessing import Process, Queue from queue import Empty import threading def play(chosen_word): l = [chosen_word, chosen_word] return l def partial_test(id, words, queue): print(f'Process {id} started and allocated {len(words)} words.') guesses = [] for word in words: guesses.append(play(word)) print(f"Process {id} has finished ALL WORDS.") # debugging only queue.put((id, guesses)) print(f'Process {id} added results to queue') print(f'Process {id} exited. Queue has approximately {queue.qsize()} elements') def full_test(): # 假设此处已定义:word_list, process_count, words_per_process # 创建结果队列 queue = Queue() # 初始化结果存储列表,按进程ID对应 results = [[] for _ in range(process_count)] # 线程锁,保证多线程操作results时的安全性 result_lock = threading.Lock() # 队列读取线程函数 def read_queue(): while True: try: # 超时等待,避免无限制阻塞 res = queue.get(timeout=5) with result_lock: results[res[0]] = res[1].copy() queue.task_done() except Empty: # 超时说明队列已空,所有数据读取完成 break # 启动守护线程读取队列 reader_thread = threading.Thread(target=read_queue, daemon=True) reader_thread.start() # 创建并启动子进程 processes = [ Process(target=partial_test, args=(x, word_list[x*words_per_process:(x+1)*words_per_process], queue)) for x in range(process_count-1) ] for process in processes: process.start() # 主进程处理最后一份任务 partial_test(process_count-1, word_list[(process_count-1)*words_per_process:], queue) # 等待所有子进程结束 for idx, process in enumerate(processes): process.join() print(f'Joined process {idx} with main thread.') # 等待队列中所有任务处理完毕 queue.join() # 等待读取线程完成收尾 reader_thread.join() print("All processes finished!") print(f"Got {len([r for r in results if r])} valid results!") # 后续处理results逻辑 # ...
关键细节说明
- 移除
cancel_join_thread():子进程退出时会默认等待队列写入线程完成数据传输,确保所有put的数据都能被主进程接收。 - 守护线程读取队列:在子进程运行期间持续读取数据,避免队列缓冲区溢出导致子进程阻塞,从根本解决了之前的卡顿问题。
queue.join()与queue.task_done():这对方法配合使用,确保队列中所有提交的任务都被处理完毕,不会遗漏任何结果。- 线程锁:因为读取线程和主进程可能同时操作
results列表,锁机制避免了数据竞争导致的结果错乱。
内容的提问来源于stack exchange,提问作者Doragon
相关产品推荐
相关产品推荐

