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

使用Dask DataFrame时临时文件过度累积问题求助

Dask外连接与聚合时临时文件暴涨问题解决

问题背景

原数据约50GB,仅处理1/10数据时临时文件体积就达到750GB。代码逻辑为从5个文件夹读取数据,调整4个文件夹的列名后与剩余文件夹做外连接,经字段处理、聚合后导出CSV文件。原代码如下:

import pandas as pd
import glob
import os
import numpy as np
import dask.dataframe as dd
from dask.diagnostics import ProgressBar
from dask.distributed import Client, LocalCluster

def rpt_dask_dataframe(folder_path, using_cols=None):
    all_files = glob.glob(os.path.join(folder_path, "*.rpt"))

    # Check if the folder is empty
    if len(all_files) == 0:
        raise ValueError("No files found in the folder")
    # Check file header row to get the column names.
    # the header start with '!' or '&' the last one is the header row
    header_row = None
    for file in all_files:
        with open(file, 'r') as f:
            for line in f:
                if line.startswith('!') or line.startswith('&'):
                    header_row = line
                    break
        if header_row is not None:
            break

    if using_cols is None:
        dask_df = dd.read_csv(all_files)
    else:
        dask_df = dd.read_csv(all_files, usecols=using_cols)
    return dask_df

def unique_list(l):
    x = []
    for a in l:
        if a not in x:
            x.append(a)
    return x


folder_path = r'C:/Users/USER/Desktop'

folder_list = os.listdir(folder_path)

folder_dict = {'BF_EB': 1, 'BF_NB1': 2, 'BF_NB2': 3, 'BF_NB3': 4, 'AF_EB': 5}

folder_list = [os.path.join(folder_path, folder) for folder in folder_list if folder in folder_dict]

result_folder = r'C:/Users/USER/Desktop/result'

if not os.path.exists(result_folder):
    os.makedirs(result_folder)

reference_folder = r'C:/Users/USER/Desktop/ref_info.fac'

reference_df = dd.read_csv(reference_folder, usecols=[1, 2], dtype=str)

reference_df = reference_df.compute()

pk_list = ['pk1', 'pk2', 'pk3', 'pk4', 'pk5', 'pk6', 'pk7', 'pk8']
fk_list = ['fk1', 'fk2']
value1_list = ['value1', 'value2', 'value3', 'value4', 'value5']
value2_list = ['value4', 'value5', 'value6', 'value7', 'value8']
main_list = pk_list + fk_list + value1_list

result_col_list = ['ref1', 'pk1', 'pk2', 'pk3', 'pk4', 'fk1', 'fk2', 'value4', 'value4_BF', 'value4_AF', 'Cvalue1_BF', 'Cvalue1_AF',
                   'value3_BF', 'value5', 'value5_BF', 'value5_AF', 'value6', 'value6_BF', 'value6_AF', 'value6_DEL', 'value6_ADD',
                   'value7_BF', 'value7_AF', 'value8_BF', 'value8_AF', 'value9_BF', 'value9_AF', 'value10_BF', 'value10_AF']

summary_columns = ['rsf1', 'pk1', 'pk2', 'pk3', 'fk1', 'fk2']
summary_key_num = len(summary_columns)

#usecols & rename
df_BF_EB = rpt_dask_dataframe(folder_list[0], unique_list(main_list + [_ + ' M3' for _ in value2_list])).rename(columns=lambda x: x.replace("_M3", "_BF"))
df_BF_NB1 = rpt_dask_dataframe(folder_list[1], unique_list(main_list + [_ + 'M3' for _ in value2_list])).rename(columns=lambda x: x.replace('M3', '_BF'))
df_BF_NB2 = rpt_dask_dataframe(folder_list[2], unique_list(main_list + [_ + ' M2' for _ in value2_list])).rename(columns=lambda x: x.replace("M2", '_BF'))
df_BF_NB3 = rpt_dask_dataframe(folder_list[3], unique_list(main_list + [_ + '_M1' for _ in value2_list])).rename(columns=lambda x: x.replace('_M1', '_BF'))
df_AF_EB = rpt_dask_dataframe(folder_list[4], unique_list(main_list + value2_list)).rename(columns=lambda x: x + '_AF' if x in value2_list else x)

df_BF = dd.concat([df_BF_EB, df_BF_NB1, df_BF_NB2, df_BF_NB3], axis=0)

joined_df = dd.merge(df_BF, df_AF_EB, how='outer', on=pk_list, suffixes=['_BF', '_AF'])

# Apply vectorized operations instead of row-wise apply
joined_df['fk1'] = joined_df['fk1_AF'].combine_first(joined_df['fk1_BF'])
joined_df['fk2'] = joined_df['fk2_AF'].combine_first(joined_df['fk2_BF'])
joined_df['Cvalue1_BF'] = joined_df['value1_BF'] + joined_df['value2_BF']
joined_df['Cvalue1_AF'] = joined_df['value1_AF'] + joined_df['value2_AF']
joined_df['value6_DEL'] = joined_df['value6_BF'].where(joined_df['value6_AF'].isna(), 0)
joined_df['value6_ADD'] = joined_df['value6_AF'].where(joined_df['value6_BF'].isna(), 0)
joined_df = joined_df.drop(columns=['fk1_AF', 'fk1_BF', 'fk2_AF', 'fk2_BF', 'value1_AF', 'value1_BF', 'value2_AF', 'value2_BF'])

