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

如何在DataFrame上使用多进程?我的实现代码报错求助

搞定多进程处理DataFrame的坑

嘿,首先得说Stack Overflow确实帮了咱们程序员不少大忙,能理解你想加速百万行DataFrame处理的需求!不过你的代码里踩了几个多进程和pandas配合的典型坑,导致了JSONDecodeError,而且就算没报错,也根本不会修改原DataFrame,咱们一步步来解决:

为什么原代码会报错?

1. 子进程拿不到父进程的原DataFrame

当你用multiprocessing.Pool时,每个子进程会复制父进程的内存空间,也就是说你在jobs函数里操作的df是子进程自己的拷贝,不是原DataFrame。而且pandas的DataFrame作为全局变量跨进程传递时,会触发序列化操作,这就是你看到JSONDecodeError的根源——序列化过程中出了问题。

2. 直接修改全局变量的方式行不通

多进程之间默认是内存隔离的,子进程里对df的修改不会同步到父进程,就算代码不报错,最后你的原DataFrame也不会有任何变化。

正确的多进程实现方式

正确的思路是:把DataFrame拆成分片,传递给子进程处理,子进程返回处理后的分片,最后合并回完整的DataFrame。这里给你修正后的代码:

import pandas as pd
from multiprocessing import Pool

# 定义处理单个分片的函数,接收分片数据,返回处理后的结果
def process_chunk(chunk):
    # 这里替换成你的实际复杂逻辑,示例是给第二列乘以40
    chunk.iloc[:, 1] = chunk.iloc[:, 1] * 40
    return chunk

if __name__ == '__main__':
    # 模拟你的百万行DataFrame
    df = pd.DataFrame({
        'col1': range(1, 1000001),
        'col2': range(1, 1000001)
    })
    
    # 将DataFrame拆分为两个分片
    mid_point = len(df) // 2
    chunk1 = df.iloc[:mid_point, :]
    chunk2 = df.iloc[mid_point:, :]
    
    # 使用进程池处理分片
    with Pool(processes=2) as p:
        processed_chunks = p.map(process_chunk, [chunk1, chunk2])
    
    # 合并处理后的分片得到最终结果
    processed_df = pd.concat(processed_chunks)
    print(processed_df.head())

更高效的替代方案:优先用矢量化操作

其实你的示例逻辑完全不需要循环和多进程!pandas的矢量化操作是底层用C实现的,速度比任何Python循环(包括多进程循环)快得多:

# 一行搞定,速度秒杀循环
df.iloc[:, 1] = df.iloc[:, 1] * 40

如果你的实际逻辑能转化为矢量化操作(比如用pandas的内置函数、布尔索引等),一定要优先用这种方式,这才是pandas的正确打开方式。

如果必须用复杂循环处理

如果你的实际逻辑真的无法矢量化(比如涉及复杂的条件判断或外部调用),除了上面的分片处理,还可以考虑:

  • 用multiprocessing.Manager创建共享数据结构,但这种方式开销较大,效率不一定高;
  • 用专门的大数据并行库,比如Dask,它可以自动帮你处理DataFrame的并行计算,语法和pandas几乎一致。

总结一下:你的原代码问题在于跨进程操作全局DataFrame导致序列化错误和无效修改,正确的方式是分片传递、处理后合并;优先选择矢量化操作,这能让你的代码既简洁又高效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:48:43