FastAPI中I/O密集型任务的线程实现方案选择咨询
FastAPI中I/O密集型任务的线程实现方案选择咨询
嘿,刚好之前我也在FastAPI里处理过类似的本地文件I/O密集场景,来给你唠唠我的经验~
首先得明确:FastAPI本身是基于asyncio异步框架的,所以处理I/O密集型任务时,优先推荐用asyncio.get_running_loop()结合run_in_executor的方式,比手动用ThreadPoolExecutor省心太多,完全符合你说的“减少样板代码”的需求。
为啥优先选run_in_executor?
FastAPI在异步模式下,默认会维护一个全局的线程池,当你调用loop.run_in_executor(None, 同步函数, 参数)时,它会自动把同步的I/O任务(比如你的本地JSON文件读取)丢到这个线程池里执行,不会阻塞主事件循环。而且你不用手动创建、关闭线程池,也不用操心线程池的资源回收,框架都帮你搞定了。
给你看我当时用的代码示例:
from fastapi import FastAPI import asyncio import json from pathlib import Path app = FastAPI() # 这是你的同步读取文件函数,纯I/O操作 def read_local_json(file_path: str) -> dict: with open(file_path, 'r', encoding='utf-8') as f: return json.load(f) @app.get("/batch-load-json/") async def batch_load_json(): # 模拟你的500-1000个文件路径列表 file_paths = [Path(f"my_data_{i}.json") for i in range(800)] loop = asyncio.get_running_loop() # 把所有读取任务提交到线程池,并发执行 tasks = [loop.run_in_executor(None, read_local_json, str(path)) for path in file_paths] # 等待所有任务完成,收集结果 all_results = await asyncio.gather(*tasks) return {"total_files": len(all_results), "first_sample": all_results[0] if all_results else None}
那ThreadPoolExecutor什么时候用?
如果你有特殊需求,比如要自定义线程池的大小(比如磁盘I/O带宽有限,想控制并发数)、给线程设置自定义名称方便调试,这时候再考虑手动创建ThreadPoolExecutor。
但要注意:别在每个请求里创建新的线程池!那样会导致资源浪费甚至内存泄漏,最好全局初始化一个,然后在应用关闭时正确销毁。比如:
from fastapi import FastAPI import asyncio import json from pathlib import Path from concurrent.futures import ThreadPoolExecutor app = FastAPI() # 全局线程池,根据你的磁盘性能调整,比如16-32个线程就够了 custom_executor = ThreadPoolExecutor(max_workers=24, thread_name_prefix="FileReaderThread") def read_local_json(file_path: str) -> dict: with open(file_path, 'r', encoding='utf-8') as f: return json.load(f) # 应用关闭时关闭线程池 @app.on_event("shutdown") def shutdown_executor(): custom_executor.shutdown(wait=True) @app.get("/batch-load-json-custom/") async def batch_load_json_custom(): file_paths = [Path(f"my_data_{i}.json") for i in range(800)] loop = asyncio.get_running_loop() # 用自定义线程池执行任务 tasks = [loop.run_in_executor(custom_executor, read_local_json, str(path)) for path in file_paths] all_results = await asyncio.gather(*tasks) return {"total_files": len(all_results), "first_sample": all_results[0] if all_results else None}
小提醒
本地文件I/O的并发数不是越大越好哦!磁盘的读写带宽是有限的,开太多线程反而会因为磁盘资源竞争导致速度变慢。我当时测试的时候,16-24个线程的速度是最快的,你可以根据自己的磁盘实际性能调整。
内容来源于stack exchange
相关产品推荐
相关产品推荐

