You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.22 03:07:43