如何退出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')
核心改动说明
- 按需提交任务:用迭代器逐个获取文件,只在进程池有空闲时才提交新任务,避免一开始就把所有任务塞进队列
- 及时取消任务:找到预期结果后,立刻停止提交新任务,同时调用
task.cancel()取消所有还没开始执行的待处理任务 - 终止循环逻辑:通过
found标记跳出循环,避免继续处理后续任务
另外补充:原代码里在with块内调用pool.shutdown()是多余的——with语句结束时会自动调用shutdown(),而且此时所有任务已经提交,所以不会起到阻止新任务的作用。
内容的提问来源于stack exchange,提问作者Jellby
相关产品推荐
相关产品推荐

