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

如何通过多线程/进程拆分大型pandas DataFrame并行处理?

问题描述

我有一个约30万行的pandas DataFrame,需要读取并逐行执行以下操作:

for index, row in df.iterrows():
   # 对row['col A']和row['col B']执行处理操作
   # 根据处理结果为初始为空的row['col C']赋值

目前已通过以下方式将DataFrame拆分为多个子DataFrame:

df = pd.read_csv('file.csv')
df_split = np.array_split(df, 10)

def one_split(sub_df): # sub_df为每个线程/进程处理的独立子DataFrame
   for index, row in sub_df.iterrows():
      # 步骤1:处理row['col A']和row['col B']
      # 步骤2:根据步骤1的结果为row['col C']赋值
      # 步骤3(待实现):将sub_df中'col C'的值复制回原DataFrame df的对应列

one_split(df_split[0])

目标:

  • 将任务拆分到10个线程/进程中并行执行,提升运行效率;
  • 每个线程/进程并行执行one_split(df_split[i]);
  • 实现one_split()中的步骤3,由于各线程/进程处理原DataFrame的不同部分,不存在race condition(竞争条件)。

补充说明:预期行为为拆分后的子DataFrame各自处理col C,最终将结果合并回原DataFrame对应位置。


解决方案

方案一:多进程(适合CPU密集型处理)

如果你的处理逻辑以CPU计算为主,多进程可以规避Python GIL限制,最大化利用多核CPU。核心思路是让子进程返回处理后的子DataFrame,最后合并回原结构。

import pandas as pd
import numpy as np
from multiprocessing import Pool

def process_sub_df(sub_df):
    # 替换为你的实际处理逻辑
    for idx, row in sub_df.iterrows():
        processed_value = row['col A'] + row['col B']  # 示例计算
        sub_df.loc[idx, 'col C'] = processed_value
    return sub_df

if __name__ == '__main__':
    df = pd.read_csv('file.csv')
    df['col C'] = np.nan  # 初始化空列
    df_split = np.array_split(df, 10)
    
    # 启动10个进程并行处理
    with Pool(10) as pool:
        processed_sub_dfs = pool.map(process_sub_df, df_split)
    
    # 合并处理后的子DataFrame,恢复原索引顺序
    df = pd.concat(processed_sub_dfs).sort_index()

方案二:多线程(适合IO密集型处理)

如果处理逻辑涉及大量IO操作(如调用外部API、读写文件),线程池的开销更低,适合这类场景。

import pandas as pd
import numpy as np
from concurrent.futures import ThreadPoolExecutor

def process_sub_df(sub_df):
    # 替换为你的实际处理逻辑
    for idx, row in sub_df.iterrows():
        processed_value = row['col A'] + row['col B']  # 示例计算
        sub_df.loc[idx, 'col C'] = processed_value
    return sub_df

df = pd.read_csv('file.csv')
df['col C'] = np.nan
df_split = np.array_split(df, 10)

# 启动10个线程并行处理
with ThreadPoolExecutor(max_workers=10) as executor:
    processed_sub_dfs = list(executor.map(process_sub_df, df_split))

# 合并结果
df = pd.concat(processed_sub_dfs).sort_index()

关键注意点

  1. 不要直接在子进程/线程中修改原DataFrame:多进程会复制数据,直接修改无法同步到主进程;线程虽共享内存,但逐行修改仍易出问题,返回子DataFrame再合并是更可靠的方式。
  2. 无竞争条件:每个任务处理独立的行区间,合并时按索引排序,不会出现数据冲突。
  3. 进阶优化:如果逐行处理效率仍低,建议将处理逻辑向量化(如用df.apply、np.vectorize),比iterrows快数倍,再结合并行效果更佳。

内容的提问来源于stack exchange,提问作者abs8090

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 10:10:34