如何通过多线程/进程拆分大型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()
关键注意点
- 不要直接在子进程/线程中修改原DataFrame:多进程会复制数据,直接修改无法同步到主进程;线程虽共享内存,但逐行修改仍易出问题,返回子DataFrame再合并是更可靠的方式。
- 无竞争条件:每个任务处理独立的行区间,合并时按索引排序,不会出现数据冲突。
- 进阶优化:如果逐行处理效率仍低,建议将处理逻辑向量化(如用
df.apply、np.vectorize),比iterrows快数倍,再结合并行效果更佳。
内容的提问来源于stack exchange,提问作者abs8090
相关产品推荐
相关产品推荐

