FastAPI中基于多进程Worker的资源锁实现方案咨询
问题:FastAPI模型更新时的并发安全与锁定实现
我要搭建一个带/get端点的FastAPI服务,返回ML模型的推理结果。基础功能已实现,但需要定期用新版本更新模型(模型获取逻辑暂不考虑)。担心旧模型替换为新模型的过程中,有请求调用旧模型会出现异常,想知道如何用asyncio实现锁定机制,确保多请求并发时,模型更新仅让请求等待而不失败(模型体积较小)。
原实现代码
import asyncio import time from concurrent.futures import ProcessPoolExecutor from fastapi import FastAPI, Request from sentence_transformers import SentenceTransformer app = FastAPI() sbertmodel = None def create_model(): global sbertmodel sbertmodel = SentenceTransformer('multi-qa-MiniLM-L6-cos-v1') # if you try to run all predicts concurrently, it will result in CPU trashing. pool = ProcessPoolExecutor(max_workers=1, initializer=create_model) def model_predict(): ts = time.time() vector = sbertmodel.encode('How big is London') return vector async def vector_search(vector): # simulate I/O call (e.g. Vector Similarity Search using a VectorDB) await asyncio.sleep(0.005) @app.get("/") async def entrypoint(request: Request): loop = asyncio.get_event_loop() ts = time.time() # worker should be initialized outside endpoint to avoid cold start vector = await loop.run_in_executor(pool, model_predict) print(f"Model : {int((time.time() - ts) * 1000)}ms") ts = time.time() await vector_search(vector) print(f"io task: {int((time.time() - ts) * 1000)}ms") return "ok"
解决方案
核心思路
利用asyncio.Lock实现全局互斥锁,确保模型更新过程中所有请求处于等待状态,直到更新完成后再继续处理,避免竞态条件导致的异常。同时调整模型管理方式,用原子操作替换模型实例,保证请求始终访问有效模型。
修改后的代码
import asyncio import time from concurrent.futures import ProcessPoolExecutor from fastapi import FastAPI, Request from sentence_transformers import SentenceTransformer app = FastAPI() # 全局模型实例 sbert_model = None # 模型更新互斥锁 update_lock = asyncio.Lock() # 推理专用进程池(避免CPU抢占) pool = ProcessPoolExecutor(max_workers=1) def init_model(): """初始化/加载模型的同步函数,供进程池调用""" return SentenceTransformer('multi-qa-MiniLM-L6-cos-v1') def model_predict(model): """模型推理函数,接收当前模型实例作为参数""" ts = time.time() vector = model.encode('How big is London') print(f"Model : {int((time.time() - ts) * 1000)}ms") return vector async def vector_search(vector): """模拟向量搜索的IO密集型任务""" await asyncio.sleep(0.005) print(f"io task: 5ms") @app.on_event("startup") async def startup_event(): """服务启动时初始化模型""" global sbert_model loop = asyncio.get_event_loop() sbert_model = await loop.run_in_executor(pool, init_model) @app.get("/") async def entrypoint(request: Request): global sbert_model # 获取锁:模型更新时请求会等待,直到锁释放 async with update_lock: loop = asyncio.get_event_loop() vector = await loop.run_in_executor(pool, model_predict, sbert_model) # IO任务无需持有锁,不影响并发性能 await vector_search(vector) return "ok" async def update_model(): """定期执行的模型更新函数""" global sbert_model async with update_lock: # 这里可替换为从外部获取新模型的逻辑 loop = asyncio.get_event_loop() new_model = await loop.run_in_executor(pool, init_model) # 原子替换模型实例,确保所有后续请求立即使用新模型 sbert_model = new_model print("模型已完成更新")
关键说明
- 互斥锁的使用:
async with update_lock会自动处理锁的获取与释放,模型更新时,所有进入/端点的请求都会等待锁释放,避免模型切换期间访问无效实例。 - 模型实例传递:将模型实例作为参数传入推理函数,替代原全局变量方式,避免进程池内全局变量的状态不一致问题。
- 原子替换:
sbert_model = new_model是原子操作,一旦执行完成,后续请求会立即使用新模型,无中间状态异常。 - 并发性能优化:IO密集型的向量搜索任务不持有锁,确保非模型操作的请求可以正常并发,不影响整体服务性能。
内容的提问来源于stack exchange,提问作者mehekek
相关产品推荐
相关产品推荐

