如何高效合并列不同的数千个Pandas DataFrame并避免内存错误
合并数千个列不固定的Pandas DataFrame(避免内存错误)
问题背景
需要合并数千个列不固定、列信息预先未知的Pandas DataFrame,直接调用pd.concat()会触发内存错误,序列化方案因列不固定无法适用。
示例代码
import pandas as pd import numpy as np df1 = pd.DataFrame(np.array([[1, 2, 3], [4, 5, 6], [7, 8, 9]]), columns=['a', 'b', 'c']) df2 = pd.DataFrame(np.array([[1, 2, 3], [4, 5, 6], [7, 8, 9]]), columns=['a', 'k', 'c'])
期望合并结果
a b c k 0 1 2.0 3 NaN 1 4 5.0 6 NaN 2 7 8.0 9 NaN 0 1 NaN 3 2.0 1 4 NaN 6 5.0 2 7 NaN 9 8.0
当前实现方案
import os import pandas as pd import numpy as np def bycsv(dfs): md, hd = 'w', True for df in dfs: df.to_csv('df_all.csv', mode=md, header=hd, index=None) md, hd = 'a', False df_all = pd.read_csv('df_all.csv', index_col=None) os.remove('df_all.csv') return df_all # 收集所有列并去重排序 cols = [] for df in [df1, df2]: cols += list(df) cols = sorted(list(set(cols))) # 统一所有DataFrame的列 dfsok = [df.reindex(columns=cols) for df in [df1, df2]] # 执行合并 result = bycsv(dfsok)
更高效的优化方案
方案1:优化CSV写入逻辑(减少IO开销)
当前方案每次调用to_csv()都会重新打开文件,改用pandas.io.parsers.TextFileWriter复用写入对象,降低IO开销;同时用集合收集列更高效:
import pandas as pd import numpy as np import os def efficient_csv_concat(dfs): # 收集所有唯一列(用集合避免重复) all_columns = set() for df in dfs: all_columns.update(df.columns) all_columns = sorted(all_columns) # 初始化CSV写入器,一次性写入表头 with open('temp_merge.csv', 'w', newline='', encoding='utf-8') as f: writer = pd.io.parsers.TextFileWriter(f, delimiter=',', quoting=0) writer.writeheader(all_columns) # 逐个处理DataFrame,对齐列后写入 for df in dfs: aligned_df = df.reindex(columns=all_columns) writer.write_rows(aligned_df.itertuples(index=False, name=None)) # 读取合并结果并清理临时文件 merged_df = pd.read_csv('temp_merge.csv') os.remove('temp_merge.csv') return merged_df # 测试调用 result = efficient_csv_concat([df1, df2]) print(result)
优势:避免多次打开/关闭文件,减少IO开销;处理一个DF就释放对应内存,内存占用更低。
方案2:使用Dask DataFrame(超大数据量首选)
如果数据量极大,Dask可分布式处理数据,无需将所有DataFrame加载到内存:
import dask.dataframe as dd import pandas as pd import numpy as np # 将每个Pandas DataFrame转为Dask分区(若DF来自文件,直接用dd.read_csv读取更高效) dask_dfs = [dd.from_pandas(df, npartitions=1) for df in [df1, df2]] # 自动对齐列并合并 merged_dask = dd.concat(dask_dfs, axis=0) # 转为Pandas DataFrame(按需选择,也可直接用Dask做后续处理) merged_df = merged_dask.compute() print(merged_df)
优势:支持并行计算,内存占用可控;若原始数据来自文件,直接用Dask读取可跳过Pandas加载步骤,效率更高。
方案3:使用Parquet格式替代CSV(更快的IO)
Parquet是列式存储格式,读写速度远快于CSV,且支持复杂数据类型,内存占用更低:
import pandas as pd import numpy as np import os def parquet_concat(dfs): # 收集所有唯一列 all_columns = set() for df in dfs: all_columns.update(df.columns) all_columns = sorted(all_columns) # 逐个写入Parquet文件(append模式) first_write = True for df in dfs: aligned_df = df.reindex(columns=all_columns) aligned_df.to_parquet( 'temp_merge.parquet', append=not first_write, engine='pyarrow', index=False ) first_write = False # 读取合并结果并清理临时文件 merged_df = pd.read_parquet('temp_merge.parquet', engine='pyarrow') os.remove('temp_merge.parquet') return merged_df # 测试调用 result = parquet_concat([df1, df2]) print(result)
优势:读写速度比CSV提升数倍,适合列数多、数据量大的场景;文件体积更小。
内容的提问来源于stack exchange,提问作者Forinstance
相关产品推荐
相关产品推荐

