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

如何用multiprocessing.Pipe/Queue获取大对象?异步+多进程是否可行?

解决multiprocessing传输大对象阻塞及异步+多进程方案优化

一、Pipe传输大对象卡住的原因及解决方法

Pipe默认创建的是双向管道,内部缓冲区存在大小限制。当子进程调用send()发送大对象(比如numpy数组)时,如果父进程未及时调用recv()读取数据,管道缓冲区被填满后,子进程会阻塞在send()调用处,无法继续执行或退出。

解决方法:

  1. 调整接收时序:确保父进程在子进程发送数据前,就处于等待接收的状态。比如启动子进程后立刻调用parent_conn.recv(),而非等子进程运行完毕再接收。
    示例代码:

    import multiprocessing
    import numpy as np
    
    def worker(conn):
        big_data = np.ones([100, 100])
        conn.send(big_data)
        conn.close()
    
    if __name__ == "__main__":
        parent_conn, child_conn = multiprocessing.Pipe()
        p = multiprocessing.Process(target=worker, args=(child_conn,))
        p.start()
        # 提前启动接收,避免子进程send阻塞
        result = parent_conn.recv()
        p.join()
        print(result.shape)
    
  2. 改用multiprocessing.Queue:Queue内部维护了后台线程负责数据传输,自带缓冲机制,能避免缓冲区满导致的阻塞。即使父进程未立即接收,子进程的put()操作也不会轻易阻塞(除非Queue最大容量被填满)。
    示例代码:

    import multiprocessing
    import numpy as np
    
    def worker(q):
        big_data = np.ones([100, 100])
        q.put(big_data)
    
    if __name__ == "__main__":
        q = multiprocessing.Queue()
        p = multiprocessing.Process(target=worker, args=(q,))
        p.start()
        result = q.get()
        p.join()
        print(result.shape)
    
  3. 使用共享内存传输numpy数组:超大numpy数组的序列化/反序列化会带来性能损耗,甚至触发管道阻塞。可以用multiprocessing.Array或numpy共享内存机制直接在进程间共享数据,避免完整对象拷贝。
    示例代码(numpy共享内存):

    import multiprocessing
    import numpy as np
    from multiprocessing import shared_memory
    
    def worker(shm_name, shape, dtype):
        existing_shm = shared_memory.SharedMemory(name=shm_name)
        arr = np.ndarray(shape, dtype=dtype, buffer=existing_shm.buf)
        arr[:] = np.ones(shape)
        existing_shm.close()
    
    if __name__ == "__main__":
        shape = (100, 100)
        dtype = np.float64
        shm = shared_memory.SharedMemory(create=True, size=np.prod(shape)*dtype.itemsize)
        arr = np.ndarray(shape, dtype=dtype, buffer=shm.buf)
        
        p = multiprocessing.Process(target=worker, args=(shm.name, shape, dtype))
        p.start()
        p.join()
        
        print(arr)
        shm.close()
        shm.unlink()
    

二、异步+多进程方案的合理性及优化

你的异步+多进程方案本身没有问题,但需要注意异步事件循环与多进程操作的交互逻辑:

  • 不能在异步函数中直接阻塞调用parent_conn.recv()或q.get(),否则会阻塞整个事件循环。应该用asyncio.to_thread()(Python 3.9+)或loop.run_in_executor()将同步接收操作放到线程池中执行。
  • 内存监控的异步函数要独立于数据接收流程,避免因监控逻辑占用资源导致接收不及时。

优化后的异步示例:

import asyncio
import multiprocessing
import numpy as np
import psutil

def worker(q):
    big_data = np.ones([100, 100])
    q.put(big_data)

async def monitor_process(p):
    while p.is_alive():
        mem_info = psutil.Process(p.pid).memory_info()
        print(f"Process {p.pid} memory usage: {mem_info.rss / 1024 / 1024:.2f} MB")
        await asyncio.sleep(0.1)

async def main():
    q = multiprocessing.Queue()
    p = multiprocessing.Process(target=worker, args=(q,))
    p.start()
    
    # 启动内存监控任务
    monitor_task = asyncio.create_task(monitor_process(p))
    
    # 在线程中执行Queue.get(),避免阻塞事件循环
    result = await asyncio.to_thread(q.get)
    
    await monitor_task
    p.join()
    print(result.shape)

if __name__ == "__main__":
    asyncio.run(main())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 21:22:49