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

多进程多线程架构下结果获取及multiprocessing Queue阻塞问题排查

问题排查:Python多进程使用Queue收集结果导致程序挂起

问题描述

我实现了一套多进程加线程的Python代码,用于处理PO文件中的翻译模式匹配。当前仅打印线程返回结果时程序能正常运行(耗时约20.5秒,所有进程可正常JOIN),但尝试用multiprocessing.Queue收集结果时,程序会在运行末尾挂起,无法正常结束。

原始代码

#!/usr/bin/env python3
import concurrent.futures
import multiprocessing
import time
import re
import os
from multiprocessing.managers import BaseManager

from sphinx_intl import catalog as c
from translation_finder import TranslationFinder
from definition import Definitions as df
from testFindResult import makeBatches, PatternFoundResult, FormatUpperCase, findEachMessage, FormatBase

class FindPatternProcess(multiprocessing.Process):
    def __init__(self,
                 start_time,
                 tf: TranslationFinder,
                 pat: re.Pattern,
                 batch,
                 exit_cond,
                 queue,
                 ): ### new code
        super(FindPatternProcess, self).__init__()
        self.batch = batch
        self.tf: TranslationFinder = tf
        self.exit_cond = exit_cond ### new code
        self.result_queue = queue
        self.start_time = start_time
        self.pat = pat
        self.tf: TranslationFinder

    def updateResult(self, result):
        (execution_time, line_number, proc_name, return_values) = result
        # self.result_queue.put(return_values)
        print(return_values)

    def run(self):
        start_time = self.start_time
        formatter = FormatUpperCase()
        with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
            thread_list = [
                executor.submit(findEachMessage, start_time, m, self.tf, formatter, self.pat)
                for m in self.batch]

            # done, not_done = concurrent.futures.wait(thread_list, return_when=concurrent.futures.ALL_COMPLETED)
            for th in concurrent.futures.as_completed(thread_list):
                try:
                    result = th.result(timeout=5)
                    list(filter(self.updateResult, result))
                except Exception as e:
                    print(e, th)
                    raise e


if __name__ == '__main__':
    start_time = time.perf_counter()
    home_dev = os.environ['DEV']
    input_path = os.path.join(home_dev, "current_blender_manual_merge_flat_0001.po")
    input_cat = c.load_po(input_path)

    BaseManager.register('PatternFoundResult', PatternFoundResult)
    BaseManager.register('TranslationFinder', TranslationFinder)
    manager = BaseManager()
    manager.start()

    tf = manager.TranslationFinder()
    pat = df.UPPERCASE_UNTRANSLATED_PATTERN

    batch_list = makeBatches(input_cat)
    exit_cond = multiprocessing.Event()
    batch_len = len(batch_list)
    queue = multiprocessing.Queue()

    result = PatternFoundResult('df.UPPERCASE_UNTRANSLATED_PATTERN')
    result.output_file = 'findNonMulti.po'
    result.output_text_file = 'findNonMulti.txt'

    procs = {}
    exit_cond = multiprocessing.Event() ### new code

    for (index, batch) in enumerate(batch_list):
        p = FindPatternProcess(start_time, tf, pat, batch, exit_cond, queue)
        procs[index] = p

    for (index, p) in procs.items():
        p.start()

    time.sleep(1)

    exit_cond.set() ### new code
    for (index, p) in procs.items():
        print('JOIN process:', p)
        p.join()

    print(f'Execution took {time.perf_counter() - start_time} seconds')

问题原因

程序挂起的核心原因是死锁:

  1. 子进程调用queue.put(return_values)时,若队列被填满,子进程会阻塞等待队列释放空间。
  2. 主进程在调用p.join()前未读取队列数据,导致子进程一直卡在put操作,主进程则等待子进程结束,形成双向阻塞。
  3. 代码中exit_cond变量未在子进程的run方法中使用,设置exit_cond.set()无法终止子进程,属于冗余代码。

