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

针对超大型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:

customeroriginallist1list2list3
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 09:30:36