如何在Python concurrent.futures所有任务完成后执行指定后续任务
完全可以用concurrent.futures的原生能力实现,不需要额外引入其他线程逻辑,有两种常见实现方式:
方案1:利用上下文管理器的自动等待特性(最简便)
ProcessPoolExecutor的with上下文管理器会在退出代码块时,隐式调用executor.shutdown(wait=True),这个操作会阻塞直到所有提交的任务全部执行完成、资源释放完成才会继续向下运行。你只需要把所有后续处理逻辑放到with代码块外面即可。
示例代码:
import os import concurrent.futures from multiprocessing import Manager # 进程间共享数据需要用Manager提供的结构 def do_the_work(item, shared_dict): # 你的业务逻辑,写入共享数据 key, val = item shared_dict[key] = val * 2 return True if __name__ == '__main__': # 初始化进程安全的共享数据结构 with Manager() as manager: shared_data = manager.dict() work_list = {"task1": 1, "task2": 2, "task3": 3} with concurrent.futures.ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: futures = [executor.submit(do_the_work, item, shared_data) for item in work_list.items()] for i, future in enumerate(concurrent.futures.as_completed(futures)): status = future.result() print(f'DONE: count:{i} result:{status}') # 走到此处所有future已经全部执行完成,直接处理共享数据即可 print("所有任务执行完毕,共享数据内容:", dict(shared_data)) # 此处写你的后续处理逻辑
方案2:手动调用shutdown方法
如果你不用上下文管理器的写法,也可以手动调用executor.shutdown(wait=True)实现等待效果:
if __name__ == '__main__': from multiprocessing import Manager with Manager() as manager: shared_data = manager.dict() work_list = {"task1": 1, "task2": 2, "task3": 3} executor = concurrent.futures.ProcessPoolExecutor(max_workers=os.cpu_count()) futures = [executor.submit(do_the_work, item, shared_data) for item in work_list.items()] for i, future in enumerate(concurrent.futures.as_completed(futures)): status = future.result() print(f'DONE: count:{i} result:{status}') # wait=True表示阻塞等待所有提交的任务执行完成 executor.shutdown(wait=True) # 后续处理逻辑 print("所有任务执行完毕,共享数据内容:", dict(shared_data))
额外注意点
因为你用的是进程池,普通Python对象无法跨进程共享,你的共享数据结构必须使用multiprocessing模块提供的Manager、Queue、Array等进程安全的共享结构,否则子进程的写入操作不会同步到主进程中。
内容的提问来源于stack exchange,提问作者Exploring
相关产品推荐
相关产品推荐

