如何高效读取并扩充大型CSV数据集?现有方案性能待优化
大型CSV按月拆分的性能优化方案
问题背景
现有13GB CSV数据集,每条样本包含数值及数月起止日期,需转换为单月对应单条样本的格式(如原样本覆盖5-7月,则拆分为3条对应5、6、7月的样本)。当前用Pandas分块+字典存储的实现,处理500k样本耗时约1小时,全量100M样本预估需200小时,硬件配置为8GB RAM、Windows 11 64bit、Python 3.8.8,急需提升性能,尤其是并行处理方案。
现有实现的核心瓶颈
- 逐行循环+低效数据转换:用
iloc遍历每行,再将单行转DataFrame后转字典,完全浪费Pandas向量化操作的性能优势。 - 字典存储冗余:用字典存储中间扩充数据,内存开销大且写入CSV时效率低。
- 单线程运行:未利用CPU多核能力,硬件资源未充分调用。
优化方案
1. Pandas单线程基础优化(先降低单块处理耗时)
通过向量化操作替代循环,大幅提升单块处理效率:
优化点说明
- 读取CSV时直接解析日期列,避免字符串切片提取年月
- 用
date_range生成每个样本覆盖的所有月份,再通过explode拆分多行 - 全程用DataFrame操作,避免字典转换的开销
优化后的代码
import pandas as pd def augment_to_monthly(chunk): # 生成每个样本对应的月份序列(按月份最后一天生成,方便提取年月) chunk['month_dates'] = chunk.apply( lambda row: pd.date_range(start=row['start date'], end=row['end date'], freq='M'), axis=1 ) # 拆分每个月份为单独行 expanded_chunk = chunk.explode('month_dates').reset_index(drop=True) # 提取年月,替换原有的month和year列 expanded_chunk['month'] = expanded_chunk['month_dates'].dt.month expanded_chunk['year'] = expanded_chunk['month_dates'].dt.year # 移除临时列 expanded_chunk = expanded_chunk.drop(columns=['month_dates']) return expanded_chunk # 读取并处理 chunk_count = 0 output_path = 'expanded_data.csv' for chunk in pd.read_csv( 'enc_star_logar_ek.csv', delimiter=';', chunksize=10000, parse_dates=['start date', 'end date'], date_parser=lambda x: pd.to_datetime(x, format='%d/%m/%Y') ): chunk_count += 1 expanded_chunk = augment_to_monthly(chunk) if chunk_count == 1: # 第一次写入表头 expanded_chunk.to_csv(output_path, sep=';', index=False) else: # 追加写入,不写表头 expanded_chunk.to_csv(output_path, sep=';', mode='a', header=False, index=False) print(f"处理完成第 {chunk_count} 块")
2. Dask并行处理方案(利用多核加速)
Dask可自动将任务拆分到多个CPU核心并行执行,适合超大数据量处理,同时自动控制内存使用:
步骤说明
- 用Dask读取CSV,自动分块
- 定义向量化的处理函数,和Pandas优化版逻辑一致
- 并行计算后导出结果
Dask实现代码
import dask.dataframe as dd import pandas as pd def augment_to_monthly_dask(df): # 生成每个样本的月份序列 df['month_dates'] = df.apply( lambda row: pd.date_range(start=row['start date'], end=row['end date'], freq='M'), meta=('month_dates', 'object'), axis=1 ) # 拆分多行 df = df.explode('month_dates') # 提取年月 df['month'] = df['month_dates'].dt.month df['year'] = df['month_dates'].dt.year # 移除临时列 df = df.drop(columns=['month_dates']) return df # 读取CSV,指定块大小(8GB内存建议设为64MB左右) dask_df = dd.read_csv( 'enc_star_logar_ek.csv', delimiter=';', blocksize='64MB', parse_dates=['start date', 'end date'], date_parser=lambda x: pd.to_datetime(x, format='%d/%m/%Y') ) # 并行处理 expanded_dask_df = dask_df.map_partitions( augment_to_monthly_dask, meta=dask_df.dtypes.append(pd.Series([int, int], index=['month', 'year'])) ) # 导出结果(自动并行写入单个CSV) expanded_dask_df.to_csv('dask_expanded_data.csv', sep=';', index=False, single_file=True)
注意事项
blocksize根据内存调整:8GB内存建议设为64-128MB,避免内存溢出meta参数必须指定,Dask需要知道处理后的数据结构- 若不需要单个大文件,可去掉
single_file=True,生成多个分块文件
额外性能提升建议
- 数据类型优化:读取时指定列的 dtype,比如数值列用
int32/float32(精度允许的情况下),减少内存占用 - 过滤无效数据:提前过滤跨多年的样本,减少后续处理量
- 减少IO操作:尽量减少文件打开/关闭次数,Dask批量写入效率优于Pandas逐块追加
内容的提问来源于stack exchange,提问作者AKΛ
相关产品推荐
相关产品推荐

