vLLM 中的流式请求与实时 API

19 分钟阅读
Meta、Mistral AI 以及 vLLM 团队

大型语言模型推理传统上基于一个简单的前提:用户提交完整的提示词(请求),模型进行处理,然后返回响应(流式传输或一次性返回)。这种范式对于基于文本的聊天机器人和批处理工作负载非常有效,但在处理实时应用(如音频或视频流)时却显得力不从心。

vLLM 最近在其引擎中增加了对可流式输入(streamable inputs)的支持,并以此为基础构建了实时 WebSocket API,在服务器中引入了一个新的 /v1/realtime 终端。

在本文中,我们将探讨实时推理的需求,并介绍 vLLM 中解锁这些功能的两个新特性:流式输入支持实时 WebSocket API

:如果您想了解如何使用 vLLM 的新流式输入或实时 API,请参考以下资料:

为什么需要实时处理

vLLM 中的传统批处理范式

传统的 LLM 推理假设完整的提示词在开始时就已经准备好。用户提交完整的请求(例如通过 ChatCompletionRequest),等待模型处理完毕,然后接收完整的响应。虽然 vLLM 长期以来一直支持输出流式传输(即在生成 Token 时即刻发出),但输入端始终是固定的:在开始推理之前,必须提供整个请求。

这种方法对于大多数应用来说已经足够。基于文本的聊天机器人、文档摘要和代码生成都自然地符合这种模型。但是,越来越多的应用无法等待输入完整后再开始处理。

流式处理的重要性

考虑用于控制电脑或手机的语音助手。所有的操作不再是通过键盘、鼠标或触摸板完成,而是由语音控制。语音由麦克风录制,并以音频流的形式发送给作为语音助手模型的 LLM。LLM 需要持续处理音频流并实时生成操作。对于此类应用,延迟(更准确地说是 首字延迟 (TTFT))非常关键——用户不希望等待超过一秒钟来打开一个应用程序、在搜索栏输入文本等。为了实现最自然、拟人的语音助手,它需要能够同时进行“听”和“说”,即 LLM 需要能够实时处理音频流并同时生成操作。

一个自然的问题是,是否可以通过将输入分块处理来使用非流式 LLM 近似实现流式行为。原则上,音频可以缓冲成片段,每个片段独立处理,并将结果拼接在一起。但在实践中,这种方法存在几个局限性。要实现亚秒级的 TTFT,需要高性能的分块检测,即准确确定何时对音频流进行分段,以确保不会丢失任何相关信息。分块不当会导致 TTFT 增加,或因碎片化了有意义的时间上下文而降低模型性能。基于分块的处理也排除了真正的双向交互:每个块必须在生成响应之前完全处理完毕,这阻碍了“听”和“说”的并发。这导致的是一种回合制交互模型,而不是人类对话中那种连续、重叠的沟通特性。

这个问题出现在许多领域:

  • 语音助手:如前所述,需要亚秒级的响应时间才能感觉自然。
  • 实时转录服务:需要在语音被识别时显示文本。
  • 机器人和具身智能:需要处理连续的传感器流(摄像头、麦克风、激光雷达)并以极小的延迟生成控制指令,以便安全地与物理世界交互。

对于这些应用,传统的批处理范式会引入不可接受的延迟。我们需要一种能够增量处理输入,并在所有输入到达之前就开始生成输出的基础设施。

:即使对于传统应用,在需要读取完整输入以生成第一个输出 Token 的情况下,利用可用输入进行流式传输依然有益。默认情况下,vLLM 使用 分块预填充(chunked prefill),因此将输入处理为个 Token 的次前向传播,其中max_num_batched_tokens 决定。在情况下,随着输入可用时即进行流式传输,可以降低整体 TTFT,因为第一次预填充前向传播可以更早地调度。

流式处理的要求

并非所有模型都能支持真正的流式推理。必须满足两个关键要求:正确的注意力模式以及针对增量处理进行的训练。

注意力模式

注意力机制决定了模型是能够增量处理输入,还是必须等待整个序列。

  • 因果注意力(单向掩码)限制每个位置仅能关注位置上的 Token,其中。由于排除了未来的 Token,模型在时刻的输出 Token 在 Token到达时即是最终结果。这使得真正的流式处理成为可能:每个新 Token 都可以立即处理,并且早期的输出永远不需要修正。

  • 双向注意力(全掩码)允许每个位置关注过去和未来的 Token。结果是,模型在位置的输出 Token 取决于尚未到达的 Token。在完整的输入序列已知之前,模型无法计算出任何位置的稳定输出,因为未来的 Token 可能会改变对早期 Token 的解读方式。

