如何以非阻塞方式使用starmap_async?FastAPI场景解惑
解决方案:非阻塞使用
starmap_async并避免FastAPI主线程阻塞 核心问题分析
你的代码阻塞主线程的原因有两个:
starmap_async的wait()/get()是同步阻塞调用,在FastAPI的异步事件循环中执行会直接卡住整个线程,导致无法处理HTTP请求。- 启动事件中直接
await build_cache()会强制服务器等待缓存构建完成后才启动服务。
解决思路
将阻塞的多进程操作从异步事件循环中剥离,用asyncio.run_in_executor调度到独立线程/进程池执行;同时在启动事件中用后台任务异步执行缓存构建,不阻塞服务器启动。
代码实现方案一:基于multiprocessing.Pool改造
import asyncio from multiprocessing.pool import Pool from multiprocessing import cpu_count from io import BytesIO from ftplib import FTP from fastapi import FastAPI import uvicorn ftp_url = "你的FTP地址" def get_from_ftp(file_name): flo = BytesIO() ftp_server = FTP(ftp_url) ftp_server.login() ftp_server.retrbinary(f'RETR {file_name}', flo.write) flo.seek(0) # 写入磁盘逻辑 with open(file_name, 'wb') as f: f.write(flo.read()) # 封装多进程下载逻辑为同步函数 def run_ftp_downloads(): with Pool(cpu_count() - 1) as pool: some_files_to_download = ['file1', 'file2'] # 单参数场景用map_async更合适,多参数才需要starmap_async(传入元组列表) result = pool.map_async(get_from_ftp, some_files_to_download) result.wait() async def build_cache(): loop = asyncio.get_running_loop() # 将阻塞的多进程操作放到executor中,不占用事件循环线程 await loop.run_in_executor(None, run_ftp_downloads) app = FastAPI() @app.on_event("startup") async def startup(): # 用后台任务启动缓存构建,不等待完成,服务器直接启动 asyncio.create_task(build_cache()) @app.get("/") async def root(): return {"message": "服务器已启动,缓存构建中..."} if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000)
代码实现方案二:改用concurrent.futures.ProcessPoolExecutor(更贴合asyncio生态)
如果不需要multiprocessing.Pool的特殊API,推荐直接用concurrent.futures的进程池,和asyncio兼容性更好:
import asyncio from concurrent.futures import ProcessPoolExecutor from multiprocessing import cpu_count from io import BytesIO from ftplib import FTP from fastapi import FastAPI import uvicorn ftp_url = "你的FTP地址" def get_from_ftp(file_name): flo = BytesIO() ftp_server = FTP(ftp_url) ftp_server.login() ftp_server.retrbinary(f'RETR {file_name}', flo.write) flo.seek(0) with open(file_name, 'wb') as f: f.write(flo.read()) async def build_cache(): loop = asyncio.get_running_loop() some_files_to_download = ['file1', 'file2'] with ProcessPoolExecutor(max_workers=cpu_count()-1) as executor: # 将进程池的阻塞map操作放到executor中执行 await loop.run_in_executor(None, lambda: list(executor.map(get_from_ftp, some_files_to_download))) app = FastAPI() @app.on_event("startup") async def startup(): asyncio.create_task(build_cache()) @app.get("/") async def root(): return {"message": "服务器已启动,缓存构建中..."} if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000)
关键说明
- 关于
starmap_async的正确用法:如果get_from_ftp需要多个参数(比如get_from_ftp(file_name, ftp_user, ftp_pwd)),则需要将参数整理为元组列表,用starmap_async,例如:params = [('file1', 'user', 'pwd'), ('file2', 'user', 'pwd')] result = pool.starmap_async(get_from_ftp, params) - Windows系统注意事项:多进程代码必须放在
if __name__ == "__main__"保护之外,否则会触发进程启动异常。
内容的提问来源于stack exchange,提问作者jossefaz
相关产品推荐
相关产品推荐

