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

ThreadPoolExecutor无法中断:shutdown参数配置失效问题排查

问题

执行以下Python脚本时,按下Ctrl+C或在PyCharm中触发Interrupt(等同于KeyboardInterrupt)后,脚本仍持续输出内容,需连续按三次Ctrl+C才能停止。配置shutdown(wait=False,cancel_futures=True)为何无法生效?尝试过shutdown和Event机制均无效,希望触发KeyboardInterrupt后立即关闭所有线程。

import urllib3
from urllib3.exceptions import HTTPError, URLError
from concurrent.futures import ThreadPoolExecutor
from datetime import datetime

urls = ["facebook.com","google.com"]
user_agent = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/71.0.3578.98 Safari/537.36"

def build_list(directory,url):
    urllist = []
    with open(directory) as file:
        for line in file:
            line = line.rstrip()
            current_url = f"{url}/{line}"
            urllist.append(current_url)
        return urllist

def process_directories(request_url):
    http = urllib3.PoolManager()
    head = {"User-Agent": user_agent}
    try:
        response = http.request("GET", request_url, headers=head)
        print(f"Status for {request_url}: {response.status}")
        return response.status
    except (URLError, HTTPError) as e:
        if hasattr(e, 'code') and e.code != 404:
            print(f"!!!!! [{e.code}] ==> {request_url}")
        pass
    except Exception as e:
        print(f"Unexpected error: {str(e)}")

def task(url):
    urllist = build_list("../dirlist.txt", url)
    print(f"Urllist Built For: {url}")
    for url in urllist:
        process_directories(url)

if __name__ == "__main__":
    with ThreadPoolExecutor(max_workers=10) as pool:
        futures = [
            pool.submit(task, url) for url in urls
        ]

        for future in futures:
            print([f._state for f in futures])
            if future.cancelled():
                continue
            try:
                n = future.result()
                print(f"{datetime.now()} - {n}")
            except KeyboardInterrupt as e:
                print(f"{datetime.now()} - EXCEPTION! {e}")
                pool.shutdown(wait=False,cancel_futures=True)
    print(f"{datetime.now()} - Run complete")
原因分析
  1. cancel_futures=True的局限性:这个参数仅能取消尚未启动的任务future。对于已经在运行的线程(比如task函数已经开始遍历URL列表并发起请求),Python无法强制终止正在执行的线程,该参数对这类线程完全无效。
  2. 中断响应滞后:主线程的future.result()是阻塞调用,只有当某个任务执行完成后才会进入下一轮循环。如果所有线程都处于运行状态,主线程会一直卡在result()调用上,无法及时响应Ctrl+C触发的中断。
  3. IO请求无超时:process_directories中的http.request("GET")是同步阻塞操作,一旦发起请求,线程会一直等待服务器响应,即使主线程触发中断,这些阻塞的IO操作也不会立刻停止,导致脚本继续输出内容。
解决方案

要实现触发中断后立即停止所有线程,需从三个层面优化:

1. 给IO请求添加超时限制

修改process_directories,为HTTP请求添加超时,避免线程无限阻塞:

def process_directories(request_url):
    http = urllib3.PoolManager()
    head = {"User-Agent": user_agent}
    try:
        # 设置5秒超时,可根据需求调整
        response = http.request("GET", request_url, headers=head, timeout=5.0)
        print(f"Status for {request_url}: {response.status}")
        return response.status
    except (URLError, HTTPError) as e:
        if hasattr(e, 'code') and e.code != 404:
            print(f"!!!!! [{e.code}] ==> {request_url}")
        pass
    except Exception as e:
        print(f"Unexpected error: {str(e)}")

2. 用全局标志位终止运行中的任务

添加线程安全的全局中断标志,在task的循环中检查该标志,一旦触发就停止执行:

import threading

# 全局中断标志,线程安全
interrupted = threading.Event()

def task(url):
    urllist = build_list("../dirlist.txt", url)
    print(f"Urllist Built For: {url}")
    for url in urllist:
        # 每次循环前检查是否触发中断
        if interrupted.is_set():
            print(f"Stopping task for {url}")
            break
        process_directories(url)

3. 主线程优化中断捕获逻辑

改用as_completed遍历future,避免阻塞在单个任务的result()调用上,同时在捕获中断时设置全局标志:

if __name__ == "__main__":
    try:
        with ThreadPoolExecutor(max_workers=10) as pool:
            futures = [pool.submit(task, url) for url in urls]
            
            from concurrent.futures import as_completed
            for future in as_completed(futures):
                print([f._state for f in futures])
                try:
                    n = future.result()
                    print(f"{datetime.now()} - {n}")
                except Exception as e:
                    print(f"{datetime.now()} - Task error: {str(e)}")
    except KeyboardInterrupt as e:
        print(f"{datetime.now()} - EXCEPTION! {e}")
        interrupted.set()
    print(f"{datetime.now()} - Run complete")

可选:激进的线程终止方式(不推荐)

如果上述方法仍无法满足需求,Python没有原生的强制终止线程的API,但可以通过ctypes调用底层方法实现。这种方式可能导致资源泄漏(如未关闭的HTTP连接),仅在极端场景下使用:

import ctypes
import inspect

def terminate_thread(thread):
    """强制终止指定线程"""
    if not thread.is_alive():
        return
    # 获取线程ID
    tid = thread.ident
    # 调用底层API抛出SystemExit异常终止线程
    exc = ctypes.py_object(SystemExit)
    res = ctypes.pythonapi.PyThreadState_SetAsyncExc(ctypes.c_long(tid), exc)
    if res == 0:
        raise ValueError("无效的线程ID")
    elif res != 1:
        # 清理异常状态
        ctypes.pythonapi.PyThreadState_SetAsyncExc(ctypes.c_long(tid), None)
        raise SystemError("线程终止失败")

# 在主线程捕获中断时,遍历线程池的工作线程并终止
# 注意:ThreadPoolExecutor的工作线程存储在内部属性中,依赖CPython实现,可能因版本变化失效
# if __name__ == "__main__":
#     try:
#         pool = ThreadPoolExecutor(max_workers=10)
#         # ... 提交任务逻辑 ...
#     except KeyboardInterrupt:
#         for thread in pool._threads:
#             terminate_thread(thread)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 07:27:06