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

为何multiprocessing.Pipe传输大numpy数组反而慢于Queue?

为什么传输大型numpy数组时multiprocessing.Pipe比Queue慢?

我原本以为multiprocessing.Queue内部是基于Pipe实现的,所以Pipe应该有更快的传输速度,但在测试大型numpy数组的传输时,结果却相反。以下是我的测试代码和结果,想请教一下问题出在哪里?

Pipe测试代码

import sys
import time
from multiprocessing import Process, Pipe
import numpy as np
NUM = 1000
def worker(conn):
    for task_nbr in range(NUM):
        conn.send(np.random.rand(400, 400, 3))
    sys.exit(1)
def main():
    parent_conn, child_conn = Pipe(duplex=False)
    Process(target=worker, args=(child_conn,)).start()
    for num in range(NUM):
        message = parent_conn.recv()
if __name__ == "__main__":
    start_time = time.time()
    main()
    end_time = time.time()
    duration = end_time - start_time
    msg_per_sec = NUM / duration
    print "Duration: %s" % duration
    print "Messages Per Second: %s" % msg_per_sec

耗时10.86秒

Queue测试代码

import sys
import time
from multiprocessing import Process
from multiprocessing import Queue
import numpy as np
NUM = 1000
def worker(q):
    for task_nbr in range(NUM):
        q.put(np.random.rand(400, 400, 3))
    sys.exit(1)
def main():
    recv_q = Queue()
    Process(target=worker, args=(recv_q,)).start()
    for num in range(NUM):
        message = recv_q.get()
if __name__ == "__main__":
    start_time = time.time()
    main()
    end_time = time.time()
    duration = end_time - start_time
    msg_per_sec = NUM / duration
    print "Duration: %s" % duration
    print "Messages Per Second: %s" % msg_per_sec

耗时6.86秒


这问题挺有意思的——虽然咱们都知道multiprocessing.Queue底层是基于Pipe实现的,但在你测试的大型numpy数组场景下,Queue反而更快,主要是因为Queue做了几个关键优化,刚好适配了大对象传输的需求:

1. Queue的缓冲机制让传输更“流畅”

Pipe的send()在单向模式下是阻塞的:当管道内部缓冲区被填满后,发送数据的worker进程会直接卡住,直到主进程调用recv()把数据取走。在你的测试里,worker每生成一个400x400x3的数组就会停下来等主进程接收,相当于“生成等传输,传输等接收”的串行循环,时间都耗在等待上了。

而Queue内部自带缓冲区(默认大小可以用maxsize调),put()方法在缓冲区没满的时候不会阻塞,worker可以一口气生成好几个数组都放进队列里,主进程的get()能跟着自己的节奏取数据,相当于流水线并行,减少了进程间互相等待的开销。

2. 大对象序列化的针对性优化

numpy数组是复杂对象,进程间传输必须序列化(默认用pickle)。虽然Pipe和Queue都用pickle,但Queue底层对支持缓冲协议的对象(比如numpy数组)做了特殊处理——会走更高效的序列化路径,甚至能做到零拷贝,避免了多余的数据复制。

反观Pipe的send()方法,对numpy数组的处理就比较“直白”,没有针对大对象做优化,序列化和传输的开销自然更高。

3. 进程调度的开销差异

虽然Queue因为要支持多生产者/多消费者,会有额外的锁,但在你这种单生产者单消费者的场景下,锁的影响几乎可以忽略。反而Pipe的频繁阻塞会导致进程来回切换(worker阻塞时主进程跑,主进程等数据时worker跑),增加了上下文切换的开销;而Queue的缓冲机制让进程切换更少,整体效率自然上来了。

可以试试这些验证方法

如果你想确认这些原因,可以做几个小测试:

  • 给Pipe的测试改一改,比如一次发送多个数组打包成列表,减少阻塞的次数,看看性能会不会提升;
  • 用multiprocessing.connection的send_bytes()和recv_bytes()直接传输numpy数组的字节数据(比如arr.tobytes()),绕开pickle,对比下速度;
  • 把Queue的maxsize设为1,让它和Pipe一样变成同步传输,这时候两者的性能应该会接近很多。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:49:20