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)
关键说明:
- 实时进度反馈:通过
as_completed遍历完成的任务,每完成一个就更新进度计数并打印,无需轮询等待。 - 异常即时处理:调用
future.result()时会立即触发子进程的异常,此时调用executor.shutdown(wait=False)会立即终止所有未完成的进程,避免继续执行新任务。 - 结果顺序保留:通过
futures字典关联每个future和任务的原始索引,将结果存入对应位置的列表,确保最终DataFrame的顺序和原df_options一致。 - 进程池关闭:异常时用
wait=False强制终止;正常完成时用wait=True等待所有进程结束。
内容的提问来源于stack exchange,提问作者Jimbo
相关产品推荐
相关产品推荐

