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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 10:27:06