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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 12:36:04