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

为何引入Pipe()后Python多进程无法并行?求解决方案

解决Pipe通信下的进程并行问题

嘿,我一眼就看到问题出在哪了——你现在根本没让funcb跑在独立子进程里,而且主进程过早的recv()调用直接把执行流程变成串行的了!

先给你上修正后的代码,保证能回到1秒左右的耗时,同时保留Pipe通信:

import numpy as np
import multiprocessing
import time

def funca(mylist, conn):
    time.sleep(1)
    mylist.append(666.6)
    conn.send(['a','a','a'])
    conn.close()  # 发完消息记得关连接,避免资源浪费

def funcb(mylist, conn):
    time.sleep(1)
    mylist.append(66.6)
    conn.send(['b','b','b'])
    conn.close()

if __name__ == "__main__":
    samples = [1,2,3]
    with multiprocessing.Manager() as MG:
        # 给两个进程各建一个管道,分别和主进程通信
        conn_a_parent, conn_a_child = multiprocessing.Pipe()
        conn_b_parent, conn_b_child = multiprocessing.Pipe()
        
        mylist = MG.list(samples)
        tic = time.time()
        
        # 启动两个独立子进程,各自拿自己的管道子端
        p1 = multiprocessing.Process(target=funca, args=(mylist, conn_a_child))
        p2 = multiprocessing.Process(target=funcb, args=(mylist, conn_b_child))
        
        p1.start()
        p2.start()
        
        # 等两个进程都启动后再接收消息,此时它们已经在后台并行跑了
        print(conn_a_parent.recv())
        print(conn_b_parent.recv())
        
        # 等待两个进程执行完毕
        p1.join()
        p2.join()
        
        print(list(mylist))
        toc = time.time()
        print('pass time = ', toc - tic)

关键问题拆解&优化点:

  • 必须给funcb创建子进程:你之前直接在主进程里调用funcb,相当于主进程先等funca发消息,收到后才会执行funcb,完全是串行执行,耗时自然翻倍。现在给funcb也创建Process实例,两个进程会同时启动并行运行。
  • 管道不能共享句柄:原代码里你把conn1同时传给p1和主进程里的funcb,这是不安全的——管道的一端只能被一个进程使用,否则会出现通信混乱。我们给每个子进程单独分配管道,主进程通过父端接收各自的消息。
  • 调整recv()的时机:先启动所有子进程,再执行接收操作。虽然recv()是阻塞的,但两个子进程已经在后台并行跑了,所以总耗时还是约1秒,而不是两个任务的时间相加。

如果你想省点事,用同一个管道让两个子进程给主进程发消息也可以,只是消息的接收顺序不确定(取决于哪个进程先跑完send()),代码大概是这样:

import numpy as np
import multiprocessing
import time

def funca(mylist, conn):
    time.sleep(1)
    mylist.append(666.6)
    conn.send(['a','a','a'])

def funcb(mylist, conn):
    time.sleep(1)
    mylist.append(66.6)
    conn.send(['b','b','b'])

if __name__ == "__main__":
    samples = [1,2,3]
    with multiprocessing.Manager() as MG:
        parent_conn, child_conn = multiprocessing.Pipe()
        mylist = MG.list(samples)
        tic = time.time()
        
        p1 = multiprocessing.Process(target=funca, args=(mylist, child_conn))
        p2 = multiprocessing.Process(target=funcb, args=(mylist, child_conn))
        
        p1.start()
        p2.start()
        
        # 接收两次,顺序可能随机
        print(parent_conn.recv())
        print(parent_conn.recv())
        
        p1.join()
        p2.join()
        
        print(list(mylist))
        toc = time.time()
        print('pass time = ', toc - tic)

这个版本也能保持并行,耗时同样约1秒,只是打印的消息顺序可能会变,看你需求选就行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 04:49:51