因此,双向注意力本质上要求在产生输出之前获取完整的输入序列,这使得它与流式或在线处理不兼容。

对于长时间或无限流式处理,标准的因果注意力是不够的。如果每个 Token 都关注整个过去,计算和内存将无限制增长,这是不切实际的。在实践中,过去的上下文必须截断。常见的架构解决方案是滑动窗口注意力,其中每个 Token 仅关注固定大小的近期 Token 窗口,从而在支持流式处理的同时保持计算和内存占用在可控范围内。因此,带有滑动窗口的因果注意力通常是现代流式模型的首选架构。

针对流式输入的训练

然而,拥有一个完全可流式的架构本身是不够的:模型还必须经过真正的流式输入训练。

为输入序列,为输出序列。在流式应用中,模型应生成对应于输入的输出,时间步为,且延迟尽可能低。具体而言,可以认为是音频帧的转录结果,该帧在时刻.

被流式传输到模型中。标准的下一 Token 训练目标通常将下一 Token 的分布条件设定在整个输入序列上:

这种表述不适合流式处理,因为生成需要处理完整的输入序列,这在实时环境中是无法获取的。

相反,流式模型必须能够仅使用过去的输入以及可选的少量未来上下文来预测其中

是一个超前参数,应尽可能小。理论上,可以设为零;在实践中,通常需要一点点延迟来获得合理的性能。因此,训练流式模型需要:

i) 对齐输入和输出序列,使得

  • ,且每个都是对应于的正确输出;;
  • ii) 使用一种能够在前序输入已经被处理的情况下,处理新输入的架构。

一种直观的架构,正如 延迟流建模 (Delayed Streams Modeling) 所开创,并被 Voxtral-Realtime 所采用的那样,将输入嵌入(如语音嵌入)和输出嵌入(如文本嵌入)求和池化为一个单一的嵌入序列。然后模型预测

是一个超前参数,应尽可能小。理论上,

这种区别对于部署很重要:不能简单地拿一个任意的因果模型就指望它在流式设置中表现良好。为了完全可流式化,模型必须明确地以上述对齐和架构约束进行训练,确保条件 i)ii) 均得到满足。

模型架构对于服务的重要性

vLLM 可以服务任何模型,但真正的流式传输需要架构上因果的模型。像 Voxtral 这样的模型就是从零开始为流式处理设计的,使用了支持增量处理的因果注意力机制。

同样重要的是,服务基础设施必须支持增量输入。即使有了流式模型,如果服务器在开始推理前要求提供完整的提示词,您也会失去延迟方面的优势。这就是为什么 vLLM 现在在其现有的输出流式能力之外,还支持流式输入的原因。

流式架构的进一步阅读

vLLM 中的流式输入支持

通过 PR #28973,vLLM 现在支持推理的流式输入。这实现了上述的增量处理,即输入随时间到达,输出连续生成。

StreamingInput 接口

核心抽象是 StreamingInput 数据类。

from dataclasses import dataclass
from vllm.inputs import PromptType
from vllm.sampling_params import SamplingParams
 
@dataclass
class StreamingInput:
    prompt: PromptType
    sampling_params: SamplingParams | None = None

您不再需要将固定提示词传递给 AsyncLLM.generate(),现在可以传递一个随时间产生 StreamingInput 对象的 AsyncGenerator。每个 StreamingInput 都包含要附加到累积提示词中的下一个输入块。以下是其用法示例:

import asyncio
from vllm.inputs.data import StreamingInput
from vllm.v1.engine.async_llm import AsyncLLM
from vllm.sampling_params import SamplingParams
 
async def streaming_input_example():
    async_llm = AsyncLLM.from_engine_args(...)
    
    # Input queue can consume inputs in separate async task
    input_queue = asyncio.Queue[list[int]]()
    
    async def input_generator():
        # Loop until empty list encountered => input finished
        while new_tokens := input_queue.get():
            yield StreamingInput(prompt=new_tokens)
 
    output_generator = async_llm.generate(
        prompt=input_generator(),
        sampling_params=SamplingParams(temperature=0.0, max_tokens=1),
    )
    
    # Consume outputs
    async for output in output_generator:
        # ...
 
asyncio.run(streaming_input_example())

您可以在发送下一个输入之前等待对应于上一个输入的输出完成,但这不是必需的(输入块在内部排队)。输入流的终止可以通过退出异步输入生成器或使用 aclose 函数来表示。返回的输出生成器只有在所有接收到的输入都已处理输入生成器已完成时才会结束。

