vLLM+FastAPI异步流响应瓶颈:多请求下FastAPI无法适配vLLM速度
解决FastAPI StreamingResponse并发瓶颈问题
问题核心
FastAPI异步后端处理多请求时,因vLLM的token流速度过快,事件循环被单个请求的流处理循环独占,无法切换到其他请求任务,最终导致请求串行处理、响应延迟。
解决方案
1. 强制让出事件循环控制权
在vLLM流处理循环中加入await asyncio.sleep(0),强制事件循环切换到其他就绪任务,避免当前请求独占资源。sleep(0)不会引入实际延迟,仅触发任务调度。
2. 优化流处理逻辑
用aiter_lines()替代aiter_bytes()处理vLLM响应,更贴合vLLM的SSE流格式(每行返回data: <payload>),减少手动解码和分割的开销,同时跳过无效行(如结尾的[DONE])。
3. 用后台任务队列解耦请求处理
启动独立的后台任务循环处理请求队列,每个vLLM调用作为单独的异步任务运行,主事件循环专注于处理前端的请求和响应流,避免阻塞。
4. 移除不必要的轮询延迟
原端点代码中await asyncio.sleep(0.1)的定时轮询会导致消息积压和延迟,直接使用await response_queue.get()等待队列消息,有消息立即处理。
修改后的代码示例
1. 优化vLLM调用函数
import asyncio import json import httpx from your_module import LlmRequestModel URL = "http://your-vllm-server-address" async def call_infer_llm(request: LlmRequestModel, response_queue): data = { "model": "/usr/Workplace/models/llama3-8b-Instruct/", "messages": request.messages, "temperature": request.temperature, "top_p": request.top_p, "stream": True, } async with httpx.AsyncClient() as client: # 直接传json参数,httpx自动序列化,无需手动json.dumps async with client.stream('POST', f"{URL}/v1/chat/completions", json=data) as resp: async for line in resp.aiter_lines(): if line.startswith('data: '): try: payload = json.loads(line[6:]) new_tokens = payload['choices'][0]['delta'].get('content', '') await response_queue.put({ 'status': 'PARTIAL_RESULT', 'result': new_tokens }) # 强制让出事件循环控制权 await asyncio.sleep(0) except json.JSONDecodeError: # 跳过无效JSON行(如[data: [DONE]]) continue # 发送任务完成信号 await response_queue.put({'status': 'DONE'})
2. 请求处理与后台队列
async def process_llm_request(data, response_queue): # 预处理:解析请求数据为LlmRequestModel c_request = LlmRequestModel(**json.loads(data)) # 启动异步任务处理vLLM调用,不阻塞当前队列处理 await call_infer_llm(c_request, response_queue)
3. 优化后的FastAPI端点
from fastapi import FastAPI, Form, StreamingResponse import asyncio import json app = FastAPI() request_queue = asyncio.Queue() # 启动后台任务循环处理请求队列 @app.on_event("startup") async def startup_event(): async def process_queue(): while True: response_queue, data = await request_queue.get() try: await process_llm_request(data, response_queue) except Exception as e: await response_queue.put({ 'status': 'ERROR', 'result': str(e) }) finally: request_queue.task_done() asyncio.create_task(process_queue()) @app.post("/api/llm-request") async def llm_endpoint(data: str = Form(...)): async def handle_llm_request(): response_queue = asyncio.Queue() request_queue.put_nowait([response_queue, data]) try: while True: message = await response_queue.get() yield json.dumps(message) if message['status'] in ['DONE', 'ERROR']: break except asyncio.CancelledError: print('task cancelled') return StreamingResponse(handle_llm_request(), media_type="application/x-ndjson")
其他优化建议
- 调整vLLM批量参数:通过API参数
max_tokens_per_batch让vLLM批量返回token,降低流的频率,减少后端处理次数(无需修改vLLM源码)。 - 配置httpx连接池:给
httpx.AsyncClient设置limits=httpx.Limits(max_connections=20),避免连接数限制影响并发。 - 前端直接调用vLLM(可选):如果预处理逻辑简单,可以让前端先调用后端完成预处理,再直接请求vLLM,彻底绕过后端的流转发瓶颈。
内容的提问来源于stack exchange,提问作者Jules Civel
相关产品推荐
相关产品推荐

