针对超大型Pandas DataFrame的多线程/多进程实现方案咨询
百万列DataFrame批量替换的并行化实现
问题背景
假设存在一个超100万列的巨型Pandas DataFrame,现有串行逻辑已实现按指定字典批量生成替换后的新列(将指定元素替换为ham),但处理速度无法满足需求,需要通过多线程/多进程加速。
示例基础代码与预期结果:
import pandas as pd import numpy as np from concurrent.futures import * import multiprocessing num_processes = multiprocessing.cpu_count() print(f'num_precesses: {num_processes}') def ReplaceItem(old_item, input_list, new_item = 'ham'): output = [] for item in input_list: if item == old_item: output.append(new_item) else: output.append(item) return output # 示例DataFrame food = pd.DataFrame({'customer':range(1,4), 'original':[['egg', 'berry', 'pork', 'tea'], ['chicken', 'beef', 'water', 'soda'], ['fish', 'chicken', 'pork', 'coffee']]}) replace_items = {'list1':'egg', 'list2':'chicken', 'list3':'fish'}
预期生成包含替换后新列的DataFrame:
| customer | original | list1 | list2 | list3 |
|---|---|---|---|---|
| 1 | ['egg', 'berry', 'pork', 'tea'] | ['ham', 'berry', 'pork', 'tea'] | ['egg', 'berry', 'pork', 'tea'] | ['egg', 'berry', 'pork', 'tea'] |
| 2 | ['chicken', 'beef', 'water', 'soda'] | ['chicken', 'beef', 'water', 'soda'] | ['ham', 'beef', 'water', 'soda'] | ['chicken', 'beef', 'water', 'soda'] |
| 3 | ['fish', 'chicken', 'pork', 'coffee'] | ['fish', 'chicken', 'pork', 'coffee'] | ['fish', 'ham', 'pork', 'coffee'] | ['ham', 'chicken', 'pork', 'coffee'] |
并行化实现方案
一、多线程实现(ThreadPoolExecutor)
适合轻计算/IO密集型场景,Pandas处理Python对象时GIL释放充分,多线程可利用CPU空闲时间。
import pandas as pd import numpy as np from concurrent.futures import ThreadPoolExecutor import multiprocessing num_workers = multiprocessing.cpu_count() # 优化ReplaceItem:用列表推导式替代逐行append,提升单任务效率 def ReplaceItem(old_item, input_list, new_item='ham'): return [new_item if item == old_item else item for item in input_list] # 定义单个替换任务 def process_replace_task(key, old_val): return (key, food['original'].apply(lambda x: ReplaceItem(old_val, x))) # 多线程执行 with ThreadPoolExecutor(max_workers=num_workers) as executor: # 提交所有替换任务 futures = [executor.submit(process_replace_task, k, v) for k, v in replace_items.items()] # 收集结果并转为字典 result_dict = {fut.result()[0]: fut.result()[1] for fut in futures} # 合并原DataFrame与结果 new_df = pd.DataFrame(result_dict) update_df = pd.concat([food, new_df], axis=1) print(update_df)
二、多进程实现(ProcessPoolExecutor)
适合CPU密集型场景,绕过GIL限制,充分利用多核CPU。注意:多进程需传递可序列化数据,避免直接传递整个大DataFrame。
import pandas as pd import numpy as np from concurrent.futures import ProcessPoolExecutor import multiprocessing num_workers = multiprocessing.cpu_count() def ReplaceItem(old_item, input_list, new_item='ham'): return [new_item if item == old_item else item for item in input_list] # 多进程批量处理任务:接收旧值与所有原始列表,返回批量处理结果 def process_replace_batch(old_val, original_list): return [ReplaceItem(old_val, lst) for lst in original_list] # 准备任务参数:提取原始列数据,避免在进程间传递整个DataFrame original_data = food['original'].tolist() tasks = [(val, original_data) for val in replace_items.values()] with ProcessPoolExecutor(max_workers=num_workers) as executor: # 执行所有任务并收集结果 results = list(executor.map(lambda args: process_replace_batch(*args), tasks)) # 将结果映射为对应列名,生成新DataFrame result_dict = {key: results[i] for i, key in enumerate(replace_items.keys())} new_df = pd.DataFrame(result_dict) update_df = pd.concat([food, new_df], axis=1) print(update_df)
额外优化建议
- 函数预优化:用列表推导式替代
ReplaceItem中的逐行append,单任务效率提升明显 - 减少内存开销:串行逻辑中每次循环都执行
concat,并行方案仅在最后合并一次,避免重复内存分配 - 内存管理:若DataFrame真达百万列规模,建议分批次处理并写入磁盘(如
to_csv分块写入),避免内存溢出 - 并行方式选择:轻计算/IO密集用多线程,CPU密集用多进程
内容的提问来源于stack exchange,提问作者DaCard
相关产品推荐
相关产品推荐

