基于APScheduler的多线程下载器:如何在指定时间终止所有线程?
如何在指定时间终止多线程下载任务(APScheduler + ThreadPoolExecutor)
我用APScheduler开发了一个下载器,要求凌晨12点启动下载,早上6点前停止。目前已经实现了12点自动启动下载,但6点的终止操作只能退出当前线程,ThreadPoolExecutor里的下载线程还在继续运行。以下是我尝试的代码:
import multiprocessing from concurrent.futures import ThreadPoolExecutor, as_completed from apscheduler.schedulers.background import BlockingScheduler executor = ThreadPoolExecutor(max_workers=multiprocessing.cpu_count() * 5) urls = [...] # 待下载的URL列表 def download(url): ... # 下载逻辑 def main_download(): futures = [executor.submit(download, url) for url in urls] for future in as_completed(futures): ... # 处理下载完成后的逻辑 scheduler = BlockingScheduler(timezone="Asia/Kolkata") job = scheduler.add_job(main_download, trigger="cron", hour=12) def kill_all(): # 尝试终止所有任务 job.remove() scheduler.remove_all_jobs() scheduler.shutdown() quit(1) # 已经试过exit、抛出KeyboardInterrupt、sys.exit,都没用 scheduler.add_job(kill_all, trigger="cron", hour=6) # 6点执行终止 scheduler.start()
请问有没有可靠的方法能终止所有线程?目前的操作只能退出当前线程,其他下载线程还在跑。
解决方案:用事件(Event)优雅终止线程
Python不推荐强制终止线程(易引发资源泄漏、数据损坏),更可靠的方式是给线程传递终止信号,让线程主动退出。可以用threading.Event实现:
步骤1:定义全局终止信号
创建一个Event对象,用来标记是否需要停止所有下载任务:
import threading stop_event = threading.Event()
步骤2:修改下载函数,定期检查终止信号
在download函数的关键节点(比如循环写入数据、网络请求间隙)检查stop_event,如果信号触发就立即退出:
def download(url): try: # 示例:以requests流式下载为例 import requests response = requests.get(url, stream=True) with open(f"./{url.split('/')[-1]}", 'wb') as f: for chunk in response.iter_content(chunk_size=1024): # 每次写入前检查终止信号 if stop_event.is_set(): print(f"终止下载:{url}") return f.write(chunk) except Exception as e: print(f"下载出错:{url},错误信息:{e}")
步骤3:修改kill_all函数,触发信号并等待线程结束
在kill_all里,先触发终止信号,然后关闭线程池并等待所有任务完成,最后关闭调度器:
def kill_all(): # 触发终止信号,通知所有下载线程停止 stop_event.set() # 关闭线程池,不再接受新任务,等待已提交的任务完成或终止 executor.shutdown(wait=True) # 清理调度器任务并关闭 scheduler.remove_all_jobs() scheduler.shutdown()
完整修改后的代码
import multiprocessing import threading from concurrent.futures import ThreadPoolExecutor, as_completed from apscheduler.schedulers.background import BlockingScheduler # 全局终止信号 stop_event = threading.Event() executor = ThreadPoolExecutor(max_workers=multiprocessing.cpu_count() * 5) urls = [...] # 待下载的URL列表 def download(url): try: import requests response = requests.get(url, stream=True) with open(f"./{url.split('/')[-1]}", 'wb') as f: for chunk in response.iter_content(chunk_size=1024): if stop_event.is_set(): print(f"终止下载:{url}") return f.write(chunk) print(f"下载完成:{url}") except Exception as e: print(f"下载出错:{url},错误信息:{str(e)}") def main_download(): # 每次启动下载前重置终止信号 stop_event.clear() futures = [executor.submit(download, url) for url in urls] for future in as_completed(futures): # 处理任务结果(如果需要) try: future.result() except Exception as e: print(f"任务执行出错:{e}") scheduler = BlockingScheduler(timezone="Asia/Kolkata") job = scheduler.add_job(main_download, trigger="cron", hour=12) def kill_all(): stop_event.set() executor.shutdown(wait=True) scheduler.remove_all_jobs() scheduler.shutdown() scheduler.add_job(kill_all, trigger="cron", hour=6) scheduler.start()
为什么之前的方法没用?
你调用的quit()、sys.exit()都是在APScheduler的工作线程里执行的,这些方法只会终止当前线程,而ThreadPoolExecutor的下载线程是独立的,不会感知到这个终止信号,所以会继续运行。用Event的方式是让每个下载线程主动检查并退出,这是Python中终止线程的标准安全做法。
内容的提问来源于stack exchange,提问作者The Unknown Bengali
相关产品推荐
相关产品推荐

