为何引入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
相关产品推荐
相关产品推荐

