使用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,提问作者기경수
相关产品推荐
相关产品推荐

