将单元格含一维数组的两个DataFrame合并为含二维数组的DataFrame的最优方案
高效合并两个大DataFrame的单元格数组
需求说明
需要创建一个每个单元格包含二维数组的DataFrame,其中每个二维数组由DataFrame A和DataFrame B对应位置单元格的一维数组组合而成。由于要处理约200GB的两个DataFrame,急需高效、简洁的实现方案,当前已尝试结合numpy的apply方法,但不确定是否为最优解。
数据示例
DataFrame A
col1 col2 0 [1, 2, 3, 4] [9, 10, 11, 12] 1 [5, 6, 7, 8] [13, 14, 15, 16]
DataFrame B
col1 col2 0 [17, 18, 19, 20] [25, 26, 27, 28] 1 [21, 22, 23, 24] [29, 30, 31, 32]
期望结果
col1 col2 0 [[1, 2, 3, 4], [17, 18, 19, 20]] [[9, 10, 11, 12], [25, 26, 27, 28]] 1 [[5, 6, 7, 8], [21, 22, 23, 24]] [[13, 14, 15, 16], [29, 30, 31, 32]]
当前实现代码
import numpy as np import pandas as pd def function(x): a[x.name] = pd.Series(np.stack([x, b[x.name]], axis=1).tolist(), name=x.name) data_a = {'col1': [[1, 2, 3, 4], [5, 6, 7, 8]], 'col2': [[9, 10, 11, 12], [13, 14, 15, 16]]} a = pd.DataFrame(data=data_a) data_b = {'col1': [[17, 18, 19, 20], [21, 22, 23, 24]], 'col2': [[25, 26, 27, 28], [29, 30, 31, 32]]} b = pd.DataFrame(data=data_b) print(a) print() print(b) a.apply(function) print(a)
优化建议
针对200GB级别的大文件,直接全量加载到内存不现实,以下是几个高效的优化方向:
1. 抛弃apply,改用向量化操作
apply本质是逐元素遍历,效率极低。可以直接对整个DataFrame做numpy批量堆叠操作,再转换为列表:
# 假设a和b已按分块加载完成 result = pd.DataFrame() for col in a.columns: # 把两列数组堆叠成二维数组后转列表 result[col] = np.stack([a[col].to_numpy(), b[col].to_numpy()], axis=1).tolist()
这种方式利用numpy的批量处理能力,比apply效率提升数倍。
2. 分块加载处理
200GB数据无法一次性入内存,必须用chunksize参数分块读取,逐块处理后合并结果:
chunk_size = 10_000 # 根据可用内存调整块大小 result_chunks = [] # 分块读取A和B for chunk_a, chunk_b in zip(pd.read_csv('data_a.csv', chunksize=chunk_size), pd.read_csv('data_b.csv', chunksize=chunk_size)): # 若原始文件存储的是字符串形式数组,先转为列表 chunk_a = chunk_a.applymap(eval) chunk_b = chunk_b.applymap(eval) # 处理当前块 processed_chunk = pd.DataFrame() for col in chunk_a.columns: processed_chunk[col] = np.stack([chunk_a[col].to_numpy(), chunk_b[col].to_numpy()], axis=1).tolist() result_chunks.append(processed_chunk) # 合并所有块并保存 final_result = pd.concat(result_chunks, ignore_index=True) final_result.to_csv('merged_result.csv', index=False)
注:若使用Parquet等二进制格式存储数据,可直接读取数组类型,无需eval转换。
3. 用Dask做并行处理
Dask可自动处理超内存数据集的分块与并行计算,适合大规模数据场景:
import dask.dataframe as dd # 读取Dask DataFrame dd_a = dd.read_csv('data_a.csv') dd_b = dd.read_csv('data_b.csv') # 转换字符串数组为列表(若需要) dd_a = dd_a.applymap(eval, meta=dd_a.dtypes) dd_b = dd_b.applymap(eval, meta=dd_b.dtypes) # 定义合并函数 def merge_arrays(x, y): return [x, y] # 逐列合并数组 result_dd = dd.concat([dd_a[col].map(merge_arrays, dd_b[col]) for col in dd_a.columns], axis=1) # 计算并保存结果 result_dd.compute().to_csv('dask_merged_result.csv', index=False)
4. 预处理为numpy三维数组(内存允许时)
若所有一维数组长度固定,可将整个DataFrame转为numpy三维数组,直接操作后转回DataFrame,效率最高:
# 假设每个一维数组长度为4 arr_a = np.array(a.to_numpy().tolist()) arr_b = np.array(b.to_numpy().tolist()) # 堆叠得到三维数组后转列表 merged_arr = np.stack([arr_a, arr_b], axis=2).tolist() # 转回DataFrame result = pd.DataFrame(merged_arr, columns=a.columns)
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

