Dask保存parquet时占满内存且运行速度慢于pandas问题求助
Dask内存占用过高、运行速度慢问题解决方案
问题1:数据处理速度远低于pandas,compute卡顿
核心原因
- 当前使用的
threading调度器无法绕过pandas操作的GIL全局解释器锁,所有CPU密集型操作实际串行执行,并行优势完全无法发挥,速度自然低于直接用pandas。 - 直接调用
compute会把选中的date、permno两列全量加载到单机内存,40G数据集对应两列总容量至少数G,序列化、数据传输开销极大。 - 100个分区对应单分区大小约400M,若上游处理存在join、groupby等shuffle操作,极易出现数据倾斜,部分分区实际大小远超平均值,拖慢整体处理速度。
- 全量float64类型空间占用过高,100余列float64的单条数据内存占用超过800字节,数据拷贝、计算的开销都会显著上升。
修复方案
- 调度器替换为
processes或者本地分布式调度器,修改compute代码为output = df[["date", "permno"]].compute(scheduler='processes'),绕开GIL限制。 - 除非必须拿到全量数据到单机处理,否则尽量保留数据在Dask DataFrame中完成所有计算,仅最终导出小结果时再调用compute。
- 做数据类型压缩:将不需要高精度的float64列转换为float32,可直接减少一半内存占用;
date列转换为datetime类型,permno若为整数则转换为int32/int64类型,进一步降低空间占用。 - 调整分区大小,执行
df = df.repartition(npartitions=200)将单分区大小控制在100-200M区间,降低单分区处理的内存压力,提升并行效率。
问题2:写parquet时内存不足,buff/cache飙升
核心原因
- fastparquet写入时默认会进行行组合并、全量列统计信息收集,单分区写入的实际内存开销是分区本身大小的2-3倍,大分区场景下内存占用极易失控。
- 系统page cache会自动缓存正在写入的parquet文件,几十G的文件写入会直接把buff/cache占满,这部分内存虽然可回收,但叠加Dask本身的处理内存占用就会触发OOM。
- 100个分区默认会并行写入,每个写入进程独立持有写入缓冲、元数据内存,叠加后的总内存占用很容易超过110G的服务器内存上限。
修复方案
- 替换写入引擎为pyarrow,pyarrow的parquet写入内存控制、速度都远优于fastparquet,修改代码为
df.to_parquet('my_data_frame', engine="pyarrow")即可。 - 写入时限制并行度,添加参数控制同时写入的进程数:
df.to_parquet('my_data_frame', engine="pyarrow", compute_kwargs={"scheduler": "processes", "num_workers": 8}),避免多进程内存叠加超限。 - 不需要按列做过滤查询的话,添加
write_statistics=False参数关闭统计信息收集,大幅降低元数据处理的内存开销。 - 若buff/cache占用持续过高,写入完成后可手动执行系统命令回收缓存,正常场景下系统会在内存不足时自动回收这部分空间,无需额外处理。
内容的提问来源于stack exchange,提问作者J.Ewa
相关产品推荐
相关产品推荐

