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

如何退出as_completed池执行器循环并取消剩余未启动任务?

问题

需求是运行一组任务,任务完成时立即处理结果,一旦有任务返回预期结果,就取消剩余任务——允许等待已启动的任务完成,但禁止启动新任务。

尝试调用pool.shutdown(cancel_futures=True)但无效,代码如下:

import time
import multiprocessing
import concurrent.futures
import random

filelist = ['001', '002', '003', '004', '005', '006', '007', '008', '009', '010']
parallel = 2

def worker(name, lock, index):
    with lock:
        index.value += 1
        pc = index.value*10
        print(f'starting: {name} ({pc})')
    d = random.randrange(1, 5)
    time.sleep(d)
    return pc == 60

print('INIT')

with multiprocessing.Manager() as manager:
    lock = manager.Lock()
    index = manager.Value('b', 0)
    if parallel > 0:
        with concurrent.futures.ProcessPoolExecutor(max_workers=parallel) as pool:
            tasks = [pool.submit(worker, filename, lock, index) for filename in filelist]
            for task in concurrent.futures.as_completed(tasks):
                rc = task.result()
                if rc:
                    print('Found!')
                    pool.shutdown(cancel_futures=True)
                    break
    else:
        for filename in filelist:
            rc = worker(filename, lock, index)
            if rc:
                print('Found!')
                break

print('END')

串行模式(parallel = 0)运行结果符合预期:

INIT
starting: 001 (10)
starting: 002 (20)
starting: 003 (30)
starting: 004 (40)
starting: 005 (50)
starting: 006 (60)
Found!
END

但并行模式下程序会运行到所有任务结束:

INIT
starting: 001 (10)
starting: 002 (20)
starting: 003 (30)
starting: 004 (40)
starting: 005 (50)
starting: 006 (60)
Found!
starting: 007 (70)
starting: 008 (80)
starting: 009 (90)
END

请问有没有办法退出as_completed循环并阻止新任务启动?

解决方案

问题根源是代码一开始就用[pool.submit(...)]一次性提交了所有任务到进程池的任务队列里。pool.shutdown(cancel_futures=True)只能取消还没开始执行的任务,但进程池会按照max_workers的数量,逐步把队列里的任务调度到空闲进程执行——所以即使调用了shutdown,队列里的剩余任务还是会被启动。

要实现需求,必须改变任务提交方式:不要一次性提交所有任务,而是按需提交,同时跟踪已提交的任务,找到目标结果后立刻停止提交新任务,并取消队列中未启动的任务。

修改后的代码如下:

import time
import multiprocessing
import concurrent.futures
import random

filelist = ['001', '002', '003', '004', '005', '006', '007', '008', '009', '010']
parallel = 2

def worker(name, lock, index):
    with lock:
        index.value += 1
        pc = index.value*10
        print(f'starting: {name} ({pc})')
    d = random.randrange(1, 5)
    time.sleep(d)
    return pc == 60

print('INIT')

with multiprocessing.Manager() as manager:
    lock = manager.Lock()
    index = manager.Value('b', 0)
    if parallel > 0:
        with concurrent.futures.ProcessPoolExecutor(max_workers=parallel) as pool:
            # 把文件列表转为迭代器,按需取任务
            task_iter = iter(filelist)
            # 先提交初始的max_workers个任务,填满进程池
            pending = []
            for _ in range(parallel):
                try:
                    filename = next(task_iter)
                    pending.append(pool.submit(worker, filename, lock, index))
                except StopIteration:
                    break
            
            found = False
            while pending and not found:
                # 遍历已完成的任务
                for task in concurrent.futures.as_completed(pending):
                    rc = task.result()
                    pending.remove(task)
                    if rc:
                        print('Found!')
                        found = True
                        break
                    # 没找到结果,继续提交下一个任务
                    try:
                        next_file = next(task_iter)
                        pending.append(pool.submit(worker, next_file, lock, index))
                    except StopIteration:
                        pass
                # 找到结果后,取消所有待执行的任务
                if found:
                    for task in pending:
                        task.cancel()
    else:
        for filename in filelist:
            rc = worker(filename, lock, index)
            if rc:
                print('Found!')
                break

print('END')

核心改动说明

  1. 按需提交任务:用迭代器逐个获取文件,只在进程池有空闲时才提交新任务,避免一开始就把所有任务塞进队列
  2. 及时取消任务:找到预期结果后,立刻停止提交新任务,同时调用task.cancel()取消所有还没开始执行的待处理任务
  3. 终止循环逻辑:通过found标记跳出循环,避免继续处理后续任务

另外补充:原代码里在with块内调用pool.shutdown()是多余的——with语句结束时会自动调用shutdown(),而且此时所有任务已经提交,所以不会起到阻止新任务的作用。

内容的提问来源于stack exchange,提问作者Jellby

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 11:02:06