joined_df = dd.merge(joined_df, reference_df, how='left', on='pk4')

joined_df = joined_df.loc[:, result_col_list]

# Persist joined_df to disk
joined_df = joined_df.persist()

# Convert to delayed and delete joined_df
df_enum = enumerate(joined_df.to_delayed())
del joined_df

for i, chunk in df_enum:
    result = chunk.groupby(summary_columns).sum().compute()    
    result.to_csv(os.path.join(result_folder, f'result_{i}.csv'), index=True)    
    del chunk
    del result


agg_data = dd.read_csv(os.path.join(result_folder, 'result_*.csv'))
agg_data = agg_data.groupby(summary_columns).sum().compute()
agg_data.to_csv(os.path.join(result_folder, 'result.csv'), index=True)

问题根源

  • 外连接数据膨胀:outer join会保留两边所有行,若pk_list存在大量不匹配的键,会产生远超原数据的冗余行,直接导致中间数据暴增。
  • 不必要的持久化:joined_df.persist()将未聚合的全量连接数据写入磁盘,此时数据未经过过滤,体积自然巨大。
  • 低效的分块聚合:逐块compute后导出临时CSV再二次聚合,每个临时文件都包含大量未聚合的原始数据,磁盘IO和存储开销翻倍。
  • 冗余列未及时清理:读取和处理过程中保留了大量最终聚合不需要的列,进一步增大了中间数据体积。

优化方案

1. 严格控制列范围,提前清理冗余数据

只保留连接和聚合必需的列,减少数据携带量:

# 定义BF和AF数据集必需的列
required_bf_cols = pk_list + fk_list + ['value1', 'value2', 'value3', 'value4', 'value5', 'value6', 'value7', 'value8']
required_af_cols = pk_list + fk_list + value2_list

# 读取数据时仅加载必需列
df_BF_EB = rpt_dask_dataframe(folder_list[0], unique_list(required_bf_cols + [_ + ' M3' for _ in value2_list])).rename(columns=lambda x: x.replace("_M3", "_BF"))
df_BF_NB1 = rpt_dask_dataframe(folder_list[1], unique_list(required_bf_cols + [_ + 'M3' for _ in value2_list])).rename(columns=lambda x: x.replace('M3', '_BF'))
df_BF_NB2 = rpt_dask_dataframe(folder_list[2], unique_list(required_bf_cols + [_ + ' M2' for _ in value2_list])).rename(columns=lambda x: x.replace("M2", '_BF'))
df_BF_NB3 = rpt_dask_dataframe(folder_list[3], unique_list(required_bf_cols + [_ + '_M1' for _ in value2_list])).rename(columns=lambda x: x.replace('_M1', '_BF'))
df_AF_EB = rpt_dask_dataframe(folder_list[4], unique_list(required_af_cols)).rename(columns=lambda x: x + '_AF' if x in value2_list else x)

2. 移除不必要的persist,直接用Dask完成全量聚合

跳过逐块导出临时文件的步骤,直接在Dask层面完成聚合,减少磁盘写入:

# 去掉persist和delayed转换逻辑,直接执行聚合
agg_result = joined_df.groupby(summary_columns).sum()
# 直接导出最终结果
agg_result.compute().to_csv(os.path.join(result_folder, 'result.csv'), index=True)

3. 优化外连接逻辑,减少数据膨胀

若业务允许,优先用inner join替代outer join;必须用外连接时,先对两边数据按主键去重:

# 对BF和AF数据按pk_list去重(根据业务选择聚合方式,比如取第一个值)
df_BF = df_BF.groupby(pk_list).first().reset_index()
df_AF_EB = df_AF_EB.groupby(pk_list).first().reset_index()

# 再执行外连接
joined_df = dd.merge(df_BF, df_AF_EB, how='outer', on=pk_list, suffixes=['_BF', '_AF'])

4. 调整Dask临时文件配置

指定大空间临时目录,并开启自动清理:

import dask
# 设置临时文件存储目录(需选择有足够空间的路径)
dask.config.set({'temporary_directory': 'C:/Users/USER/Desktop/dask_temp'})
# 开启临时文件自动清理
dask.config.set({'delayed.cleanup': True})

5. 合理设置分区大小

根据数据量调整分区,避免过多小分区带来的开销:

# 合并后调整分区为1GB/分区(可根据内存情况调整)
df_BF = df_BF.repartition(partition_size='1GB')
joined_df = joined_df.repartition(partition_size='1GB')

内容的提问来源于stack exchange,提问作者기경수

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 01:20:54