Python中从多文件读取大数据并聚合的最快实现方式是什么?
多CSV文件批量读取聚合效率优化方案
适用Jupyter环境的优化实现方案
CSV读取的瓶颈以IO密集型为主,CPU解析开销为辅,可按场景选择对应方案:
- 优先选择多线程+分块IO方案,完全规避进程间序列化开销:CSV读取过程中IO等待占比很高,GIL会在IO等待时自动释放,多线程足以跑满磁盘/网络带宽,不需要引入多进程。示例代码如下:
import pandas as pd from concurrent.futures import ThreadPoolExecutor import os # 按需生成所有CSV文件路径 file_list = [f"{i}.csv" for i in range(1, N+1) if os.path.exists(f"{i}.csv")] def read_csv(file_path): # 提前指定dtype可减少类型推断开销,同时降低内存占用 return pd.read_csv(file_path, dtype={"列名1": str, "列名2": int}) # 线程数可按存储介质调整:SSD建议设为8~32,机械硬盘不超过8避免随机IO开销过高 with ThreadPoolExecutor(max_workers=16) as executor: dfs = list(executor.map(read_csv, file_list)) # 合并所有DataFrame df_total = pd.concat(dfs, ignore_index=True)
- 如果单CSV文件体积大于1GB,字符串解析、类型推断等CPU开销占比很高,确需使用多进程时,不要将大DataFrame作为参数传入子进程,仅将文件路径字符串传给子进程,子进程自行读取文件返回结果,此时序列化的只有短字符串,开销可忽略不计。
Queue、joblib、Ray三种方案的效率对比
针对你描述的场景(仅传入文件路径、返回单个DataFrame),三种方案的性能表现如下:
- 原生Queue:效率最低。原生multiprocessing的Queue使用pickle序列化,没有针对科学计算对象做优化,返回大量DataFrame时序列化/反序列化开销极高,且需要自行处理进程管理、异常捕获,开发成本高,不推荐使用。
- joblib:性价比最高。默认使用loky后端,针对numpy数组、pandas对象做了专门的序列化优化,速度比原生pickle高3~10倍,且API封装完善,不需要手动管理进程、队列,Jupyter环境兼容性好,适合单机内存可以容纳所有聚合数据的场景。示例代码如下:
from joblib import Parallel, delayed # n_jobs设为CPU核心数即可 dfs = Parallel(n_jobs=8, verbose=10)(delayed(read_csv)(f) for f in file_list) df_total = pd.concat(dfs, ignore_index=True)
- Ray:峰值性能最高但启动开销大。Ray是分布式计算框架,使用自研的Plasma共享内存和Apache Arrow格式传递数据,大型pandas对象传递时不需要序列化/反序列化,进程间数据传递开销几乎为0。但框架本身初始化需要启动集群服务,单机场景下启动开销比joblib高1~2秒,更适合总数据量超过单机内存、或后续需要进行分布式计算的场景,仅做CSV聚合的话属于大材小用。Jupyter环境下使用Ray需要先执行初始化:
ray.init(ignore_reinit_error=True),任务函数需要加@ray.remote装饰器。
补充说明:你之前的经验是正确的,如果将大体积DataFrame作为参数传入子进程,任何方案的序列化开销都会很高,但你的场景完全可以规避该操作,仅返回处理后的结果时,优化后的序列化方案开销会远低于多进程带来的CPU解析加速收益。
内容的提问来源于stack exchange,提问作者YNX
相关产品推荐
相关产品推荐

