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

使用multiprocessing.Pipe.send()时出现BrokenPipeError的解决咨询

问题:多进程Pipe通信触发BrokenPipeError

我开发了一款调用在线OCR API的程序,单张图片处理耗时2-5秒。为避免用户等待所有图片处理完成,用multiprocessing实现多进程:先处理第一张图片,其余图片在子进程并行处理,依赖multiprocessing.Pipe()完成进程间通信。

核心代码如下:

import multiprocessing as mp
# importing cv2, PIL, os, json, other stuff

def image_processor():
    # 先处理列表中第一张图,剩余图片交给子进程处理
    p_conn, c_conn = mp.Pipe()
    p = mp.Process(target=Processing.worker, args=([c_conn, images, path], 5))
    p.start()
    
    while True:
        out = p_conn.recv()
        if not out:
            break
        else:
            im_data.append(out)
            p_conn.send(True)


class Processing:
    def worker(data, mode, headers=0):
        # 其他分支逻辑省略
        elif mode == 5:
            print(data[0])
            for im_name in data[1]:
                if data[1].index(im_name) != 0:
                    im_path = f'{data[2]}\{im_name}'  # 获取图片路径
                    im = pil_img.open(im_path).convert('L')  # PIL打开并转为灰度图
                    os.rename(im_path, f'{data[2]}\Archive\{im_name}')  # 原图移至归档目录
                    im_grayscale = f'{data[2]}\g_{im_name}'  # 灰度图保存路径
                    im.save(im_grayscale)  # 保存灰度图
                        
                    ocr_data = json.loads(bl.Visual.OCR.ocr_space_file(im_grayscale)).get('ParsedResults')[0].get('ParsedText').splitlines()
                    print(ocr_data)
                    data[0].send([im_name, f'{data[2]}\Archive\{im_name}', ocr_data])
                    data[0].recv()
            
            data[0].send(False)

运行后触发报错:

Process Process-1:
Traceback (most recent call last):
  File "C:\Users\BruhK\AppData\Local\Programs\Python\Python310\lib\multiprocessing\process.py", line 315, in _bootstrap
    self.run()
  File "C:\Users\BruhK\AppData\Local\Programs\Python\Python310\lib\multiprocessing\process.py", line 108, in run
    self._target(*self._args, **self._kwargs)
  File "c:\Users\BruhK\PycharmProjects\pythonProject\FleetFeet-OCR-Final.py", line 275, in worker
    data[0].send([{im_name}, f'{data[2]}\{im_name}', ocr_data])
  File "C:\Users\BruhK\AppData\Local\Programs\Python\Python310\lib\multiprocessing\connection.py", line 211, in send
    self._send_bytes(_ForkingPickler.dumps(obj))
  File "C:\Users\BruhK\AppData\Local\Programs\Python\Python310\lib\multiprocessing\connection.py", line 285, in _send_bytes
    ov, err = _winapi.WriteFile(self._handle, buf, overlapped=True)
BrokenPipeError: [WinError 232] The pipe is being closed

注:测试环境中父子进程可正常收发二维/三维数组,测试代码如下:

import multiprocessing as mp
import random
import time


def hang(p):
  hang_time = random.randint(1, 5)
  time.sleep(hang_time)
  print(p)
  p.send(hang_time)
  time.sleep(1)


class Child:
  def process():
    start = time.time()
    p_conn, c_conn = mp.Pipe()
    p = mp.Process(target=hang, args=(c_conn,))
    p.start()
    out = p_conn.recv()
    print(f'Waited for {time.time() - start}')
    p.join()
    print(f'New time: {time.time() - start}')
    return out


class Parent:
  def run():
    # 其他逻辑省略
    print(f'Hang time: {Child.process()}')
    # 其他逻辑省略


if __name__ == '__main__':
  Parent.run()

解决思路与方案

1. 修正Pipe连接的传递方式

Windows系统中multiprocessing采用spawn模式创建进程,参数需支持pickle序列化。你将c_conn嵌套在列表中传递,可能导致连接序列化/反序列化异常,进而使管道提前关闭。

修改方式:将Pipe连接作为独立参数传递,不要嵌套在列表中:

# 原代码
p = mp.Process(target=Processing.worker, args=([c_conn, images, path], 5))
# 修改后
p = mp.Process(target=Processing.worker, args=(c_conn, images, path, 5))

同步调整worker函数的参数定义:

class Processing:
    def worker(c_conn, images, path, mode, headers=0):
        # 其他分支逻辑省略
        elif mode == 5:
            print(c_conn)
            for im_name in images:
                if images.index(im_name) != 0:
                    im_path = f'{path}\{im_name}'
                    # 其余代码保持不变,将data[0]替换为c_conn,data[1]替换为images,data[2]替换为path
                    c_conn.send([im_name, f'{path}\Archive\{im_name}', ocr_data])
                    c_conn.recv()
            
            c_conn.send(False)

2. 优化同步逻辑,避免不必要的阻塞

当前子进程每发送一次数据后,会调用data[0].recv()等待父进程的True响应,这种强同步逻辑可能导致管道超时或异常断开。

  • 如果不需要父进程确认,直接删除子进程中的recv()调用;
  • 如果必须保留确认机制,确保父进程的send()逻辑不会被其他操作阻塞,能及时响应子进程。

3. 确保父进程不会提前退出

检查image_processor()函数执行完毕后,父进程是否会直接退出,导致Pipe被强制关闭。可以在循环结束后调用p.join()等待子进程完成:

def image_processor():
    p_conn, c_conn = mp.Pipe()
    p = mp.Process(target=Processing.worker, args=(c_conn, images, path, 5))
    p.start()
    
    while True:
        out = p_conn.recv()
        if not out:
            break
        else:
            im_data.append(out)
            # 若无需确认,可删除下面的send语句
            p_conn.send(True)
    # 等待子进程执行完毕
    p.join()

4. 验证数据的可序列化性

虽然测试中数组能正常传递,但实际场景中的ocr_data可能包含无法被pickle序列化的对象。可以:

  • 确保ocr_data是纯Python基础类型(如字符串列表);
  • 或者将ocr_data用json.dumps()转为字符串后发送,接收时再用json.loads()解析。

内容的提问来源于stack exchange,提问作者BappoHotel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 13:01:09