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

如何高效读取并扩充大型CSV数据集?现有方案性能待优化

大型CSV按月拆分的性能优化方案

问题背景

现有13GB CSV数据集,每条样本包含数值及数月起止日期,需转换为单月对应单条样本的格式(如原样本覆盖5-7月,则拆分为3条对应5、6、7月的样本)。当前用Pandas分块+字典存储的实现,处理500k样本耗时约1小时,全量100M样本预估需200小时,硬件配置为8GB RAM、Windows 11 64bit、Python 3.8.8,急需提升性能,尤其是并行处理方案。

现有实现的核心瓶颈

  1. 逐行循环+低效数据转换:用iloc遍历每行,再将单行转DataFrame后转字典,完全浪费Pandas向量化操作的性能优势。
  2. 字典存储冗余:用字典存储中间扩充数据,内存开销大且写入CSV时效率低。
  3. 单线程运行:未利用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Λ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 20:31:35