Dask调用dd.read_parquet内存溢出但dd.from_pandas运行正常原因求助
问题原因分析
1. 不必要的shuffle操作触发全量数据加载
你在dd.read_parquet后直接调用repartition,会触发全量数据的shuffle重分区操作。Dask的执行是基于计算图的,哪怕你后续只调用head()取前几行,也需要先完成repartition的全量计算步骤,导致worker需要一次性加载所有解压缩后的parquet数据。
源文件磁盘占用914MB是压缩后的大小,加载到内存解压缩后体积会膨胀3~5倍,刚好接近你单worker 7.8GB的内存阈值,触发Dask的内存保护机制杀掉worker。
而用pd.read_parquet+dd.from_pandas的流程没有shuffle操作:pandas先把全量数据加载到驱动进程内存,dd.from_pandas只是在内存中对数据做逻辑切分,不需要跨worker传输数据,head()也只需要取第一个分区的前几行,所以不会触发内存超限。
2. 单文件parquet的默认分区问题
如果你的源数据是单个parquet文件,dd.read_parquet默认只会生成1个分区,你直接重分区到100份的操作会强制把单分区全量数据打散,进一步加剧内存占用。
解决方案
- 读文件时直接指定分区大小,避免额外repartition
不需要单独调用repartition,在dd.read_parquet时通过blocksize参数控制每个分区的大小,自动生成符合要求的分区,不会触发shuffle:
# 按每个分区10MB大小自动拆分,不需要额外repartition df1 = dd.read_parquet('..', blocksize='10MB') df1.head()
- 调整worker内存参数避免误杀
初始化集群时调整内存阈值,让worker在内存占用过高时先溢写磁盘,而不是直接被杀掉:
client = Client( n_workers=4, memory_limit="8GB", memory_target_fraction=0.8, # 内存占比到80%开始溢写磁盘 memory_spill_fraction=0.9, memory_pause_fraction=0.95 )
- 仅在必要时做重分区
如果确实需要调整分区数,先把数据持久化到内存/磁盘后再做重分区,避免每次计算都触发shuffle。
dd.read_parquet的存在意义
它的核心使用场景是处理远大于单机内存的数据集:当你的数据量达到几十GB、上TB级别时,pandas无法一次性把全量数据加载到内存,dd.read_parquet可以按需分批加载分片计算,全程不需要占满整机内存,完全可以处理超过物理内存上限的数据量。你当前的数据量较小,所以pandas加载的方案可行,换成超大数据量后就能体现dd.read_parquet的价值。
内容的提问来源于stack exchange,提问作者Tim
相关产品推荐
相关产品推荐

