Pandas DataFrame转换性能优化:替换For Loop处理亿级多文件数据集
问题背景
我日常使用Python开展研究工作,首次处理规模超亿行的多文件数据集,所用设备为配置Xeon E5-2637 v4 CPU、Quadro K420 GPU的老旧但性能尚可的工作站。目前通过嵌套For Loop实现Pandas DataFrame的数据转换,但处理速度无法满足需求,已查阅性能优化文档及尝试使用groupby替代循环但未取得效果,恳请提供加速该算法的方案。
数据格式(所有文件一致)
C:/../data1.csv -- col1 col2 col3 parent abcde NaN child d3d a1a child s2s f4f parent fghij NaN child g5g h6h child j7j k8k
原实现代码
#list of file locations filelist = {'files': ['C:/../data1.csv', 'C:/../data2.csv', 'C:/../data3.csv']} filelist_df = pd.DataFrame(data=filelist) filelist_df = filelist_df["files"].str.strip("[]") #data transformation column_names=['1', '2', '3', '4'] temp_parent=[] for i in range(3): new_df=pd.DataFrame(columns=column_names) data_df=pd.read_csv(filelist_df[i], skiprows=1, names=column_names) for j in range(len(data_df)): if data_df['1'][j]=='parent': temp_parent=data_df['2'][j] else: data_df['4'][j]=temp_parent temp_row=data_df.loc[j,:] new_df = new_df.append(temp_row, ignore_index=True) new_df.to_csv('C:/../new%d' % i + '.csv', index=False, header=False) del new_df, data_df, temp_parent, temp_row
输出示例(仅data1.csv的输出)
C:/../new0.csv -- child d3d a1a abcde child s2s f4f abcde child g5g h6h fghij child j7j k8k fghij
加速方案
1. 替换逐行循环与append操作(核心优化)
原代码的嵌套循环和DataFrame.append是性能瓶颈——append每次都会生成新的DataFrame对象,时间复杂度极高。可以利用Pandas的矢量化操作和向前填充(ffill)实现需求:
import pandas as pd filelist = ['C:/../data1.csv', 'C:/../data2.csv', 'C:/../data3.csv'] for idx, filepath in enumerate(filelist): # 读取数据,跳过表头并指定列名 df = pd.read_csv(filepath, skiprows=1, names=['1', '2', '3', '4']) # 提取parent行的col2值到col4 df['4'] = df.loc[df['1'] == 'parent', '2'] # 向前填充col4,让所有child行继承最近的parent值 df['4'] = df['4'].ffill() # 过滤掉parent行,只保留child数据 new_df = df[df['1'] != 'parent'].reset_index(drop=True) # 保存结果 new_df.to_csv(f'C:/../new{idx}.csv', index=False, header=False)
2. 分块读取处理超大文件
如果单文件体积超过内存容量,使用chunksize分块读取处理,避免内存溢出:
import pandas as pd filelist = ['C:/../data1.csv', 'C:/../data2.csv', 'C:/../data3.csv'] chunk_size = 10**6 # 根据内存情况调整,比如100万行/块 for idx, filepath in enumerate(filelist): output_path = f'C:/../new{idx}.csv' first_write = True # 逐块读取并处理 for chunk in pd.read_csv(filepath, skiprows=1, names=['1', '2', '3', '4'], chunksize=chunk_size): chunk['4'] = chunk.loc[chunk['1'] == 'parent', '2'] chunk['4'] = chunk['4'].ffill() # 过滤parent行 chunk = chunk[chunk['1'] != 'parent'] # 写入文件,第一块不写表头,后续追加 chunk.to_csv(output_path, index=False, header=first_write, mode='a') first_write = False
3. 用Dask处理超大规模数据集
如果Pandas单进程处理仍不够快,Dask可以实现并行计算,支持超出内存的数据集处理:
import dask.dataframe as dd filelist = ['C:/../data1.csv', 'C:/../data2.csv', 'C:/../data3.csv'] for idx, filepath in enumerate(filelist): # 用Dask读取数据 ddf = dd.read_csv(filepath, skiprows=1, names=['1', '2', '3', '4']) # 同样用ffill填充parent值 ddf['4'] = ddf.loc[ddf['1'] == 'parent', '2'].ffill() # 过滤parent行 ddf = ddf[ddf['1'] != 'parent'] # 保存为单文件 ddf.to_csv(f'C:/../new{idx}.csv', index=False, header=False, single_file=True)
4. 额外小优化
- 读取时指定
usecols只加载需要的列,减少内存占用:pd.read_csv(..., usecols=['1','2','3'])(后续再添加列4) - 将字符串列转为
category类型(如果重复值较多):df['1'] = df['1'].astype('category'),降低内存消耗 - 关闭Pandas警告:
pd.set_option('mode.chained_assignment', None),减少不必要的开销
内容的提问来源于stack exchange,提问作者byeme
相关产品推荐
相关产品推荐

