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

如何分割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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 20:20:35