使用multiprocessing.Pool传递Pandas DataFrame时速度异常缓慢的原因
我尝试用multiprocessing.Pool并行创建约100MB级别的Pandas DataFrame,结果发现虽然处理环节并行执行,但IPC(进程间通信)开销导致整体速度远慢于单进程。更奇怪的是,把DataFrame序列化到磁盘再读取的方案反而更快——理论上内存内传递应该比磁盘操作快才对。
以下是测试两种方案及对比纯字节数组的结果和脚本,手动序列化到磁盘的方式始终快得多。
测试结果
Mac OS Ventura(Python 3.10)
Pandas frame: 56.739380359998904 Bytes: 100.31256382900756 Pandas hard disk: 2.5154516759794205
同一Mac的Docker容器内(Python 3.10)
Pandas frame: 30.639809517015237 Bytes: 17.284717564994935 Pandas hard disk: 4.386503675952554
测试脚本
import multiprocessing import pandas as pd import numpy as np import time def creator(id_): # ~230 Mb in size return pd.DataFrame(np.ones((3000, 10000))) def creator_bytes(id_): return bytearray(3000 * 10000 * 8) def creator_pkl(id_): filename = f"pickled_{id_}.pkl" pd.DataFrame(np.ones((3000, 10000))).to_pickle(filename) return filename if __name__ == "__main__": timestamp = time.perf_counter() fake_input = list(range(10)) with multiprocessing.Pool(1) as p: results = p.map(creator, fake_input) print(f"Pandas frame: {time.perf_counter() - timestamp}") ### timestamp = time.perf_counter() with multiprocessing.Pool(1) as p: results = p.map(creator_bytes, fake_input) print(f"Byte array: {time.perf_counter() - timestamp}") ### timestamp = time.perf_counter() with multiprocessing.Pool(1) as p: results = p.map(creator_pkl, fake_input) print(f"Pandas hard disk: {time.perf_counter() - timestamp}")
默认序列化机制的低效性:
multiprocessing默认使用的pickle(或兼容序列化器)对Pandas DataFrame处理效率极低。DataFrame不仅包含原始数值数据,还有大量元数据(索引、列类型、内存布局信息等),默认pickle会递归序列化所有细节;而Pandas自带的to_pickle做了针对性优化——比如直接批量序列化底层NumPy数组,跳过冗余的对象封装层,序列化速度快得多。IPC的内存拷贝开销:多进程间传递数据时,操作系统需要把序列化后的完整数据从子进程内存复制到父进程内存,这一步的拷贝开销对于大对象来说非常可观。而磁盘方案里,子进程只返回一个极小的文件名字符串,几乎没有IPC开销;序列化和后续反序列化都在本地进程完成,避免了跨进程的大内存拷贝。
大对象传递的固有成本:测试里纯bytearray的耗时也很高,说明不管对象类型是什么,只要体积大,跨进程传递的序列化+拷贝时间都会远超现代SSD的磁盘IO速度。现代NVMe SSD的写入速度轻松达到GB/s级别,而跨进程的IPC拷贝受限于内存总线和进程间通信机制,实际吞吐量反而更低。
内容的提问来源于stack exchange,提问作者Yourstruly

