使用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)
核心问题分析
所有问题的根源都是跨进程操作的竞态条件:多个进程同时读写共享资源时,没有同步机制保证操作的原子性,导致数据被覆盖或修改失效。
- Manager.Namespace的竞态:
pd.concat+赋值不是原子操作,多个Worker同时读取当前DataFrame、拼接新数据、再赋值,后执行的Worker会覆盖先执行的结果,丢失部分行。 - SharedMemory的无同步写入:多个进程同时对共享内存中的numpy数组执行
+=1,字节级的并发写入导致部分修改被覆盖,只有部分元素累加成功。 - 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
相关产品推荐
相关产品推荐

