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

为何ThreadPoolExecutor会提前执行未触发的Futures?

问题:ThreadPoolExecutor为何会在futures被触发前提前执行?

代码示例

import concurrent.futures
import time

def sleep_test(order_number):
    num_seconds = 0.5
    print(f"Order {order_number} - Sleeping {num_seconds} seconds")
    time.sleep(num_seconds)
    print(f"Order {order_number} - Slept {num_seconds} seconds")
    if order_number == 4:
        raise Exception("Reached order #4")

def main():
    order_numbers = [i for i in range(10_000)]
    max_number_of_threads = 2

    with concurrent.futures.ThreadPoolExecutor(max_workers=max_number_of_threads) as executor:
        futures = []
        for order in order_numbers:
            futures.append(executor.submit(sleep_test, order_number=order))

        for future in futures:
            if future.cancelled():
                continue
            try:
                _ = future.result()
            except Exception:
                print("Caught Exception, stopping all future orders")
                executor.shutdown(wait=False, cancel_futures=True)

if __name__ == "__main__":
    main()

执行输出示例

$ python3 thread_pool_test.py
Order 0 - Sleeping 0.5 seconds
Order 1 - Sleeping 0.5 seconds
Order 0 - Slept 0.5 seconds
Order 1 - Slept 0.5 seconds
Order 2 - Sleeping 0.5 seconds
Order 3 - Sleeping 0.5 seconds
Order 2 - Slept 0.5 seconds
Order 4 - Sleeping 0.5 seconds
Order 3 - Slept 0.5 seconds
Order 5 - Sleeping 0.5 seconds
Order 4 - Slept 0.5 seconds
Order 6 - Sleeping 0.5 seconds
Caught Exception, stopping all future orders
Order 5 - Slept 0.5 seconds
Order 4706 - Sleeping 0.5 seconds
Order 6 - Slept 0.5 seconds
Order 4706 - Slept 0.5 seconds

执行过程中出现Order 4706被调用,不符合预期。本以为线程会在触发异常的Order 5或6左右停止,但有时脚本运行正常,有时却会调用数千个futures之后的任务。请问这是什么原因?如何阻止这种情况发生?


原因分析
  1. 任务提交后立即进入执行队列:调用executor.submit()时,任务会被立即添加到线程池的任务队列中,工作线程会持续从队列中取出任务执行——不管你是否调用future.result()去获取结果。你的代码先一次性提交了10000个任务,之后才开始遍历futures处理结果,这意味着在处理到Order 4的异常前,线程池已经在后台执行了大量任务,甚至可能已经将后面的任务(比如4706)调度到工作线程中等待执行。

  2. shutdown只能取消未启动的任务:当你调用executor.shutdown(wait=False, cancel_futures=True)时,它只会取消尚未被工作线程取出执行的任务。如果某个任务已经被工作线程从队列中取出(哪怕还没开始打印日志),这个任务会继续执行完毕,无法被取消。你看到的Order 4706就是这种情况——在触发shutdown之前,它已经被工作线程调度,所以会继续执行。

  3. 遍历futures的顺序不是任务完成顺序:你按提交顺序遍历futures调用result(),但任务的完成顺序可能和提交顺序不一致。Order 4的任务抛出异常时,Order 5、6甚至更后面的任务可能已经被工作线程启动了,这些任务都会继续执行。


解决方案

方案1:使用concurrent.futures.as_completed()处理任务

as_completed()会返回已完成的future(不管成功还是失败),这样你可以在第一个异常出现时立即停止提交新任务,并取消未启动的任务,避免大量任务提前进入执行队列。

修改后的代码:

import concurrent.futures
import time

def sleep_test(order_number):
    num_seconds = 0.5
    print(f"Order {order_number} - Sleeping {num_seconds} seconds")
    time.sleep(num_seconds)
    print(f"Order {order_number} - Slept {num_seconds} seconds")
    if order_number == 4:
        raise Exception("Reached order #4")

def main():
    order_numbers = [i for i in range(10_000)]
    max_number_of_threads = 2

    with concurrent.futures.ThreadPoolExecutor(max_workers=max_number_of_threads) as executor:
        # 逐个提交任务并加入集合,避免一次性提交所有
        futures = {executor.submit(sleep_test, order): order for order in order_numbers}
        
        for future in concurrent.futures.as_completed(futures):
            try:
                _ = future.result()
            except Exception:
                print("Caught Exception, stopping all future orders")
                executor.shutdown(wait=False, cancel_futures=True)
                return  # 立即退出循环,不再处理后续任务

if __name__ == "__main__":
    main()

方案2:控制任务提交速率,避免一次性塞满队列

如果不需要提交所有任务,或者希望更精细控制,可以在提交任务时配合as_completed(),只提交一定数量的任务,后续根据完成情况补充提交,这样在异常发生时,未提交的任务不会进入队列:

import concurrent.futures
import time

def sleep_test(order_number):
    num_seconds = 0.5
    print(f"Order {order_number} - Sleeping {num_seconds} seconds")
    time.sleep(num_seconds)
    print(f"Order {order_number} - Slept {num_seconds} seconds")
    if order_number == 4:
        raise Exception("Reached order #4")

def main():
    order_numbers = [i for i in range(10_000)]
    max_number_of_threads = 2
    stop_flag = False

    with concurrent.futures.ThreadPoolExecutor(max_workers=max_number_of_threads) as executor:
        # 先提交初始批次的任务
        futures = {executor.submit(sleep_test, order): order for order in order_numbers[:max_number_of_threads]}
        remaining_orders = iter(order_numbers[max_number_of_threads:])

        while futures and not stop_flag:
            for future in concurrent.futures.as_completed(futures):
                try:
                    _ = future.result()
                    # 提交下一个任务(如果还有剩余)
                    try:
                        next_order = next(remaining_orders)
                        new_future = executor.submit(sleep_test, next_order)
                        futures[new_future] = next_order
                    except StopIteration:
                        pass
                except Exception:
                    print("Caught Exception, stopping all future orders")
                    executor.shutdown(wait=False, cancel_futures=True)
                    stop_flag = True
                    break
            del futures[future]

if __name__ == "__main__":
    main()

关键说明

  • 两种方案的核心都是避免一次性提交所有任务,或者尽早处理完成的任务,这样在异常发生时,未提交或未被调度的任务可以被及时取消。
  • shutdown(cancel_futures=True)只能取消那些还在队列中等待的任务,已经被工作线程取出的任务会继续执行,这是线程池的设计特性,无法强制终止正在执行的线程(除非使用更底层的线程API,但不推荐,可能导致资源泄漏)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 11:25:16