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

Python中任一进程出错时如何终止异步starmap多进程池?

问题:multiprocessing starmap_async异常延迟抛出,迁移到concurrent.futures并保留实时进度反馈

使用Python的multiprocessing.Pool.starmap_async时,发现子进程抛出的异常要等所有进程执行完毕才会触发。当前代码中,若某个子进程出错,其他进程不仅会完成当前任务,还会启动新任务,无法及时终止。希望迁移到concurrent.futures实现异常触发后立即终止所有进程,同时保留原有的实时进度反馈功能。

原代码如下:

from multiprocessing import Pool, cpu_count
import datetime
import itertools
import time

with Pool(max(cpu_count()//2, 1)) as p:
    df_iter = df_options.iterrows()
    ir = itertools.repeat
    results = p.starmap_async(_run, zip(df_iter, ir(fixed_options), ir(outputs_grab)), chunksize=1)
    p.close()  # no more jobs to submit
    
    # Printing progress
    n_remaining = results._number_left + 1
    while (not results.ready()):
        time.sleep(1)
        # Check for errors here ... How ????
        # Then what? call terminate()????? 
        if verbose:
            if results._number_left < n_remaining:
                now = datetime.datetime.now()
                n_remaining = results._number_left
                print('%d/%d  %s' % (n_remaining, n_rows, str(now)[11:]))
                
    print('joining')
    p.join()
    
    all_results = results.get()

df = pd.DataFrame(all_results)

解决方案:使用concurrent.futures.ProcessPoolExecutor

concurrent.futures的ProcessPoolExecutor配合as_completed方法,可以逐个获取完成的任务结果,一旦捕获到异常就能立即终止所有进程,同时轻松实现实时进度反馈。

迁移后的代码示例:

from concurrent.futures import ProcessPoolExecutor, as_completed
import datetime
import itertools
import time
import pandas as pd

def main():
    total_tasks = len(df_options)
    completed_tasks = 0
    all_results = [None] * total_tasks  # 用列表保存结果,保持原顺序
    
    # 创建进程池,数量和原代码一致
    with ProcessPoolExecutor(max_workers=max(cpu_count()//2, 1)) as executor:
        # 提交所有任务,保存future和对应的索引
        futures = {}
        for idx, (row, _, _) in enumerate(zip(df_options.iterrows(), itertools.repeat(fixed_options), itertools.repeat(outputs_grab))):
            future = executor.submit(_run, row, fixed_options, outputs_grab)
            futures[future] = idx
        
        # 遍历已完成的任务
        for future in as_completed(futures):
            task_idx = futures[future]
            try:
                # 获取任务结果,若有异常会在此抛出
                result = future.result()
                all_results[task_idx] = result
                completed_tasks += 1
                
                # 实时打印进度
                if verbose:
                    now = datetime.datetime.now()
                    print(f'{completed_tasks}/{total_tasks}  {str(now)[11:]}')
            except Exception as e:
                # 捕获到异常,立即终止所有进程
                executor.shutdown(wait=False)
                # 重新抛出异常,终止主进程
                raise e
        
        # 正常完成所有任务后关闭进程池
        executor.shutdown(wait=True)
    
    # 整理结果为DataFrame
    df = pd.DataFrame(all_results)

关键说明:

  1. 实时进度反馈:通过as_completed遍历完成的任务,每完成一个就更新进度计数并打印,无需轮询等待。
  2. 异常即时处理:调用future.result()时会立即触发子进程的异常,此时调用executor.shutdown(wait=False)会立即终止所有未完成的进程,避免继续执行新任务。
  3. 结果顺序保留:通过futures字典关联每个future和任务的原始索引,将结果存入对应位置的列表,确保最终DataFrame的顺序和原df_options一致。
  4. 进程池关闭:异常时用wait=False强制终止;正常完成时用wait=True等待所有进程结束。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.24 23:45:40