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

Python multiprocessing.Queue数据丢失问题求助

解决multiprocessing队列数据丢失与进程阻塞问题

问题原因分析

你的问题根源有两个:

  1. 错误使用queue.cancel_join_thread():这个方法会让子进程退出时直接终止队列内部的写入线程,导致未完成写入的数据被丢弃,这就是子进程结果丢失的核心原因。
  2. 进程阻塞死锁:原生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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 07:05:28