工作原理

在内部,vLLM 通过将每个块视为具有累积提示词的独立请求来处理流式输入。随着新块到达,引擎会:

  1. 用生成的 max_tokens - 1output_tokens 以及新的传入 prompt_token_ids 扩展 prompt_token_ids
  2. 重用所有已缓存的 KV 值。
  3. 根据当前累积提示词和指定的 max_tokens 生成输出 Token。
  4. 当新输入到达时,选择性地丢弃输出。

这种设计意味着输入块之间生成的输出 Token 可能会随着更多上下文的可用而进行修正。最终输出反映的是完整的输入。

在内部,vLLM 通过粘性会话(sticky session)机制实现流式输入。第一个输入块创建一个在会话期间持续存在的锚点请求(anchor request)。后续具有相同内部请求 ID 的块会被排队并按顺序处理。

锚点请求模式

┌─────────────────────────────────────────────────────────────────────────────┐
│                           STREAMING SESSION                                 │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│   User's AsyncGenerator              Scheduler                              │
│   ═══════════════════               ═════════                               │
│                                                                             │
│   ┌──────────────┐                                                          │
│   │   Chunk 1    │ ──────────────►  Add ANCHOR REQUEST                      │
│   │   [A, B, C]  │                  ┌────────────────────────────────┐      │
│   └──────────────┘                  │  Request (id="session_1")      │      │
│                                     │  ├── resumable: true           │      │
│                                     │  ├── max_tokens: 2             │      │
│                                     │  ├── streaming_queue: deque()  │      │
│                                     │  ├── status: RUNNING           │      │
│                                     │  └── prompt_token_ids: [A,B,C] │      │
│                                     └────────────────────────────────┘      │
│                                              │                              │
│                                              ▼                              │
│   ┌──────────────┐                  ┌────────────────┐                      │
│   │   Chunk 2    │                  │    ENGINE      │  Generating...       │
│   │   [D, E]     │ ─────┐           │  Processing    │  ──► Output: [X, Y]  │
│   └──────────────┘      │           └────────────────┘                      │
│                         │                                                   │
│                         ▼           Anchor busy? Queue it!                  │
│   ┌──────────────┐      │           ┌────────────────────────────────┐      │
│   │   Chunk 3    │      └────────►  │  streaming_queue:              │      │
│   │   [F, G]     │ ─────────────►   │  ┌───────┐ ┌───────┐           │      │
│   └──────────────┘                  │  │[D, E] │→│[F, G] │→ ...      │      │
│                                     │  └───────┘ └───────┘           │      │
│                                     └────────────────────────────────┘      │
│                                                                             │
├─────────────────────────────────────────────────────────────────────────────┤
│                     WHEN ANCHOR FINISHES CURRENT CHUNK                      │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│   Engine signals: chunk complete (stopped = True)                           │
│          │                                                                  │
│          ▼                                                                  │
│   ┌────────────────────────────────────────────────────────────────┐        │
│   │  _handle_stopped_request() pops first item from queue          │        │
│   │                                                                │        │
│   │  streaming_queue: [[D,E], [F,G]]  ──►  [[F,G]]                 │        │
│   │                      ▲                                         │        │
│   │                      │                                         │        │
│   │                    pop!                                        │        │
│   └──────────────────────┬─────────────────────────────────────────┘        │
│                          │                                                  │
│                          ▼                                                  │
│   ┌────────────────────────────────────────────────────────────────┐        │
│   │  _update_request_as_session(anchor, update=[D, E])             │        │
│   │                                                                │        │
│   │  BEFORE:                       AFTER:                          │        │
│   │  ┌───────────────────────┐     ┌───────────────────────────┐   │        │
│   │  │ prompt_token_ids:     │     │ prompt_token_ids:         │   │        │
│   │  │   [A, B, C]           │     │   [A, B, C, X, D, E]      │   │        │
│   │  │ _output_token_ids:    │ ──► │ _output_token_ids:        │   │        │
│   │  │   [X, Y]              │     │   []                      │   │        │
│   │  │ _all_token_ids:       │     │ _all_token_ids:           │   │        │
│   │  │   [A, B, C, X, Y]     │     │   [A, B, C, X, D, E]      │   │        │
│   │  │ num_computed_tokens: 4│     │ num_computed_tokens: 4    │   │        │
│   │  │ status: RUNNING       │     │ status: WAITING           │   │        │
│   │  └───────────────────────┘     └───────────────────────────┘   │        │
│   │                                                                │        │
│   │  Note: Y is DISCARDED (last sampled token, not yet computed)   │        │
│   │        Only X is kept (num_computed_tokens = 4, so [A,B,C,X])  │        │
│   └────────────────────────────────────────────────────────────────┘        │
│                          │                                                  │
│                          ▼                                                  │
│   ┌────────────────────────────────────────────────────────────────┐        │
│   │  Anchor returns to waiting queue → scheduled again → ENGINE    │        │
│   └────────────────────────────────────────────────────────────────┘        │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

