FastAPI中ProcessPoolExecutor与run_in_executor联用致API挂起问题
FastAPI进程池阻塞问题修复方案
问题根源
你遇到的API挂起、卡在with语句的问题,主要是这几个原因:
- 每次请求临时创建
ProcessPoolExecutor,进程启停的开销会直接阻塞异步事件循环,甚至触发死锁 - 进程池传递的函数/数据不可被pickle序列化(Python进程间通信依赖pickle),导致任务静默卡住
- 异步函数里误用同步的进程池上下文管理器,破坏了异步执行逻辑
分步修复
1. 全局初始化进程池
不要在请求处理函数里创建进程池,而是在FastAPI启动时初始化全局实例:
from fastapi import FastAPI import asyncio from concurrent.futures import ProcessPoolExecutor import re # 根据CPU核心数设置进程数,建议是核心数的1-2倍 global_executor = ProcessPoolExecutor(max_workers=4) app = FastAPI() # 程序关闭时优雅销毁进程池 @app.on_event("shutdown") def shutdown_executor(): global_executor.shutdown(wait=True)
2. 保证并行函数可序列化
你的split_on_whitespace和run_regex_on_content_chunk必须满足:
- 不能是嵌套函数(嵌套函数无法被pickle)
- 不要引用无法序列化的闭包变量或类实例
示例可序列化的函数:
def split_on_whitespace(content: str) -> list[str]: # 按空白符拆分内容,返回分块列表 return re.split(r'\s+', content.strip()) def run_regex_on_content_chunk(chunk: str) -> list[str]: # 替换成你的实际正则匹配逻辑 target_pattern = re.compile(r'[a-zA-Z0-9_-]+') return target_pattern.findall(chunk)
3. 正确在异步上下文调用进程池
直接用全局进程池,通过asyncio.run_in_executor提交任务,避免用同步的with语句:
async def process_content(content: str) -> list[str]: # 轻量的拆分逻辑直接在异步线程执行,不用放进程池 chunks = split_on_whitespace(content) # 获取当前事件循环,批量提交分块任务到进程池 loop = asyncio.get_running_loop() task_list = [ loop.run_in_executor(global_executor, run_regex_on_content_chunk, chunk) for chunk in chunks ] # 等待所有并行任务完成,合并结果 all_results = await asyncio.gather(*task_list) return [match for sub_result in all_results for match in sub_result] @app.post("/process-content") async def handle_process_request(content: str): match_results = await process_content(content) return {"total_matches": len(match_results), "matches": match_results}
额外优化点
- 对1MB级的大内容,拆分时控制分块大小(比如每块10KB),避免创建过多进程任务导致资源耗尽
- 可以监控进程池任务队列长度(
global_executor._work_queue.qsize()),防止任务积压 - 正则表达式提前编译(比如示例里的
re.compile),避免每次匹配重复编译浪费资源
内容的提问来源于stack exchange,提问作者ntriisii
相关产品推荐
相关产品推荐

