Python多进程调用apply_async无法更新DataFrame且进程结束后仍运行问题
问题核心原因
- 第一,多进程内存空间相互隔离:你在子进程
f里修改的df是主进程df拷贝到子进程内存空间的副本,子进程对副本的任何修改都不会同步回主进程的原始df,所以你最后查看主进程的df肯定没有更新。 - 第二,多进程并发写同一个文件存在冲突:所有子进程都同时往
dftest.csv写数据,没有加锁的情况下会出现写入覆盖、写入中断的问题,这就是你只看到2个值的原因,严重时还会导致进程IO阻塞,内核一直挂着。 - 第三,你没有等待所有异步任务执行完成就结束主流程:
apply_async是异步提交任务,提交后主进程不会主动等待任务跑完,如果你后续没有调用pool.close()+pool.join(),主进程可能提前退出,或者进程池一直持有资源不释放,导致内核持续运行。 - 图像向量化操作本身不会触发多进程异常,只要你的
g函数里没有依赖不可序列化的全局变量、没有多线程冲突的逻辑,单进程能跑通的计算逻辑多进程也可以正常跑。
修正方案
你不需要在子进程里修改df或者写文件,直接让子进程返回计算结果,主进程收集完所有结果后统一更新df、写文件即可,修正后的代码逻辑如下:
1. 调整子进程执行函数,只返回计算结果
import multiprocessing as mp import pandas as pd def g(path1, path2): # 你的原有图像向量化、计算逻辑 return value # float类型计算结果 def f(row, index): # 仅做计算,返回索引和对应结果,不需要操作df或者写文件 value = g(row['image1'], row['image2']) return (index, value)
2. 主进程收集结果统一处理
# spawn模式必须加这个入口判断,否则会循环启动进程报错 if __name__ == '__main__': # 这里写你加载原始df的逻辑 df['value'] = 0.0 # 提前创建结果存储列 N_CORES = mp.cpu_count() with mp.get_context("spawn").Pool(N_CORES) as pool: # 批量提交所有任务 tasks = [pool.apply_async(f, args=(row, idx)) for idx, row in df.iterrows()] # 等待所有任务完成,获取结果更新主进程df for task in tasks: idx, val = task.get() df.at[idx, 'value'] = val # 所有结果更新完成后统一写文件,避免并发写冲突 df.to_csv('dftest.csv', index=False)
如果你的数据量非常大,iterrows遍历效率较低,可以把df按行转成参数列表,用pool.map或者pool.starmap来提交任务,执行效率会更高。
内容的提问来源于stack exchange,提问作者hbillie
相关产品推荐
相关产品推荐