为什么下一个 prompt_token_ids 要丢弃最后一个 Token (Y)?

接收可恢复请求时,我们真正关心的是计算所有 prompt_token_ids 以及 max_tokens - 1 个已生成 Token 的 KV 缓存。注意,用户通过 max_tokens 指示第一个请求生成的 max_tokens - 1 个 Token 是最终的且应被重用。output_token_ids 张量的最后一个 Token 仅仅是最近一次前向传播的结果,它尚无对应的 KV 缓存状态。由用户决定如何处理它,由于它没有任何对应的 KV 缓存状态,它将被更新后的锚点请求的 prompt_token_ids 丢弃。丢弃它本质上是“免费”的:我们没有失效任何缓存状态,如果保留它,它在下一次迭代中也需要重新计算。

正如在 Realtime API 中所做的那样,在大多数应用中,max_tokens 应设为 1,以便每个可恢复请求仅计算 prompt_token_ids 的 KV 缓存状态,并由用户决定如何利用 output_token_ids 中生成的单个 Token。以 Voxtral Realtime 为例,output_token_ids 中生成的单个 Token 将与传入的新音频块结合,形成下一个可恢复请求。

注意:有些模型会发出模型需要正确继续生成所需的特殊停止 Token。在这种情况下,调度逻辑需要容纳 +1 Token,以便在处理新输入块之前重新计算停止 Token。

示例流程

为了说明目的,具有不同 max_tokens 的多个可恢复请求可以作为输入进行流式传输。在这种情况下,生成逻辑将如下运作:

Input chunks: ([A1, B1, C1], max_tokens=1), ([A2, B2], max_tokens=2), ([A3], max_tokens=2)

1. First chunk [A1, B1, C1] arrives
   -> Model generates [D1]

2. Second chunk [A2, B2] arrives
   -> Cumulative prompt: [A1, B1, C1, A2, B2] (D1 discarded)
   -> Model generates [C2, D2, E2]

3. Third chunk [A3] arrives
   -> Cumulative prompt: [A1, B1, C1, A2, B2, C2, D2, A3] (E2 discarded)
   -> Model generates [C3, D3]

Output stream: D1, C2, D2, E2, C3, D3

基于 WebSocket 的实时 API

虽然流式输入支持提供了核心能力,但生产应用需要一个方便的实时通信 API。PR #33187 引入了一个受 OpenAI 实时 API 启发的基于 WebSocket 的实时 API。

架构

实时 API 提供了一个 WebSocket 终端,使客户端和 vLLM 服务器之间能够进行双向流式传输。客户端发送音频数据,服务器则返回转录文本和模型输出。

其架构包括:

  1. WebSocket 客户端:从麦克风捕获音频,将数据块发送到服务器。
  2. 实时处理器:接收 WebSocket 消息,并将其转换为 StreamingInput
  3. AsyncLLM:处理流式输入,生成输出。
  4. 响应流:将生成的 Token 通过 WebSocket 发回。

服务器设置

启动支持实时 API 的 vLLM 服务器:

vllm serve mistralai/Voxtral-Mini-4B-Realtime-2602 --enforce-eager

服务器在 ws://:8000/v1/realtime 处暴露一个 WebSocket 终端。

客户端示例

这是一个流式传输音频文件并接收转录结果的基础客户端示例:

import asyncio
import base64
import json
import librosa
import numpy as np
import websockets
 
def load_audio_as_pcm16(audio_path: str) -> bytes:
    """Load audio file and convert to PCM16 @ 16kHz."""
    audio, _ = librosa.load(audio_path, sr=16000, mono=True)
    return (audio * 32767).astype(np.int16).tobytes()
 
