如何将ProcessPoolExecutor的submit改为阻塞调用,等待可用工作进程?
实现ProcessPoolExecutor提交任务时阻塞等待可用进程
要实现你要的效果——每提交2个任务后等待5秒,再继续提交下一组,核心是控制任务提交的节奏,让新任务只在进程池有空闲工作进程时才提交。
下面是修改后的代码:
import concurrent.futures import multiprocessing import time def wait_f(): time.sleep(5) return 1 if __name__ == '__main__': multiprocessing.freeze_support() global_results = [] max_workers = 2 with concurrent.futures.ProcessPoolExecutor(max_workers=max_workers) as executor: futures = [] for j in range(10): # 当未完成任务数达到进程池上限时,等待第一个任务完成 if len(futures) >= max_workers: done, _ = concurrent.futures.wait(futures, return_when=concurrent.futures.FIRST_COMPLETED) # 处理完成的任务,移除已完成的future for future in done: global_results.append(future.result()) futures.remove(future) # 提交新任务并打印序号 future = executor.submit(wait_f) futures.append(future) print(j + 1) # 输出1、2、3...而非0、1、2... # 处理剩余未完成的任务 for future in concurrent.futures.as_completed(futures): global_results.append(future.result())
代码说明
- 用
max_workers变量统一管理进程池的工作数,方便后续调整。 - 在任务提交循环中,先检查未完成任务的数量:如果已经占满所有进程,就调用
concurrent.futures.wait等待第一个任务完成,空出进程后再提交新任务。 - 提交任务后立即打印对应序号,这样就能实现先输出1、2,等待5秒后输出3、4的效果。
- 循环结束后,处理所有剩余的未完成任务,确保所有结果都被收集。
内容的提问来源于stack exchange,提问作者Eric
相关产品推荐
相关产品推荐

