01

在线服务架构概览

vLLM 的在线服务栈从外到内分为四层:

graph TD Client["HTTP Client"] --> FastAPI["FastAPI Server
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() 异步生成器
  • _AsyncLLMEngineLLMEngine 的异步扩展,核心方法 step_async() 使用 await 调用 Executor
为什么需要 AsyncLLMEngine
LLMEngine 本身是同步的——它的 step() 方法会阻塞直到 GPU 执行完成。在在线服务场景中,需要同时处理多个并发的 HTTP 请求,不能用同步阻塞模式。AsyncLLMEngine 通过后台循环 + asyncio.Queue 的设计,将同步的引擎调用隐藏在一个 task 中,对外暴露 async 接口。
02

AsyncStream 流式输出

AsyncStream(L49-78)是连接引擎输出和 HTTP 响应的桥梁。每个请求对应一个 AsyncStream 实例:

async_llm_engine.py — AsyncStream L49-78
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() 自然地让出事件循环。
03

RequestTracker 请求管理

RequestTracker(L81-193)是 AsyncLLMEngine 的请求生命周期管理器,在异步世界和同步引擎之间架设桥梁:

async_llm_engine.py — RequestTracker 核心结构 L81-93
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_streamsDict[str, AsyncStream]活跃请求的 stream 映射
_new_requestsQueue待添加到引擎的新请求队列
_finished_requestsQueue已完成/已取消的请求 ID 队列

关键方法

async_llm_engine.py — add_request & process_request_output L134-121
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)

请求的完整生命周期:

  1. add_request() → 创建 AsyncStream,放入 _new_requests 队列,设置 new_requests_event
  2. 后台循环调用 get_new_and_finished_requests() → 取出新请求,注册到 _request_streams
  3. 引擎 step 产出 RequestOutputprocess_request_output() 推送到对应 stream
  4. 请求完成(finished=True)→ abort_request() 清理 stream
04

_AsyncLLMEngine 异步扩展

_AsyncLLMEngine(L196-281)继承 LLMEngine,提供异步版本的核心方法:

async_llm_engine.py — step_async L199-233
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_llm_engine.py — add_request_async L250-277
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 队列。

05

AsyncLLMEngine 核心

AsyncLLMEngine(L283-697)是对外的异步引擎接口。它不继承 LLMEngine,而是通过组合持有一个 _AsyncLLMEngine 实例:

async_llm_engine.py — AsyncLLMEngine.__init__ L312-335
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,引擎运行在独立进程中
async_llm_engine.py — from_engine_args 工厂方法 L338-377
@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 实现。

06

run_engine_loop 主循环

run_engine_loop(L490-508)是整个在线服务的心脏——一个永不停止的后台协程,持续驱动引擎执行推理:

async_llm_engine.py — 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)
graph TD A["run_engine_loop 启动"] --> B{"有进行中的请求?"} B -->|"否"| C["await wait_for_new_requests()
阻塞等待新请求"] 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_llm_engine.py — engine_step L444-482
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 执行三件事:

  1. 同步请求状态:从 RequestTracker 取出新请求和已完成请求
  2. 驱动引擎:调用 step_async()(调度 + 模型前向 + 采样)
  3. 分发输出:将 RequestOutput 推送到各请求的 AsyncStream
await asyncio.sleep(0) 的作用
asyncio.sleep(0) 是一个显式的协程切换点。它确保事件循环在每轮引擎 step 之间有机会处理其他任务(如接收新的 HTTP 请求、发送 streaming 响应)。没有它,引擎循环可能会"饿死"其他协程。
07

generate 方法

generate()(L574-666)是 AsyncLLMEngine 对外最重要的接口——一个异步生成器,将请求提交给引擎并逐步产出结果:

async_llm_engine.py — generate L574-666
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
sequenceDiagram participant Client as HTTP Handler participant Gen as generate() participant AT as add_request() participant RT as RequestTracker participant BG as Background Loop participant Engine as _AsyncLLMEngine Client->>Gen: async for output in generate(...) Gen->>AT: await add_request() AT->>RT: add_request() → AsyncStream AT-->>Gen: return stream loop 每轮 engine step BG->>RT: get_new_and_finished_requests() BG->>Engine: await step_async() Engine-->>BG: [RequestOutput, ...] BG->>RT: process_request_output() RT->>Gen: stream.put(output) Gen-->>Client: yield output end

关键设计:

  • 惰性启动add_request() 中检查 is_running,如果后台循环未启动则调用 start_background_loop()
  • 自动中止:如果 generate() 被取消(如客户端断开),catch 异常并调用 _abort() 通知引擎释放资源
  • 背压控制AsyncStream 使用无界队列,引擎输出速度受 step 频率限制(通常每秒数十次),不会造成内存问题
08

OpenAI API Server

api_server.py 使用 FastAPI 构建 OpenAI 兼容的 HTTP API。核心路由:

api_server.py — 路由定义 L77-123
@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 启动流程

api_server.py — 启动入口 L126-186
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, ...)

启动顺序:

  1. 解析命令行参数 → 创建 AsyncEngineArgs
  2. 创建 AsyncLLMEngine(此时模型加载、KV Cache 分配完成)
  3. 创建 OpenAIServingChatOpenAIServingCompletion 服务层
  4. uvicorn.run() 启动 HTTP 服务,进入事件循环
API 认证
vLLM 支持通过 --api-keyVLLM_API_KEY 环境变量设置 Bearer token 认证。认证中间件只拦截 /v1/* 路径的请求,健康检查等路径不受影响。
09

Chat Completion 实现

OpenAIServingChatserving_chat.py)继承 OpenAIServing 基类,处理 /v1/chat/completions 请求。

serving_chat.py — create_chat_completion L72-143
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)

处理流程:

  1. Chat Template:使用 tokenizer 的 apply_chat_template 将消息列表转换为模型能理解的 prompt 字符串
  2. Tokenization:将 prompt 字符串转为 token IDs,同时检查长度限制
  3. Guided Decoding:如果请求了 JSON mode 或正则约束,构造 guided_decode_logits_processor
  4. 提交引擎:调用 engine.generate() 获得异步迭代器
  5. 响应分流:根据 request.stream 选择流式或完整响应
10

Streaming 与 Non-Streaming

Streaming 响应

chat_completion_stream_generator(L151-289)是一个异步生成器,逐 token 将结果以 SSE 格式发送给客户端:

serving_chat.py — streaming 核心逻辑 L151-289
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_reasonusage 信息
  • data: [DONE]\n\n 结束流

Non-Streaming 响应

chat_completion_full_generator(L291-350+)等待所有 token 生成完毕后,一次性返回完整结果:

serving_chat.py — non-streaming 核心逻辑 L291-307
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 统计
Streaming 与 Non-Streaming 的引擎行为
无论客户端选择 streaming 还是 non-streaming,引擎内部的行为完全相同——都是逐 token 生成。区别仅在 HTTP 层:streaming 模式逐步发送 SSE chunk,non-streaming 模式在内存中累积到完成后一次发送。这意味着 non-streaming 请求的首字节延迟 = 整个生成时间,但总吞吐量与 streaming 相同。

Completion API

OpenAIServingCompletionserving_completion.py)处理 /v1/completions 端点,逻辑类似但有以下差异:

  • 支持多种 prompt 格式:字符串、字符串数组、token 数组、token 数组的数组
  • 多 prompt 情况下,为每个 prompt 创建独立的 generate() 调用,使用 merge_async_iterators 合并输出
  • n != best_of 或使用 beam search 时,自动禁用 streaming