如何让pool.apply_async在提交新任务前等待当前任务完成
用apply_async实现内层循环并行且等待完成后进入外层下一次迭代
问题场景
原串行嵌套循环代码:
for i in range(N): for k in range(L): x[k] = Do_Something(x[k])
外层i循环的每次迭代依赖前一次的执行结果,无法并行;内层k循环的任务相互独立,属于易并行场景。希望用apply_async实现内层循环并行,同时保证每次外层i循环必须等待所有内层k任务完成后,再进入下一次i迭代,且不想每次创建新的multiprocessing.Pool(避免进程创建销毁的高开销)。
错误写法的问题
最初的写法会导致外层i循环不等待内层任务完成就继续执行,打乱依赖顺序:
pool = mp.Pool(Nworkers) for i in range(N): for k in range(L): pool.apply_async(Do_Something, args=(k), callback=getresults) pool.join() pool.close()
解决方案:收集AsyncResult并等待完成
核心思路是在每次外层i循环中,收集所有apply_async返回的AsyncResult对象,等待所有对象对应的任务完成后,再进入下一次i循环。
完整示例代码:
import multiprocessing as mp def Do_Something(k): # 模拟业务逻辑:处理对应位置的数据,返回(k, 处理后结果) processed_val = k * 2 # 替换为你的实际处理逻辑 return (k, processed_val) def getresults(result): # 按k的位置更新全局/共享的x列表,保证结果顺序正确 k, val = result x[k] = val if __name__ == "__main__": N = 5 # 外层循环次数 L = 10 # 内层循环次数 Nworkers = 4 # 进程池大小 x = [0] * L # 存储数据的列表 # 只初始化一次进程池,避免重复创建开销 pool = mp.Pool(Nworkers) for i in range(N): async_tasks = [] for k in range(L): # 注意args必须是元组,所以args=(k,) 末尾要加逗号 task = pool.apply_async(Do_Something, args=(k,), callback=getresults) async_tasks.append(task) # 等待当前外层循环的所有内层任务完成 for task in async_tasks: task.wait() # 或用task.get(),若不需要返回值用wait更轻量 # 所有内层任务完成,x已更新,可进入下一次外层循环 print(f"第{i}次外层循环执行完成") # 关闭进程池并等待所有任务结束 pool.close() pool.join()
关键说明
- 收集AsyncResult对象:每次外层循环中,把每个
apply_async返回的任务对象存入列表,用于后续等待。 - 等待任务完成:通过遍历任务列表调用
wait(),确保当前外层循环的所有内层任务都执行完毕,再进入下一次迭代。 - callback保证结果顺序:让
Do_Something返回(k, 处理结果),在getresults中根据k更新对应位置的x,即使任务完成顺序混乱,结果也能正确对应到原位置。 - 复用进程池:只初始化一次
Pool,避免每次外层循环创建销毁进程带来的性能开销。
内容的提问来源于stack exchange,提问作者patrick7
相关产品推荐
相关产品推荐

