在线服务架构概览
vLLM 的在线服务栈从外到内分为四层:
api_server.py"] FastAPI --> SC["OpenAIServingChat
serving_chat.py"] FastAPI --> SComp["OpenAIServingCompletion
serving_completion.py"] SC --> AE["AsyncLLMEngine
async_llm_engine.py"] SComp --> AE AE --> RT["RequestTracker
管理请求生命周期"] AE --> BG["Background Loop
run_engine_loop()"] BG --> Engine["_AsyncLLMEngine
继承 LLMEngine"] Engine --> Scheduler["Scheduler"] Engine --> Executor["Executor"]
- FastAPI Server:HTTP 入口,定义
/v1/chat/completions、/v1/completions等路由 - OpenAIServing*:将 OpenAI 协议请求转换为 vLLM 内部格式,处理 tokenization、chat template 等
- AsyncLLMEngine:异步包装层,通过后台循环驱动
LLMEngine,提供generate()异步生成器 - _AsyncLLMEngine:
LLMEngine的异步扩展,核心方法step_async()使用await调用 Executor
LLMEngine 本身是同步的——它的 step() 方法会阻塞直到 GPU 执行完成。在在线服务场景中,需要同时处理多个并发的 HTTP 请求,不能用同步阻塞模式。AsyncLLMEngine 通过后台循环 + asyncio.Queue 的设计,将同步的引擎调用隐藏在一个 task 中,对外暴露 async 接口。
AsyncStream 流式输出
AsyncStream(L49-78)是连接引擎输出和 HTTP 响应的桥梁。每个请求对应一个 AsyncStream 实例:
class AsyncStream:
"""A stream of RequestOutputs for a request that can be
iterated over asynchronously."""
def __init__(self, request_id: str) -> None:
self.request_id = request_id
self._queue: asyncio.Queue = asyncio.Queue()
self._finished = False
def put(self, item: Union[RequestOutput, Exception]) -> None:
if self._finished:
return
self._queue.put_nowait(item)
def finish(self) -> None:
self._queue.put_nowait(StopAsyncIteration())
self._finished = True
@property
def finished(self) -> bool:
return self._finished
def __aiter__(self):
return self
async def __anext__(self) -> RequestOutput:
result = await self._queue.get()
if isinstance(result, Exception):
raise result
return result
设计要点:
- 异步迭代器:实现了
__aiter__和__anext__,可以用async for遍历 - Queue 驱动:后台循环通过
put()推送RequestOutput,前端通过__anext__消费 - 结束信号:
finish()放入StopAsyncIteration,自动终止async for - 异常传播:可以
put(Exception)将引擎异常传递给等待方
AsyncStream 本质是一个单生产者(后台引擎循环)-单消费者(HTTP handler)的异步队列。使用 put_nowait 确保生产者不会阻塞,而消费者通过 await queue.get() 自然地让出事件循环。
RequestTracker 请求管理
RequestTracker(L81-193)是 AsyncLLMEngine 的请求生命周期管理器,在异步世界和同步引擎之间架设桥梁:
class RequestTracker:
"""Synchronous abstraction for tracking requests."""
def __init__(self) -> None:
self._request_streams: Dict[str, AsyncStream] = {}
self._finished_requests: asyncio.Queue[str] = asyncio.Queue()
self._new_requests: asyncio.Queue[Tuple[AsyncStream,
dict]] = asyncio.Queue()
self.new_requests_event = asyncio.Event()
RequestTracker 管理三个核心数据结构:
| 属性 | 类型 | 作用 |
|---|---|---|
_request_streams | Dict[str, AsyncStream] | 活跃请求的 stream 映射 |
_new_requests | Queue | 待添加到引擎的新请求队列 |
_finished_requests | Queue | 已完成/已取消的请求 ID 队列 |
关键方法
def add_request(self, request_id: str,
**engine_add_request_kwargs) -> AsyncStream:
"""Add a request to be sent to the engine on the next background
loop iteration."""
if request_id in self._request_streams:
raise KeyError(f"Request {request_id} already exists.")
stream = AsyncStream(request_id)
self._new_requests.put_nowait((stream, {
"request_id": request_id,
**engine_add_request_kwargs
}))
self.new_requests_event.set()
return stream
def process_request_output(self, request_output: RequestOutput,
*, verbose: bool = False) -> None:
"""Process a request output from the engine."""
request_id = request_output.request_id
self._request_streams[request_id].put(request_output)
if request_output.finished:
self.abort_request(request_id)
请求的完整生命周期:
add_request()→ 创建AsyncStream,放入_new_requests队列,设置new_requests_event- 后台循环调用
get_new_and_finished_requests()→ 取出新请求,注册到_request_streams - 引擎 step 产出
RequestOutput→process_request_output()推送到对应 stream - 请求完成(
finished=True)→abort_request()清理 stream
_AsyncLLMEngine 异步扩展
_AsyncLLMEngine(L196-281)继承 LLMEngine,提供异步版本的核心方法:
class _AsyncLLMEngine(LLMEngine):
"""Extension of LLMEngine to add async methods."""
async def step_async(self) -> List[RequestOutput]:
seq_group_metadata_list, scheduler_outputs = self.scheduler.schedule()
if not scheduler_outputs.is_empty():
execute_model_req = ExecuteModelRequest(
seq_group_metadata_list=seq_group_metadata_list,
blocks_to_swap_in=scheduler_outputs.blocks_to_swap_in,
blocks_to_swap_out=scheduler_outputs.blocks_to_swap_out,
blocks_to_copy=scheduler_outputs.blocks_to_copy,
num_lookahead_slots=scheduler_outputs.num_lookahead_slots,
running_queue_size=scheduler_outputs.running_queue_size,
)
output = await self.model_executor.execute_model_async(
execute_model_req)
else:
output = []
request_outputs = self._process_model_outputs(
output, scheduler_outputs.scheduled_seq_groups,
scheduler_outputs.ignored_seq_groups, seq_group_metadata_list)
self.do_log_stats(scheduler_outputs, output)
return request_outputs
step_async 与同步版 LLMEngine.step() 的唯一区别是:调用 await self.model_executor.execute_model_async() 而非同步的 execute_model()。这使得 GPU 执行期间事件循环可以处理其他协程(如新请求到达、HTTP 响应发送等)。
async def add_request_async(
self, request_id: str, prompt: Optional[str],
sampling_params: SamplingParams,
prompt_token_ids: Optional[List[int]] = None,
arrival_time: Optional[float] = None,
lora_request: Optional[LoRARequest] = None,
multi_modal_data: Optional[MultiModalData] = None,
) -> None:
if arrival_time is None:
arrival_time = time.time()
prompt_token_ids = await self.encode_request_async(
request_id=request_id, prompt=prompt,
prompt_token_ids=prompt_token_ids, lora_request=lora_request)
return self.add_request(request_id, prompt=prompt,
prompt_token_ids=prompt_token_ids,
sampling_params=sampling_params,
arrival_time=arrival_time,
lora_request=lora_request,
multi_modal_data=multi_modal_data)
add_request_async 先异步 tokenize prompt(encode_request_async),然后调用同步的 LLMEngine.add_request() 将请求加入 Scheduler 的 waiting 队列。
AsyncLLMEngine 核心
AsyncLLMEngine(L283-697)是对外的异步引擎接口。它不继承 LLMEngine,而是通过组合持有一个 _AsyncLLMEngine 实例:
class AsyncLLMEngine:
_engine_class: Type[_AsyncLLMEngine] = _AsyncLLMEngine
def __init__(self, worker_use_ray: bool, engine_use_ray: bool,
*args, log_requests: bool = True,
max_log_len: Optional[int] = None,
start_engine_loop: bool = True, **kwargs) -> None:
self.worker_use_ray = worker_use_ray
self.engine_use_ray = engine_use_ray
self.log_requests = log_requests
self.max_log_len = max_log_len
self.engine = self._init_engine(*args, **kwargs)
self.background_loop: Optional[asyncio.Future] = None
self._background_loop_unshielded: Optional[asyncio.Task] = None
self.start_engine_loop = start_engine_loop
self._errored_with: Optional[BaseException] = None
# Lazy initialized fields
self._request_tracker: RequestTracker
Engine 初始化策略
_init_engine(L425-442)根据配置选择不同的引擎部署方式:
- 本地模式(
engine_use_ray=False):直接实例化_AsyncLLMEngine,引擎和 API server 在同一进程 - Ray 模式(
engine_use_ray=True):将_AsyncLLMEngine包装为 Ray Actor,引擎运行在独立进程中
@classmethod
def from_engine_args(cls, engine_args: AsyncEngineArgs,
start_engine_loop: bool = True,
usage_context: UsageContext = ...) -> "AsyncLLMEngine":
engine_config = engine_args.create_engine_config()
if engine_config.device_config.device_type == "neuron":
executor_class = NeuronExecutorAsync
elif engine_config.device_config.device_type == "cpu":
executor_class = CPUExecutorAsync
elif engine_config.parallel_config.worker_use_ray:
initialize_ray_cluster(engine_config.parallel_config)
executor_class = RayGPUExecutorAsync
else:
assert engine_config.parallel_config.world_size == 1
executor_class = GPUExecutorAsync
return cls(worker_use_ray, engine_use_ray,
**engine_config.to_dict(),
executor_class=executor_class, ...)
工厂方法根据设备类型和并行配置,自动选择正确的异步 Executor 实现。
run_engine_loop 主循环
run_engine_loop(L490-508)是整个在线服务的心脏——一个永不停止的后台协程,持续驱动引擎执行推理:
async def run_engine_loop(self):
has_requests_in_progress = False
while True:
if not has_requests_in_progress:
logger.debug("Waiting for new requests...")
await self._request_tracker.wait_for_new_requests()
logger.debug("Got new requests!")
# Abort if iteration takes too long due to unrecoverable errors
try:
has_requests_in_progress = await asyncio.wait_for(
self.engine_step(), ENGINE_ITERATION_TIMEOUT_S)
except asyncio.TimeoutError as exc:
logger.error(
"Engine iteration timed out. This should never happen!")
self.set_errored(exc)
raise
await asyncio.sleep(0)
阻塞等待新请求"] B -->|"是"| D["engine_step()
执行一轮推理"] C --> D D --> E["await asyncio.sleep(0)
让出事件循环"] E --> F{"engine_step 返回值"} F -->|"True: 还有请求"| B F -->|"False: 无请求"| B
engine_step 详解
async def engine_step(self) -> bool:
"""Kick the engine to process the waiting requests.
Returns True if there are in-progress requests."""
new_requests, finished_requests = (
self._request_tracker.get_new_and_finished_requests())
for new_request in new_requests:
try:
if self.engine_use_ray:
await self.engine.add_request.remote(**new_request)
else:
await self.engine.add_request_async(**new_request)
except ValueError as e:
self._request_tracker.process_exception(
new_request["request_id"], e, verbose=self.log_requests)
if finished_requests:
await self._engine_abort(finished_requests)
if self.engine_use_ray:
request_outputs = await self.engine.step.remote()
else:
request_outputs = await self.engine.step_async()
# Put the outputs into the corresponding streams.
for request_output in request_outputs:
self._request_tracker.process_request_output(
request_output, verbose=self.log_requests)
return len(request_outputs) > 0
每一轮 engine_step 执行三件事:
- 同步请求状态:从
RequestTracker取出新请求和已完成请求 - 驱动引擎:调用
step_async()(调度 + 模型前向 + 采样) - 分发输出:将
RequestOutput推送到各请求的AsyncStream
asyncio.sleep(0) 是一个显式的协程切换点。它确保事件循环在每轮引擎 step 之间有机会处理其他任务(如接收新的 HTTP 请求、发送 streaming 响应)。没有它,引擎循环可能会"饿死"其他协程。
generate 方法
generate()(L574-666)是 AsyncLLMEngine 对外最重要的接口——一个异步生成器,将请求提交给引擎并逐步产出结果:
async def generate(
self, prompt: Optional[str],
sampling_params: SamplingParams,
request_id: str,
prompt_token_ids: Optional[List[int]] = None,
lora_request: Optional[LoRARequest] = None,
multi_modal_data: Optional[MultiModalData] = None
) -> AsyncIterator[RequestOutput]:
"""Generate outputs for a request.
...
Details:
- If the engine is not running, start the background loop,
which iteratively invokes engine_step to process requests.
- Add the request to the engine's RequestTracker.
- Wait for the request outputs from AsyncStream and yield them.
"""
arrival_time = time.time()
try:
stream = await self.add_request(
request_id, prompt, sampling_params,
prompt_token_ids=prompt_token_ids,
arrival_time=arrival_time,
lora_request=lora_request,
multi_modal_data=multi_modal_data,
)
async for request_output in stream:
yield request_output
except (Exception, asyncio.CancelledError) as e:
self._abort(request_id)
raise e
关键设计:
- 惰性启动:
add_request()中检查is_running,如果后台循环未启动则调用start_background_loop() - 自动中止:如果
generate()被取消(如客户端断开),catch 异常并调用_abort()通知引擎释放资源 - 背压控制:
AsyncStream使用无界队列,引擎输出速度受 step 频率限制(通常每秒数十次),不会造成内存问题
OpenAI API Server
api_server.py 使用 FastAPI 构建 OpenAI 兼容的 HTTP API。核心路由:
@app.get("/health")
async def health() -> Response:
await openai_serving_chat.engine.check_health()
return Response(status_code=200)
@app.get("/v1/models")
async def show_available_models():
models = await openai_serving_chat.show_available_models()
return JSONResponse(content=models.model_dump())
@app.post("/v1/chat/completions")
async def create_chat_completion(request: ChatCompletionRequest,
raw_request: Request):
generator = await openai_serving_chat.create_chat_completion(
request, raw_request)
if isinstance(generator, ErrorResponse):
return JSONResponse(content=generator.model_dump(),
status_code=generator.code)
if request.stream:
return StreamingResponse(content=generator,
media_type="text/event-stream")
else:
return JSONResponse(content=generator.model_dump())
@app.post("/v1/completions")
async def create_completion(request: CompletionRequest, raw_request: Request):
generator = await openai_serving_completion.create_completion(
request, raw_request)
if request.stream:
return StreamingResponse(content=generator,
media_type="text/event-stream")
else:
return JSONResponse(content=generator.model_dump())
Server 启动流程
if __name__ == "__main__":
args = parse_args()
# CORS, authentication middleware ...
engine_args = AsyncEngineArgs.from_cli_args(args)
engine = AsyncLLMEngine.from_engine_args(
engine_args, usage_context=UsageContext.OPENAI_API_SERVER)
openai_serving_chat = OpenAIServingChat(
engine, served_model_names, args.response_role,
args.lora_modules, args.chat_template)
openai_serving_completion = OpenAIServingCompletion(
engine, served_model_names, args.lora_modules)
uvicorn.run(app, host=args.host, port=args.port,
timeout_keep_alive=TIMEOUT_KEEP_ALIVE, ...)
启动顺序:
- 解析命令行参数 → 创建
AsyncEngineArgs - 创建
AsyncLLMEngine(此时模型加载、KV Cache 分配完成) - 创建
OpenAIServingChat和OpenAIServingCompletion服务层 uvicorn.run()启动 HTTP 服务,进入事件循环
--api-key 或 VLLM_API_KEY 环境变量设置 Bearer token 认证。认证中间件只拦截 /v1/* 路径的请求,健康检查等路径不受影响。
Chat Completion 实现
OpenAIServingChat(serving_chat.py)继承 OpenAIServing 基类,处理 /v1/chat/completions 请求。
async def create_chat_completion(
self, request: ChatCompletionRequest, raw_request: Request
) -> Union[ErrorResponse, AsyncGenerator[str, None],
ChatCompletionResponse]:
# 1. 模型检查
error_check_ret = await self._check_model(request)
if error_check_ret is not None:
return error_check_ret
# 2. 解析消息,应用 chat template
conversation: List[ConversationMessage] = []
for m in request.messages:
messages, _ = self._parse_chat_message_content(
m["role"], m["content"])
conversation.extend(messages)
prompt = self.tokenizer.apply_chat_template(
conversation=conversation, tokenize=False,
add_generation_prompt=request.add_generation_prompt)
# 3. Tokenize 并构造采样参数
request_id = f"cmpl-{random_uuid()}"
prompt_ids, prompt_text = self._validate_prompt_and_tokenize(
request, prompt=prompt)
sampling_params = request.to_sampling_params()
# 4. 调用 engine.generate()
result_generator = self.engine.generate(
prompt_text, sampling_params, request_id, prompt_ids,
lora_request)
# 5. 分流: streaming vs non-streaming
if request.stream:
return self.chat_completion_stream_generator(
request, result_generator, request_id, conversation)
else:
return await self.chat_completion_full_generator(
request, raw_request, result_generator, request_id,
conversation)
处理流程:
- Chat Template:使用 tokenizer 的
apply_chat_template将消息列表转换为模型能理解的 prompt 字符串 - Tokenization:将 prompt 字符串转为 token IDs,同时检查长度限制
- Guided Decoding:如果请求了 JSON mode 或正则约束,构造
guided_decode_logits_processor - 提交引擎:调用
engine.generate()获得异步迭代器 - 响应分流:根据
request.stream选择流式或完整响应
Streaming 与 Non-Streaming
Streaming 响应
chat_completion_stream_generator(L151-289)是一个异步生成器,逐 token 将结果以 SSE 格式发送给客户端:
async def chat_completion_stream_generator(
self, request, result_generator, request_id, conversation
) -> AsyncGenerator[str, None]:
model_name = self.served_model_names[0]
created_time = int(time.time())
previous_texts = [""] * request.n
previous_num_tokens = [0] * request.n
async for res in result_generator:
if first_iteration:
# 发送 role delta
for i in range(request.n):
chunk = ChatCompletionStreamResponse(
choices=[ChatCompletionResponseStreamChoice(
index=i, delta=DeltaMessage(role=role),
finish_reason=None)])
yield f"data: {chunk.model_dump_json(...)}\n\n"
first_iteration = False
for output in res.outputs:
delta_text = output.text[len(previous_texts[i]):]
# 发送 content delta
if output.finish_reason is None:
chunk = ... # delta with content
yield f"data: {chunk.model_dump_json(...)}\n\n"
else:
chunk = ... # final chunk with finish_reason + usage
yield f"data: {chunk.model_dump_json(...)}\n\n"
yield "data: [DONE]\n\n"
Streaming 响应遵循 SSE(Server-Sent Events)协议:
- 每个 chunk 以
data: {JSON}\n\n格式发送 - 第一个 chunk 发送 role 信息(如
"assistant") - 后续 chunk 发送增量文本(
delta_text = 新文本 - 已发送文本) - 最后一个 chunk 包含
finish_reason和usage信息 - 以
data: [DONE]\n\n结束流
Non-Streaming 响应
chat_completion_full_generator(L291-350+)等待所有 token 生成完毕后,一次性返回完整结果:
async def chat_completion_full_generator(
self, request, raw_request, result_generator, request_id,
conversation
) -> Union[ErrorResponse, ChatCompletionResponse]:
final_res: Optional[RequestOutput] = None
async for res in result_generator:
if await raw_request.is_disconnected():
# Abort the request if the client disconnects.
await self.engine.abort(request_id)
return self.create_error_response("Client disconnected")
final_res = res
assert final_res is not None
# 构造完整的 ChatCompletionResponse ...
return ChatCompletionResponse(choices=choices, usage=usage, ...)
Non-streaming 模式的关键特性:
- 消费所有
RequestOutput,只保留最后一个(final_res,包含完整生成文本) - 每次迭代检查客户端连接状态(
raw_request.is_disconnected()),提前中止已断开的请求 - 一次性返回
ChatCompletionResponse,包含所有 choices 和 usage 统计
Completion API
OpenAIServingCompletion(serving_completion.py)处理 /v1/completions 端点,逻辑类似但有以下差异:
- 支持多种 prompt 格式:字符串、字符串数组、token 数组、token 数组的数组
- 多 prompt 情况下,为每个 prompt 创建独立的
generate()调用,使用merge_async_iterators合并输出 - 当
n != best_of或使用 beam search 时,自动禁用 streaming