为何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之后的任务。请问这是什么原因?如何阻止这种情况发生?
任务提交后立即进入执行队列:调用
executor.submit()时,任务会被立即添加到线程池的任务队列中,工作线程会持续从队列中取出任务执行——不管你是否调用future.result()去获取结果。你的代码先一次性提交了10000个任务,之后才开始遍历futures处理结果,这意味着在处理到Order 4的异常前,线程池已经在后台执行了大量任务,甚至可能已经将后面的任务(比如4706)调度到工作线程中等待执行。shutdown只能取消未启动的任务:当你调用
executor.shutdown(wait=False, cancel_futures=True)时,它只会取消尚未被工作线程取出执行的任务。如果某个任务已经被工作线程从队列中取出(哪怕还没开始打印日志),这个任务会继续执行完毕,无法被取消。你看到的Order 4706就是这种情况——在触发shutdown之前,它已经被工作线程调度,所以会继续执行。遍历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