async def stream_audio_file(audio_path: str, server_url: str = "ws://:8000/v1/realtime"):
    async with websockets.connect(server_url) as ws:
        response = json.loads(await ws.recv())
 
        # Load and convert audio to PCM16
        pcm_audio = load_audio_as_pcm16(audio_path)
 
        # Validate model
        await ws.send(json.dumps({"type": "session.update", "model": model}))
 
        # Signal start of audio stream
        await ws.send(json.dumps({"type": "input_audio_buffer.commit"}))
 
        # Stream audio in 4KB chunks
        for i in range(0, len(pcm_audio), 4096):
            chunk = pcm_audio[i:i + 4096]
            await ws.send(json.dumps({
                "type": "input_audio_buffer.append",
                "audio": base64.b64encode(chunk).decode()
            }))
 
        # Signal end of audio stream
        await ws.send(json.dumps({"type": "input_audio_buffer.commit", "final": True}))
 
        # Receive transcription
        async for message in ws:
            data = json.loads(message)
            if data["type"] == "transcription.delta":
                print(data["delta"], end="", flush=True)
            elif data["type"] == "transcription.done":
                break
 
asyncio.run(stream_audio_file("audio.wav"))

该示例展示了实时音频流的核心工作流:

  • 加载并转换音频:音频文件被加载并转换为实时 API 预期的 16kHz PCM16 格式。
  • 建立 WebSocket 连接:连接到服务器的 /v1/realtime 终端,并发送一个 session.update 消息以验证模型。
  • 分块流式传输音频:使用 input_audio_buffer.append 消息以 4KB 为单位发送音频,并使用 input_audio_buffer.commit 信号标记流的开始和结束。
  • 增量接收转录:服务器返回包含部分转录内容的 transcription.delta 消息,这些内容会实时打印,直到接收到 transcription.done
  • 关于实时行为的说明:虽然此示例为了简单起见在监听转录之前发送了所有音频,但 WebSocket 协议实现了完全异步的通信——可以同时发送音频块和接收转录。在生产级的实时服务中,转录会在第一个音频块到达时立即开始,发送和接收并行发生,从而实现真正的低延迟语音识别。

消息类型

实时 API 使用基于消息的协议。关键消息类型包括:

客户端到服务器:

  • session.create:初始化新会话。
  • input_audio_buffer.append:发送音频数据。
  • input_audio_buffer.commit:发出音频输入结束信号。
  • response.create:请求模型响应。

服务器到客户端:

  • session.created:会话初始化确认。
  • response.text.delta:增量文本输出。
  • response.audio.delta:增量音频输出(用于 TTS 模型)。
  • response.done:响应完成。
  • error:发生错误。

示例脚本

vLLM 代码库中包含了可直接使用的示例客户端:

这些示例演示了如何捕获系统麦克风的音频并将其实时流式传输到 vLLM。

性能考量

使用基于 AsyncGenerator 的专用会话接口相比仅仅发送独立请求的一个优势是,会话的 KV 缓存被原样保留。这比依赖 vLLM 的自动前缀缓存(prefix caching)更可取,因为:

  • 它确保在等待下一个输入块时,相应的缓存块不会被驱逐。
  • 前缀缓存以块级别(通常为 16 个 Token)工作,这意味着如果使用其他方式,每次新输入都会重新计算一小部分现有的 Token。

然而,这也意味着必须格外小心,避免保持会话无限期开放,因为它们会占用相应的内存,无法被其他请求使用,这可能会损害整体容量/吞吐量。目前,vLLM 不会抢占“空闲”的流式输入会话——此行为将在未来的更新中得到改进。

未来方向

我们对 vLLM 中流式输入支持的潜力感到兴奋。随着越来越多的模型提供商开源与我们输入流式设计兼容的完全可流式模型权重,我们预计实时应用生态系统将显著增长。

由于流式输入在 LLM 服务中仍然是一种新颖的功能,我们预计将适应并扩展我们的实现,以支持尽可能多样的架构和用例。这包括探索与各种音频和视频编码器更紧密的集成,针对不同延迟需求优化锚点请求模式,以及扩展对多模态流式场景的支持。

参与其中

我们鼓励您试用 vLLM 的输入流式功能和实时 API。您的反馈对我们改进这些功能至关重要。请在 vLLM GitHub 仓库中分享您的使用经验、报告问题或建议改进。

在不断发展 vLLM 实时能力的同时,我们欢迎反馈和贡献。

致谢

流式输入支持和实时 API 是通过多个团队的协作努力实现的:

Meta: Joshua Deng, Jiatong Zhou, Zhuohan Li, Yu Luo, Jeremy Teboul

Mistral AI: Patrick von Platen, Andy Lo

vLLM 团队: Nick Hill, Roger Wang, Cyrus Leung, Nicolò Lucchesi, Woosuk Kwon

我们还要感谢其他在 vLLM 中实现流式输入贡献者:Tao He (Alibaba Qwen), Edward Wibowo (Brown University), Deepti Raghavan (Brown University), 以及 Luis Gaspar Schroeder (UC Berkeley)。