如何用multiprocessing.Pipe/Queue获取大对象?异步+多进程是否可行?
解决multiprocessing传输大对象阻塞及异步+多进程方案优化
一、Pipe传输大对象卡住的原因及解决方法
Pipe默认创建的是双向管道,内部缓冲区存在大小限制。当子进程调用send()发送大对象(比如numpy数组)时,如果父进程未及时调用recv()读取数据,管道缓冲区被填满后,子进程会阻塞在send()调用处,无法继续执行或退出。
解决方法:
调整接收时序:确保父进程在子进程发送数据前,就处于等待接收的状态。比如启动子进程后立刻调用
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)改用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)使用共享内存传输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
相关产品推荐
相关产品推荐

