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

使用Python Pool多进程处理含pd.pickle操作的函数时出现EOFError的原因咨询

问题分析与解决思路

我之前在多进程处理Pandas数据的时候也碰到过一模一样的EOFError,尤其是任务量上去之后更容易触发。咱们来梳理一下可能的原因和对应的解决办法:

核心原因:多进程下的文件操作冲突

当你用multiprocessing.Pool启动多个进程时,如果这些进程同时对同一个pickle文件进行读/写操作,就会出现资源竞争问题:

  • 多个进程同时写入同一个文件:会导致文件内容被截断、覆盖或者写入不完整,后续读取的进程拿到的是损坏的文件,自然就抛出EOFError(因为pickle数据没写完就被读了)。
  • 读和写同时发生:比如进程A正在写入文件,进程B同时去读,此时文件处于半写入状态,B读取到的内容不完整,也会触发EOF错误。

为什么超过10个元素才出现?因为默认情况下Pool的进程数等于CPU核心数,当任务数少于进程数时,进程是串行处理或者并发度低,冲突概率小;当任务数超过进程数,任务开始排队,进程反复复用,文件操作的冲突频率就会急剧上升,错误就显现出来了。

解决办法

1. 给每个进程分配独立的文件路径

这是最稳妥的方案,彻底避免文件竞争。让每个任务对应一个唯一的输入/输出文件,不要让多个进程共享同一个文件路径:

import pandas as pd
from multiprocessing import Pool

def my_function(item):
    # 每个item对应唯一的文件名,避免冲突
    input_file = f"input_data_{item}.pkl"
    output_file = f"processed_data_{item}.pkl"
    
    # 读取数据
    df = pd.read_pickle(input_file)
    # 你的数据处理逻辑
    df['new_col'] = item * 2
    # 写入处理后的数据
    pd.to_pickle(df, output_file)

if __name__ == "__main__":
    items = list(range(15))  # 超过10个元素的任务列表
    with Pool() as pool:
        pool.map(my_function, items)

2. 使用文件锁控制并发访问

如果必须共享同一个文件(比如需要累加更新数据),可以通过文件锁机制确保同一时间只有一个进程操作文件:

  • Linux/macOS下可以用fcntl模块加锁
  • Windows下可以用win32file模块

示例代码(Linux/macOS):

import pandas as pd
import fcntl
from multiprocessing import Pool

def safe_read_pickle(file_path):
    with open(file_path, 'rb') as f:
        # 加共享读锁,允许多个进程同时读,但阻止写操作
        fcntl.flock(f, fcntl.LOCK_SH)
        df = pd.read_pickle(f)
        # 释放锁
        fcntl.flock(f, fcntl.LOCK_UN)
    return df

def safe_write_pickle(df, file_path):
    with open(file_path, 'wb') as f:
        # 加独占写锁,阻止其他进程读/写
        fcntl.flock(f, fcntl.LOCK_EX)
        pd.to_pickle(df, f)
        fcntl.flock(f, fcntl.LOCK_UN)

def my_function(item):
    df = safe_read_pickle("shared_data.pkl")
    # 处理逻辑
    df = df[df['value'] > item]
    safe_write_pickle(df, "shared_data.pkl")

if __name__ == "__main__":
    items = list(range(15))
    with Pool() as pool:
        pool.map(my_function, items)

3. 把文件读写移到主进程

如果数据量不大,可以在主进程统一读取所有数据,然后把DataFrame作为参数传递给子进程处理,处理完后再由主进程统一写入文件。这样完全避免子进程的文件操作:

import pandas as pd
from multiprocessing import Pool

def process_data(df_slice, item):
    # 只做数据处理,不涉及文件读写
    df_slice['processed'] = df_slice['col'] * item
    return df_slice

if __name__ == "__main__":
    # 主进程读取数据
    main_df = pd.read_pickle("all_data.pkl")
    items = list(range(15))
    
    # 拆分数据或者直接传递(根据你的处理逻辑调整)
    with Pool() as pool:
        # 这里假设每个item对应处理数据的一部分,或者传递整个df
        results = pool.starmap(process_data, [(main_df, item) for item in items])
    
    # 主进程合并结果并写入
    final_df = pd.concat(results)
    pd.to_pickle(final_df, "processed_all_data.pkl")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 08:23:10