如何在不关闭Python multiprocessing.pool.Pool的情况下等待工作进程完成?
修改多进程基准测试脚本,确保等待所有任务完成
问题分析
当前par函数仅向进程池提交任务就返回,导致基准测试只统计了任务提交耗时,而非任务执行的完整耗时。需要修改par函数,使其在返回前等待所有工作进程完成任务,同时保留进程池的复用能力。
解决方案
修改par函数,收集所有异步任务的AsyncResult对象,然后等待所有任务执行完成。以下是修改后的完整代码:
import numpy as np from multiprocessing import Pool import timeit as ti def foo(n): return -np.sort(-np.arange(n))[-1] def par(reps, bigNum, pool): # 收集所有异步任务的结果对象 results = [] for i in range(bigNum, bigNum + reps): res = pool.apply_async(foo, args=(i,)) results.append(res) # 等待所有任务完成 for res in results: res.get() def ser(reps, bigNum): for i in range(bigNum, bigNum + reps): foo(i) if __name__ == '__main__': bigNum = 9_000_000 reps = 6 fun = f'par(reps, bigNum, pool)' t = 1000 * np.array(ti.repeat(stmt=fun, setup='pool=Pool(reps)', globals=globals(), number=1, repeat=10)) print(f'{fun}: {np.amin(t):6.3f}ms {np.median(t):6.3f}ms') fun = f'ser(reps, bigNum)' t = 1000 * np.array(ti.repeat(stmt=fun, setup='', globals=globals(), number=1, repeat=10)) print(f'{fun}: {np.amin(t):6.3f}ms {np.median(t):6.3f}ms')
关键修改点说明
- 收集任务结果对象:在
par函数中,每次调用pool.apply_async后,将返回的AsyncResult对象存入列表,追踪每个任务的状态。 - 等待任务完成:遍历结果列表,调用每个对象的
get()方法——该方法会阻塞直到对应任务完成并返回结果,确保所有任务执行完毕后par函数才返回。 - 优化timeit的setup:原setup中额外执行了一次
par函数,会干扰测试结果,修改后仅在setup中创建进程池,保证测试准确性。
补充说明
- 使用
res.get()可同时捕获任务执行中抛出的异常(如果存在)。 - 进程池
pool未被关闭或终止,多次调用par函数时会被复用,符合需求。
内容的提问来源于stack exchange,提问作者Paul Jurczak
相关产品推荐
相关产品推荐

