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

为何子进程未正常运行?multiprocessing队列异常排查

多进程队列处理异常排查:EOFError、BrokenPipe及输出队列无数据问题

问题场景

某应用用3个子进程并行处理数据:从输入队列取数,经process_the_data耗时计算后写入输出队列,最终由write_to_sql写入数据库。运行时出现EOFError、BrokenPipe错误,且输出队列从未有数据写入。

问题原因分析

  1. 子进程无终止逻辑,导致无限阻塞
    process_the_data函数使用while True无限循环,没有退出条件。当主进程尝试join子进程时,子进程会一直阻塞在in_q.get()上。当主进程结束时,Manager创建的队列会被关闭,子进程此时再操作队列就会触发EOFError或BrokenPipe错误。

  2. 输入队列填充逻辑完全错误
    原代码中,遍历input_data的每个元素n时,会不断将同一个n放入输入队列直到队列满。这导致输入队列中全是重复的单个值,100个原始数据根本没被正常传入,子进程一直在处理重复数据,后续队列关闭时自然无法产生有效输出。

  3. 输出队列处理时机不合理
    主进程在填充单个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 15:13:11