FastAPI场景下如何去除轮询共享字典的while循环,优雅监听线程任务状态
优化方案:用asyncio.Future替换轮询逻辑
你当前的轮询逻辑本质是主动查询结果,改成后台计算完成后主动通知的模式即可完全消除while循环,具体实现依赖asyncio提供的Future对象,是异步场景下这类需求的标准实现方案。
核心改动点
- 新增存储Future的共享字典,每个请求生成task_id的同时创建一个Future对象存入字典,请求侧直接await这个Future即可
- 后台计算线程完成批量推理后,通过主线程事件循环的
call_soon_threadsafe方法,给对应task_id的Future设置结果 - 原有存储计算结果的共享字典、轮询逻辑可以完全移除
完整修改后代码
import asyncio import uuid from typing import Union, List import threading from queue import Queue from fastapi import FastAPI, Request, Body, APIRouter import uvicorn import time import logging import datetime logger = logging.getLogger(__name__) app = APIRouter() # 新增:存储主线程事件循环,供后台线程跨线程操作Future使用 main_loop: asyncio.AbstractEventLoop = None def feed_data_into_model(queue, future_dict, lock): if queue.qsize() != 0: data = [] ids = [] while queue.qsize() != 0: task = queue.get() task_id = task[0] ids.append(task_id) text = task[1] data.append(text) result = model_work(data) # 计算完成后给每个Future设结果 for index,task_id in enumerate(ids): value = result[index] with lock: if task_id in future_dict: future = future_dict.pop(task_id) # 跨线程操作Future必须用call_soon_threadsafe保证线程安全 main_loop.call_soon_threadsafe(future.set_result, value) class TestThreading(object): def __init__(self, interval, queue,future_dict,lock): self.interval = interval thread = threading.Thread(target=self.run, args=(queue,future_dict,lock)) thread.daemon = True thread.start() def run(self,queue,future_dict,lock): while True: feed_data_into_model(queue,future_dict,lock) time.sleep(self.interval) if __name__ != "__main__": global queue, future_dict, lock queue = Queue(maxsize=64) future_dict = {} # 现在存储Future对象,不再直接存计算结果 lock = threading.Lock() tr = TestThreading(0.9, queue,future_dict,lock) # 服务启动时获取主线程的事件循环 main_loop = asyncio.get_running_loop() def model_work(x:Union[str,List[str]]): time.sleep(3) if isinstance(x,str): result = [len(x)] else: result = [len(_) for _ in x] return result # 原有get_response完全简化,无任何轮询逻辑 async def get_response(task_id, lock, future_dict, timeout: int = 10): future = main_loop.create_future() with lock: future_dict[task_id] = future # 增加超时逻辑防止请求无限挂起 try: return await asyncio.wait_for(future, timeout=timeout) except asyncio.TimeoutError: with lock: future_dict.pop(task_id, None) raise TimeoutError("任务执行超时") @app.get("/{text}") async def demo(text:str): global queue, future_dict, lock task_id = str(uuid.uuid4()) logger.info(task_id) item= [task_id,text] queue.put(item) # 直接await结果即可,无轮询 value = await get_response(task_id, lock, future_dict) return {"length": value} if __name__ == "__main__": uvicorn.run("api:app", host="0.0.0.0", port=5555)
客户端测试命令保持不变
for n in {1..5}; do curl http://localhost:5555/a & ; done
方案优势
- 完全消除轮询的CPU空转开销,也没有轮询间隔带来的结果返回延迟,性能和响应速度都有提升
- 代码逻辑更符合异步编程的最佳实践,可读性和可维护性更高
- 支持异常传递,如果后台计算出现错误,可以调用
future.set_exception()把异常直接传递到请求侧抛出,不需要额外的状态字段判断
内容的提问来源于stack exchange,提问作者swordHeart
相关产品推荐
相关产品推荐

