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

寻求Pandas DataFrame行级处理的简易高效并行替代方案(替代apply方法)

简单实现Pandas Apply并行化的方案

针对你遇到的apply单线程处理慢的问题,有两个非常易用的第三方库可以快速实现并行化,几乎不需要修改原有代码,完美匹配你的需求:


方案一:使用pandarallel(推荐,直接替换apply)

这个库专门为Pandas的apply方法做了并行封装,配置简单,开箱即用。

步骤:

  1. 先安装库:
pip install pandarallel
  1. 修改你的代码,只需要初始化并行环境,然后把apply换成parallel_apply:
from pandarallel import pandarallel
import pandas as pd
import multiprocessing as mp

# 初始化并行,默认会使用所有CPU核心,也可以手动指定进程数(比如nb_workers=4)
# 加上progress_bar=True可以看到处理进度,很实用
pandarallel.initialize(nb_workers=mp.cpu_count(), progress_bar=True)

def _format(data: pd.DataFrame, context: pd.DataFrame):
    # 直接替换apply为parallel_apply即可,原有lambda和get_context_value完全不用改
    data['context'] = data.parallel_apply(lambda row: get_context_value(context, row), axis=1)

优势:

  • 零侵入:不需要修改你的get_context_value函数逻辑,完全兼容原有代码
  • 自动拆分数据:底层会把DataFrame分成多个chunk分配到不同进程,无需手动处理数据分片
  • 进度可视化:开启进度条后能实时看到处理进度,避免等待焦虑

方案二:使用swifter(智能自动选择最优方式)

这个库会自动判断你的函数和数据规模,选择最快的执行方式——如果适合并行就用多进程/Dask加速,否则用普通apply,非常省心。

步骤:

  1. 安装库:
pip install swifter
  1. 修改代码:
import swifter
import pandas as pd

def _format(data: pd.DataFrame, context: pd.DataFrame):
    # 把apply换成swifter.apply,其他完全不变
    data['context'] = data.swifter.apply(lambda row: get_context_value(context, row), axis=1)

优势:

  • 智能判断:不需要手动配置并行参数,库会自动选择最优方案
  • 兼容更多场景:如果后续数据量变小或者函数逻辑变化,它会自动切换回单线程,避免不必要的并行开销

注意事项

  1. 函数可序列化:因为多进程依赖pickle序列化,确保你的get_context_value函数和context DataFrame是可以被pickle的——如果函数里有一些无法序列化的对象(比如自定义的非pickle类、网络连接等),需要调整逻辑,比如把context的必要数据提前提取为普通字典/数组,或者在进程内部重新初始化这些对象。
  2. 进程数不要超配:一般设置为CPU核心数即可(比如mp.cpu_count()),过多的进程会导致CPU上下文切换开销剧增,反而变慢。
  3. 数据大小考量:如果context DataFrame非常大,建议把它做成全局变量(在进程初始化时加载),避免每个进程都复制一份,减少内存开销。

手动实现多进程(备选,不推荐除非不想用第三方库)

如果不想依赖第三方库,也可以用Python自带的multiprocessing手动实现,但代码会繁琐一些:

import multiprocessing as mp
import pandas as pd

# 注意:如果context很大,建议把它设为全局变量,或者用初始化函数传入进程
def process_row(row):
    return get_context_value(context, row)

def _format(data: pd.DataFrame, context_df: pd.DataFrame):
    global context
    context = context_df  # 把context设为全局变量,让子进程能访问
    rows = [row for _, row in data.iterrows()]
    
    with mp.Pool(mp.cpu_count()) as pool:
        results = pool.map(process_row, rows)
    
    data['context'] = results

不过这个方法需要处理全局变量或者进程初始化,不如前面两个第三方库省心,所以优先推荐前两个方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 10:37:41