解决方案

  1. 主进程提前消费队列数据:在join子进程前,启动线程异步读取队列,避免子进程因队列满阻塞。
  2. 修复updateResult方法:替换打印逻辑为队列存入操作。
  3. 移除无效的exit_cond逻辑:子进程未监听该事件,无需保留。

修改后的代码

#!/usr/bin/env python3
import concurrent.futures
import multiprocessing
import time
import re
import os
from multiprocessing.managers import BaseManager

from sphinx_intl import catalog as c
from translation_finder import TranslationFinder
from definition import Definitions as df
from testFindResult import makeBatches, PatternFoundResult, FormatUpperCase, findEachMessage, FormatBase

class FindPatternProcess(multiprocessing.Process):
    def __init__(self,
                 start_time,
                 tf: TranslationFinder,
                 pat: re.Pattern,
                 batch,
                 queue,
                 ):
        super(FindPatternProcess, self).__init__()
        self.batch = batch
        self.tf: TranslationFinder = tf
        self.result_queue = queue
        self.start_time = start_time
        self.pat = pat

    def updateResult(self, result):
        (execution_time, line_number, proc_name, return_values) = result
        self.result_queue.put(return_values)  # 改为存入队列
        # print(return_values)  # 可选保留打印

    def run(self):
        start_time = self.start_time
        formatter = FormatUpperCase()
        with concurrent.futures.ThreadPoolExecutor(max_workers=2) as executor:
            thread_list = [
                executor.submit(findEachMessage, start_time, m, self.tf, formatter, self.pat)
                for m in self.batch]

            for th in concurrent.futures.as_completed(thread_list):
                try:
                    result = th.result(timeout=5)
                    list(filter(self.updateResult, result))
                except Exception as e:
                    print(e, th)
                    raise e


if __name__ == '__main__':
    start_time = time.perf_counter()
    home_dev = os.environ['DEV']
    input_path = os.path.join(home_dev, "current_blender_manual_merge_flat_0001.po")
    input_cat = c.load_po(input_path)

    BaseManager.register('PatternFoundResult', PatternFoundResult)
    BaseManager.register('TranslationFinder', TranslationFinder)
    manager = BaseManager()
    manager.start()

    tf = manager.TranslationFinder()
    pat = df.UPPERCASE_UNTRANSLATED_PATTERN

    batch_list = makeBatches(input_cat)
    queue = multiprocessing.Queue()

    result = PatternFoundResult('df.UPPERCASE_UNTRANSLATED_PATTERN')
    result.output_file = 'findNonMulti.po'
    result.output_text_file = 'findNonMulti.txt'

    procs = {}
    for (index, batch) in enumerate(batch_list):
        p = FindPatternProcess(start_time, tf, pat, batch, queue)
        procs[index] = p

    # 启动所有进程
    for p in procs.values():
        p.start()

    # 启动线程读取队列,避免主进程阻塞
    collected_results = []
    def consume_queue():
        while True:
            try:
                item = queue.get(timeout=2)  # 超时时间可根据实际调整
                collected_results.append(item)
            except multiprocessing.queues.Empty:
                # 检查所有进程是否已结束,若结束则退出循环
                if all(not p.is_alive() for p in procs.values()):
                    break

    consumer_thread = multiprocessing.Process(target=consume_queue)
    consumer_thread.start()

    # 等待所有工作进程结束
    for idx, p in procs.items():
        print('JOIN process:', p)
        p.join()

    # 等待消费者线程结束
    consumer_thread.join()

    # 处理收集到的结果,比如写入文件
    print(f"共收集到 {len(collected_results)} 条结果")
    # 这里可以添加将collected_results写入result.output_file或output_text_file的逻辑

    print(f'Execution took {time.perf_counter() - start_time} seconds')

关键修改点说明

  • 移除无用的exit_cond相关代码,简化逻辑。
  • 修复updateResult方法,将结果存入队列而非仅打印。
  • 添加消费者线程,在主进程join工作进程的同时异步读取队列,避免子进程因队列满阻塞。
  • 消费者线程通过检查队列空且所有工作进程结束来退出,确保所有结果都被收集。

内容的提问来源于stack exchange,提问作者Hoang Duy Tran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 01:40:11