如何分割Python DataFrame为k份,迭代/并行处理后合并?
解决方案
一、迭代处理所有子集(替代硬编码)
你已经用np.array_split得到了存储子集的列表dfs,直接遍历这个列表即可完成自动化处理,完全不需要globals()或exec()这类不规范的写法。以下是完整的迭代处理代码:
import pandas as pd import numpy as np from tqdm import tqdm # 假设你的空间函数已定义:foo.bar()、baz() dfs = np.array_split(dftest, 6) processed_dfs = [] for idx, df_subset in enumerate(tqdm(dfs)): # 匹配对应参数:param0~param5 current_param = f"param{idx}" # 复制子集避免修改原数据 df_processed = df_subset.copy() # 先初始化新增列,提升性能 df_processed['anotherColumn'] = np.nan for row in df_processed.itertuples(): x = row.type y = foo.bar(x, param=current_param) # 用行索引赋值 df_processed.loc[row.Index, 'anotherColumn'] = baz(y) # 其他新增列的处理逻辑同理 # 导出当前子集结果 csv_path = f"/projectPath/dfs{idx}.csv" df_processed.to_csv(csv_path, index=False) # 存入列表用于后续合并 processed_dfs.append(df_processed) # 合并所有处理后的子集 dfresult = pd.concat(processed_dfs, ignore_index=True)
关键说明:
- 用
enumerate同时获取子集索引和子集本身,自动对应param0到param5 - 复制子集后再处理,避免意外修改原列表中的数据
- 不管分割成6份还是100份,循环逻辑完全通用,无需手动编写所有子集索引
二、并行处理(加速耗时任务)
针对单份处理耗时久的场景,可通过并行计算利用多核CPU提升效率,推荐用joblib实现,代码简洁易维护:
from joblib import Parallel, delayed # 定义单个子集的处理函数 def process_subset(df_subset, param, save_path): df_processed = df_subset.copy() df_processed['anotherColumn'] = np.nan for row in df_processed.itertuples(): x = row.type y = foo.bar(x, param=param) df_processed.loc[row.Index, 'anotherColumn'] = baz(y) # 其他处理步骤 df_processed.to_csv(save_path, index=False) return df_processed # 准备参数和存储路径列表 params = [f"param{i}" for i in range(6)] save_paths = [f"/projectPath/dfs{i}.csv" for i in range(6)] # 并行处理,n_jobs=-1表示使用所有CPU核心 processed_dfs = Parallel(n_jobs=-1, verbose=10)( delayed(process_subset)(dfs[i], params[i], save_paths[i]) for i in range(len(dfs)) ) # 合并结果 dfresult = pd.concat(processed_dfs, ignore_index=True)
注意事项:
- 确保
foo.bar()和baz()是可序列化的(picklable),如果是在Jupyter Notebook中,建议将这些函数放在单独的.py文件中导入 verbose=10会显示处理进度,方便跟踪任务状态
三、优化建议
如果你的空间函数支持批量处理,尽量用矢量化操作替代逐行循环,能大幅提升性能:
def batch_process_type(series_type, param): # 批量处理type列,生成结果 y_series = series_type.apply(lambda x: foo.bar(x, param=param)) return y_series.apply(baz) # 在迭代处理中替换逐行循环: df_processed['anotherColumn'] = batch_process_type(df_processed['type'], current_param)
内容的提问来源于stack exchange,提问作者Lizardie
相关产品推荐
相关产品推荐

