多进程多线程架构下结果获取及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')
问题原因
程序挂起的核心原因是死锁:
- 子进程调用
queue.put(return_values)时,若队列被填满,子进程会阻塞等待队列释放空间。 - 主进程在调用
p.join()前未读取队列数据,导致子进程一直卡在put操作,主进程则等待子进程结束,形成双向阻塞。 - 代码中
exit_cond变量未在子进程的run方法中使用,设置exit_cond.set()无法终止子进程,属于冗余代码。
解决方案
- 主进程提前消费队列数据:在join子进程前,启动线程异步读取队列,避免子进程因队列满阻塞。
- 修复
updateResult方法:替换打印逻辑为队列存入操作。 - 移除无效的
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
相关产品推荐
相关产品推荐

