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

如何解决Dask DataFrame分组自定义计算时磁盘占用过高问题

解决Dask GroupBy Apply处理大CSV时磁盘空间溢出的问题

我完全理解你现在的困扰——用Dask处理15-20GB的CSV文件时,为了计算Dask原生不支持的峰度(kurtosis)和偏度(skew),不得不使用groupby().apply(),结果临时文件直接暴涨到150GB,触发磁盘空间不足导致脚本崩溃。别着急,咱们可以从几个方向优化代码,避免大量临时文件生成:

方案1:用agg结合自定义聚合函数(最推荐)

Dask的groupby.agg()其实支持自定义统计函数,只要你正确指定meta信息。相比apply(),agg()是按列独立处理的,Dask能更智能地划分任务、减少中间数据写入,从根源上降低临时文件的产生。

改写后的代码如下:

import dask.dataframe as dd
import numpy as np
import scipy.stats as sps
import pandas as pd

# 先定义需要的自定义统计函数
def pearson_kurt(col):
    # 调整为皮尔逊峰度(scipy默认返回费希尔峰度,加3转换)
    return sps.kurtosis(col) + 3

def calculate_skew(col):
    return sps.skew(col)

# 读取大CSV时可指定blocksize,减少分区数量(根据内存调整)
ddf = dd.read_csv('18_GB_csv_file.csv', blocksize='2GB')

segmentations = {
 'seg1' : ['col1', 'col2'],
 'seg2' : ['col1', 'col2', 'col3', 'col4'],
 'seg3' : ['col3', 'col4'],
 'seg4' : ['col1', 'col2', 'col5']
}

data_cols = [ 'datacol1', 'datacol2', 'datacol3' ]

# 构建聚合字典:每个数据列需要计算的统计量
agg_dict = {
    col: ['mean', 'std', 'min', 'max', pearson_kurt, calculate_skew]
    for col in data_cols
}

dd_comp = {}
for seg_group, seg_cols in segmentations.items():
    # 手动构造meta,避免Dask自动扫描数据(节省时间和资源)
    meta_index = pd.MultiIndex.from_product(
        [data_cols, ['mean', 'std', 'min', 'max', 'pearson_kurt', 'calculate_skew']]
    )
    meta = pd.DataFrame(dtype='float64', index=meta_index).T
    
    # 使用agg替代apply,性能和磁盘占用都会优化很多
    df_grouped = ddf.groupby(seg_cols).agg(agg_dict, meta=meta)
    dd_comp[seg_group] = df_grouped

with dd.ProgressBar():
    segmented_stats = dd.compute(dd_comp)

为什么这招管用?

  • agg()针对单列处理,不需要把整个分组的所有列数据都加载到内存/磁盘,中间数据量大幅减少;
  • 手动指定meta让Dask不用提前扫描数据集推断结构,避免额外的IO开销;
  • 调整blocksize减少分区数量,降低任务数和对应的临时文件数量。

方案2:优化Dask临时文件配置

如果方案1还不够,你可以调整Dask的临时文件相关设置,进一步减少磁盘压力:

  • 更换临时目录到更大磁盘:
    import dask.config
    dask.config.set(temporary_directory='/path/to/your/larger/disk')
    
  • 启用临时文件自动清理:
    Dask默认会在任务成功完成后清理临时文件,但异常崩溃时可能残留。可以手动配置:
    dask.config.set({'temporary-directory': {'cleanup': 'on-success'}})
    
  • 减少分区数量:
    像方案1里那样,在read_csv时设置合适的blocksize(比如2GB-4GB,根据你的内存大小调整),避免过多小分区导致大量临时文件。

方案3:拆分任务,减少apply的负载

如果你因为某些原因必须使用apply(),可以拆分统计任务:先用agg()计算Dask原生支持的统计量(mean、std等),再用apply()只计算峰度和偏度,最后合并结果。这样能大幅减少apply()处理的数据量:

import dask.dataframe as dd
import numpy as np
import scipy.stats as sps
import pandas as pd

ddf = dd.read_csv('18_GB_csv_file.csv', blocksize='2GB')

segmentations = {
 'seg1' : ['col1', 'col2'],
 'seg2' : ['col1', 'col2', 'col3', 'col4'],
 'seg3' : ['col3', 'col4'],
 'seg4' : ['col1', 'col2', 'col5']
}

data_cols = [ 'datacol1', 'datacol2', 'datacol3' ]

# 第一步:用agg计算基础统计量(无磁盘压力)
base_agg_dict = {col: ['mean', 'std', 'min', 'max'] for col in data_cols}
base_comp = {}
for seg_group, seg_cols in segmentations.items():
    base_comp[seg_group] = ddf.groupby(seg_cols).agg(base_agg_dict)

# 第二步:用apply只计算峰度和偏度
def kurt_skew_agg(segment):
    res = pd.DataFrame()
    for col in data_cols:
        res[(col, 'kurt')] = [sps.kurtosis(segment[col]) + 3]
        res[(col, 'skew')] = [sps.skew(segment[col])]
    return res

ks_comp = {}
for seg_group, seg_cols in segmentations.items():
    # 构造meta
    meta_cols = pd.MultiIndex.from_product([data_cols, ['kurt', 'skew']])
    meta = pd.DataFrame(dtype='float64', columns=meta_cols)
    ks_comp[seg_group] = ddf.groupby(seg_cols)[data_cols].apply(kurt_skew_agg, meta=meta)

# 第三步:合并结果
with dd.ProgressBar():
    base_stats = dd.compute(base_comp)
    ks_stats = dd.compute(ks_comp)

final_stats = {}
for seg_name in segmentations.keys():
    final_stats[seg_name] = pd.concat(
        [base_stats[0][seg_name], ks_stats[0][seg_name]],
        axis=1
    )

这个方法的核心是把重负载的apply()任务拆分成只处理必要的统计量,减少中间数据的写入量。

内容的提问来源于stack exchange,提问作者Gijs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 08:06:20