为何子进程未正常运行?multiprocessing队列异常排查
多进程队列处理异常排查:EOFError、BrokenPipe及输出队列无数据问题
问题场景
某应用用3个子进程并行处理数据:从输入队列取数,经process_the_data耗时计算后写入输出队列,最终由write_to_sql写入数据库。运行时出现EOFError、BrokenPipe错误,且输出队列从未有数据写入。
问题原因分析
子进程无终止逻辑,导致无限阻塞
process_the_data函数使用while True无限循环,没有退出条件。当主进程尝试join子进程时,子进程会一直阻塞在in_q.get()上。当主进程结束时,Manager创建的队列会被关闭,子进程此时再操作队列就会触发EOFError或BrokenPipe错误。输入队列填充逻辑完全错误
原代码中,遍历input_data的每个元素n时,会不断将同一个n放入输入队列直到队列满。这导致输入队列中全是重复的单个值,100个原始数据根本没被正常传入,子进程一直在处理重复数据,后续队列关闭时自然无法产生有效输出。输出队列处理时机不合理
主进程在填充单个n到队列满后才尝试读取输出队列,但此时子进程可能还未处理完数据,输出队列为空。同时因为输入队列被重复数据占满,子进程持续处理重复任务,没有机会产生可被读取的输出。
修复方案及代码
核心修复点
- 给子进程添加终止信号(如
None),收到信号后退出循环 - 修正输入队列填充逻辑,逐个传入原始数据
- 调整输出队列读取逻辑,确保所有处理结果都被读取
- 遵循多进程编程规范,添加
if __name__ == "__main__":保护
修复后的代码
import multiprocessing import random import time def process_the_data(in_q, out_q): while True: data = in_q.get(block=True) # 收到终止信号,退出循环 if data is None: break time.sleep(random.choice([1,2,3])) out_q.put(data + 1) def write_to_sql(n): print(n) if __name__ == "__main__": multiprocessing.set_start_method('fork') mgr = multiprocessing.Manager() input_q = mgr.Queue(maxsize=5) output_q = mgr.Queue() input_data = range(100) p1 = multiprocessing.Process(target=process_the_data, args=(input_q, output_q)) p2 = multiprocessing.Process(target=process_the_data, args=(input_q, output_q)) p3 = multiprocessing.Process(target=process_the_data, args=(input_q, output_q)) p1.start() p2.start() p3.start() # 正确填充所有输入数据 for n in input_data: input_q.put(n) # 给每个子进程发送终止信号,确保处理完数据后退出 for _ in range(3): input_q.put(None) # 循环读取输出队列,直到所有子进程退出且队列为空 while True: try: result = output_q.get(timeout=1) write_to_sql(result) except multiprocessing.queues.Empty: # 检查所有子进程是否已终止 if not p1.is_alive() and not p2.is_alive() and not p3.is_alive(): break p1.join() p2.join() p3.join() print('finished')
内容的提问来源于stack exchange,提问作者gph
相关产品推荐
相关产品推荐

