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

使用ProcessPoolExecutor在Worker间共享Pandas DataFrame的问题

多进程共享Pandas DataFrame与并发更新问题解决方案

问题概述

使用ProcessPoolExecutor时遇到多Worker间共享Pandas DataFrame的需求:所有Worker仅加载一次CSV数据,且能在新增文件时更新DataFrame。但测试中出现以下问题:

  • 用multiprocessing.Manager共享DataFrame时,最终行数远少于预期(65-80行,预期100行),且每次结果不同
  • 用SharedMemory共享numpy数组时,仅部分元素被正确累加
  • 用ShareableList做计数时,1000次迭代后结果少1次(仅999)

核心问题分析

所有问题的根源都是跨进程操作的竞态条件:多个进程同时读写共享资源时,没有同步机制保证操作的原子性,导致数据被覆盖或修改失效。

  1. Manager.Namespace的竞态:pd.concat+赋值不是原子操作,多个Worker同时读取当前DataFrame、拼接新数据、再赋值,后执行的Worker会覆盖先执行的结果,丢失部分行。
  2. SharedMemory的无同步写入:多个进程同时对共享内存中的numpy数组执行+=1,字节级的并发写入导致部分修改被覆盖,只有部分元素累加成功。
  3. ShareableList的计数丢失:多个进程同时读取并修改列表元素,两次读取相同值后各自加1,最终只生效一次,导致计数少1。

解决方案

1. 修复Manager共享DataFrame的竞态问题:加跨进程锁

通过Manager.Lock同步读写操作,保证每次只有一个Worker修改DataFrame:

import pandas as pd
import multiprocessing
from concurrent.futures import ProcessPoolExecutor
import os

def process(par):
    x, shared, lock = par
    new_data = pd.DataFrame([[x, int(os.getpid())]], columns=['A', 'B'])
    # 加锁确保读写原子性
    with lock:
        shared.df = pd.concat((shared.df, new_data), ignore_index=True)

if __name__ == "__main__":
    with multiprocessing.Manager() as manager:
        shared = manager.Namespace()
        shared.df = pd.DataFrame()
        lock = manager.Lock()  # 跨进程有效锁

        with ProcessPoolExecutor(max_workers=os.cpu_count()) as exe:
            nx = range(0, 100)
            par = [[x, shared, lock] for x in nx]
            exe.map(process, par)
        
        print("最终行数:", len(shared.df))  # 稳定输出100行

2. 修复SharedMemory的数组累加问题:加锁保护

用multiprocessing.Lock同步共享内存的写入操作:

import numpy as np
import multiprocessing
from concurrent.futures import ProcessPoolExecutor
from multiprocessing.shared_memory import SharedMemory
import os

def worker_function(args):
    name, lock = args
    existing_shm = SharedMemory(name=name)
    c = np.ndarray((6,), dtype=np.int64, buffer=existing_shm.buf)
    with lock:
        c += 1
    existing_shm.close()

if __name__ == '__main__':
    a = np.array([1, 1, 2, 3, 5, 8])
    shm = SharedMemory(create=True, size=a.nbytes)
    b = np.ndarray(a.shape, dtype=a.dtype, buffer=shm.buf)
    b[:] = a[:]

    lock = multiprocessing.Lock()
    with ProcessPoolExecutor(max_workers=os.cpu_count()) as exe:
        exe.map(worker_function, [(shm.name, lock)]*100)

    print(b)  # 稳定输出[101, 101, 102, 103, 105, 108]
    shm.unlink()

3. 修复ShareableList的计数问题:加锁同步

用锁保证每次只有一个进程修改列表元素:

import multiprocessing
from concurrent.futures import ProcessPoolExecutor
from multiprocessing.managers import SharedMemoryManager
import os

def worker_function(args):
    sl, lock = args
    with lock:
        sl[0] = sl[0] + 1

if __name__ == '__main__':
    with SharedMemoryManager() as smm:
        init = [0, 1, 2, 3, 4, 5, 6]
        sl = smm.ShareableList(init)   
        print("初始值:", sl)

        lock = multiprocessing.Lock()
        with ProcessPoolExecutor(max_workers=os.cpu_count()) as exe:
            exe.map(worker_function, [(sl, lock)]*1000 )

        print("最终值:", sl)  # 稳定输出[1000,1,2,3,4,5,6]

4. 高效共享Pandas DataFrame的最优方案

只读场景(仅加载一次供所有Worker读取)

利用进程池initializer在每个Worker启动时加载一次DataFrame,避免跨进程共享的开销:

import pandas as pd
from concurrent.futures import ProcessPoolExecutor
import os

# 每个Worker进程的全局变量,启动时初始化
worker_df = None

def init_worker(csv_path):
    global worker_df
    # 每个Worker启动时加载一次CSV
    worker_df = pd.read_csv(csv_path)

def process_task(x):
    # 使用本地Worker内存中的DataFrame处理
    pid = os.getpid()
    result = worker_df[worker_df['A'] == x]
    return (pid, result.shape[0])

if __name__ == "__main__":
    csv_path = "data.csv"
    with ProcessPoolExecutor(
        max_workers=os.cpu_count(),
        initializer=init_worker,
        initargs=(csv_path,)
    ) as exe:
        tasks = range(10)
        results = exe.map(process_task, tasks)
        for res in results:
            print(f"进程{res[0]}处理结果:{res[1]}行")

读写场景(需要动态更新DataFrame)

优先选择Manager+Lock方案,若数据量较大,推荐用SQLite等轻量数据库作为中间存储,避免内存共享的竞态和开销。

总结

  • 跨进程共享可变数据时,必须加锁保证操作原子性,否则会出现数据丢失或错误
  • 只读场景下,用进程池initializer在Worker本地加载数据,比跨进程共享更高效
  • 读写场景下,Manager+Lock是最直接的内存共享方案,外部存储适合大规模数据更新

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 22:14:56