配套代码:网关的流式端点(src/anna/server/app.py)、SSE 组装与解析 (src/anna/server/openai_format.pysrc/anna/providers/base.py)、 首页可视化测试台(src/anna/server/static/index.html

学完本节,你应该能回答:为什么要把一次 LLM 调用切成一块一块返回? SSE 协议长什么样?FastAPI 怎么保持低延迟转发?前端怎么逐块消费? 用户点"停止"会发生什么?流到一半断了怎么处理?长文本怎么续写?

目录

  1. 为什么需要 Streaming
  2. Server-Sent Events(SSE)
  3. FastAPI 流式接口:StreamingResponse
  4. 前端 / CLI 消费流式输出
  5. 中断与取消
  6. 流式输出异常处理
  7. 长文本输出的状态保存

1. 为什么需要 Streaming

LLM 推理不是"秒回":一次回答要先经过提示词处理、逐 token 自回归生成,短则几秒、 长则几十秒。非流式 = 客户端一直等到全部生成完才拿到第一个字,用户看到的是 一片空白;流式 = 模型每吐出一个 token 就立刻送出去,用户看到的是"边生成边打字"。

衡量这个差距有两个指标:

  • 首字节耗时(time to first byte, TTFB):从发出请求到收到第一块内容;
  • 总耗时:从发出请求到收到最后一块内容。

非流式请求的 TTFB ≈ 总耗时(必须等完整结果);流式请求的 TTFB 可以小到 “模型生成第一个 token 的时间”,总耗时不变,但感知延迟大幅下降。

本仓库首页的"流式 vs 非流式"对比面板实测过一组数据:

非流式 · 总耗时 20662 ms
流式   · 首字节 20655 ms · 总耗时 20662 ms · 分片 39 个

注意:这里流式的首字节几乎等于总耗时——这是"上游被整体缓冲后再回放" 的表现(详见第 3 节),恰好从反面说明:如果不在传输层做真正的流式, “stream: true” 就只是把最终结果切成小块,用户体验没有任何提升。

用 curl 也能直观对比。先起服务:

conda run -n llmops python main.py

非流式(不带 -N,curl 默认缓冲,只有连接结束才打印):

curl -s -o /dev/null -w '总耗时 %{time_total}s\n' \
  http://127.0.0.1:8000/v1/chat/completions \
  -H 'Content-Type: application/json' \
  -d '{"model":"auto","messages":[{"role":"user","content":"用 80 字介绍流式输出"}],"stream":false}'

流式(-N 关闭 curl 缓冲,逐块打印):

curl -N http://127.0.0.1:8000/v1/chat/completions \
  -H 'Content-Type: application/json' \
  -d '{"model":"auto","messages":[{"role":"user","content":"用 80 字介绍流式输出"}],"stream":true,"pacing_ms":200}'

pacing_ms 是本网关的演示参数:每个分片之间停顿指定毫秒数,把流式效果放慢到 肉眼可见(生产环境不需要它)。


2. Server-Sent Events(SSE)

SSE 是"基于 HTTP 长连接的文本流协议":服务器不结束连接,持续往响应体里写 data: <内容>\n\n 块,客户端一行一行读。它的几个特征:

  • Content-Type: text/event-stream,HTTP/1.1 下用 chunked 传输,连接保持打开;
  • 每个事件由空行 \n\n 分隔,事件内容以 data: 开头;
  • 结束标记是 data: [DONE](OpenAI 兼容约定,不是 SSE 标准的一部分);
  • 逐块 flush,不能积压在缓冲区里。

OpenAI 兼容的流式响应,每个 data: 里就是一个 chat.completion.chunk。 抓包看到的样子(网关 sse_chunk 的输出):

HTTP/1.1 200 OK
Content-Type: text/event-stream
data: {"id":"chatcmpl-xxx","object":"chat.completion.chunk","choices":[{"index":0,"delta":{"role":"assistant","content":""},"finish_reason":null}]}
data: {"id":"chatcmpl-xxx","object":"chat.completion.chunk","choices":[{"index":0,"delta":{"content":"流"},"finish_reason":null}]}
data: {"id":"chatcmpl-xxx","object":"chat.completion.chunk","choices":[{"index":0,"delta":{"content":"式"},"finish_reason":null}]}
data: [DONE]

组装端在 src/anna/server/openai_format.py,核心就几行:

def sse_chunk(request_id, model_label, delta, finish_reason=None) -> str:
    data = {
        "id": request_id,
        "object": "chat.completion.chunk",
        "created": created(),
        "model": model_label,
        "choices": [{"index": 0, "delta": delta, "finish_reason": finish_reason}],
    }
    return f"data: {json.dumps(data, ensure_ascii=False)}\n\n"
SSE_DONE = "data: [DONE]\n\n"

消费端(provider 层)在 src/anna/providers/base.py_iter_sse,逐行解析、 跳过空行和 [DONE]、容错跳过非 JSON 行:

def _iter_sse(self, response: httpx.Response) -> Iterator[Dict[str, Any]]:
    """从 SSE 响应解析 data: JSON 行(兼容各家事件格式)。"""
    for line in response.iter_lines():
        line = (line or "").strip()
        if not line.startswith("data:"):
            continue
        data = line[len("data:"):].strip()
        if not data or data == "[DONE]":
            continue
        try:
            yield json.loads(data)
        except json.JSONDecodeError:
            continue

一个常见坑:多字节字符被截断。 SSE 是行协议,但 TCP 分片不会管 “一个汉字是两个字节”——"流" 的 UTF-8 编码可能被切成两半到达。 所以消费端不能用"按字节切字符串"的方式,要用带 stream: true 的解码器 (见第 4 节前端代码)。


3. FastAPI 流式接口:StreamingResponse

FastAPI 里做流式接口的核心是 StreamingResponse:给它一个生成器,FastAPI 会逐块把生成器的产出写回客户端,而不是等生成器结束再一次性返回。

本仓库的流式端点 src/anna/server/app.py

if req.stream:
    return _stream_response(
        provider=provider,
        provider_name=provider_name,
        model_label=model_label,
        request_id=request_id,
        messages=messages,
        temperature=req.temperature,
        max_tokens=req.max_tokens,
        response_format=req.response_format,
        pacing_ms=req.pacing_ms,   # 演示用:放慢分片
        usage_store=usage,
        started_at=start,
    )

_stream_response 的骨架:

def _stream_response(*, provider, model_label, request_id, messages,
                     temperature, max_tokens, response_format,
                     pacing_ms, usage_store, started_at) -> StreamingResponse:
    def generate():
        # 1) 先发一个带 role 的空块(OpenAI 兼容格式要求)
        yield sse_chunk(request_id, model_label, {"role": "assistant", "content": ""})
        try:
            # 2) 把上游增量逐块转发
            for delta in provider.stream(
                messages,
                temperature=temperature,
                max_tokens=max_tokens,
                response_format=response_format,
            ):
                if pacing_ms:
                    time.sleep(max(0, pacing_ms) / 1000.0)  # 演示模式
                yield sse_chunk(request_id, model_label, {"content": delta})
            yield sse_chunk(request_id, model_label, {}, finish_reason="stop")
            # 3) 用量、成本、延迟记录(省略)
        except ProviderError as exc:
            usage_store.record(UsageRecord(..., status="error", error=str(exc)))
            yield f"data: {error_payload(str(exc), 'provider_error', exc.status_code)}\n\n"
        finally:
            yield SSE_DONE
    return StreamingResponse(
        generate(),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )

三个容易被忽略的细节:

  • Cache-Control: no-cache:防止浏览器/代理把整个响应缓冲起来;
  • X-Accel-Buffering: no:告诉 nginx 这类反向代理不要缓冲 SSE;
  • 当前用的是同步生成器def generate + time.sleep),FastAPI 会把它放到 线程池跑;如果要严格低开销,可以改成 async def + await asyncio.sleep

3.1 真正的低延迟,还需要上游也流式

StreamingResponse 只是"网关 → 浏览器"这一段。如果"上游模型 → 网关"用的是 整体缓冲的请求,首字节照样要等上游全部生成完。本仓库把"上游 → 网关"也做成了 真流式:base.py_stream_ssehttpx.Client.stream 逐行读取,第一个分片 一到就 yield,而不是等完整响应体:

def _stream_sse(self, url, payload, *, headers=None) -> Iterator[Dict[str, Any]]:
    """流式 SSE 请求:真低延迟转发 + 首字节前指数退避重试。"""
    for attempt in range(self.settings.max_retries + 1):
        if attempt:
            time.sleep(0.5 * (2 ** (attempt - 1)))
        yielded = False
        try:
            with self._client.stream("POST", url, json=payload, headers=headers) as response:
                if response.status_code >= 300:
                    err = self._response_error(response)
                    if (
                        response.status_code not in TRANSIENT_STATUS
                        or attempt == self.settings.max_retries
                    ):
                        raise err
                    continue
                for chunk in self._iter_sse(response):
                    yield chunk
                    yielded = True
                return
        except httpx.HTTPError as exc:
            if yielded:
                raise ProviderError(f"流式中断:{exc}")  # 已流出内容不再重试
            if attempt == self.settings.max_retries:
                raise ProviderError(f"请求失败:{exc}")
            continue

openai_compatible.py / anthropic.py / gemini.py 三个 adapter 的 stream 都改用它;client.streamwith 块退出即关闭上游连接——这也是第 5 节"取消 传播"能生效的基础。

两种方式的差距可以量化(5 个分片、每个间隔 0.2s,MockTransport 实测):

方式首行耗时说明
client.request()1.02s(整体完成后)服务端生成再快,首字节也要等全部
client.stream()0.00s第一块一到就可见,后续按上游节奏到

改这个之前,首页对比面板实测"首字节 20655ms ≈ 总耗时 20662ms";改成 client.stream 之后,首字节应该从"≈ 总耗时"变成"远小于总耗时"。


4. 前端 / CLI 消费流式输出

4.1 为什么不用 EventSource

浏览器原生 EventSource 能消费 SSE,但有三个限制:只能 GET、不能带自定义 请求头、不能带请求体。聊天请求要 POST JSON,所以本仓库前端用的是 fetch + ReadableStreamresponse.body.getReader()),EventSource 适合 “服务端主动推送"的场景(通知、状态广播),且它自带自动重连。

4.2 前端逐块消费(首页 index.html

核心循环:getReader().read()TextDecoder(stream: true) 解码 → 按行切分 → 只处理 data: 行 → JSON.parse 取出 delta.content

const res = await fetch("/v1/chat/completions", {
  method: "POST",
  headers: { "Content-Type": "application/json" },
  body: JSON.stringify({ ...base, stream: true, pacing_ms: 80 }),
});
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buf = "";
while (true) {
  const { done, value } = await reader.read();
  if (done) break;
  buf += decoder.decode(value, { stream: true }); // 关键:处理多字节字符被截断
  const lines = buf.split("\n");
  buf = lines.pop(); // 最后一段可能不完整,留到下一轮
  for (const line of lines) {
    const t = line.trim();
    if (!t.startsWith("data:")) continue;
    const payload = t.slice(5).trim();
    if (payload === "[DONE]") continue;
    const obj = JSON.parse(payload);
    const delta = obj.choices?.[0]?.delta?.content;
    if (delta) typewriter.push(delta);
  }
}

“打字机"效果:如果直接 textContent += delta,速度快时看起来仍像一次性输出, 所以页面用一个队列逐片渲染(每片 40ms),即使浏览器把整段缓冲了,也会按顺序 逐字"回放”,保证流式效果可见:

function createTypewriter(appendFn, intervalMs = 40) {
  const queue = [];
  let timer = null;
  function drain() {
    if (!queue.length) { timer = null; return; }
    appendFn(queue.shift());
    timer = setTimeout(drain, intervalMs);
  }
  return {
    push(delta) {
      queue.push(delta);
      if (!timer) timer = setTimeout(drain, intervalMs);
    },
    pending() { return queue.length > 0 || timer !== null; },
    flush() { while (queue.length) appendFn(queue.shift()); timer = null; },
  };
}

注意设计选择:页面故意不调用 flush(),让尾部几片也按队列节奏显示, 并显示"正在逐片回放…",这样用户能明确看到流式是分片的。

4.3 CLI 消费

命令行最轻量的方式是 curl -N(见第 1 节)。Python 客户端用 httpx 的流式 API, 和服务端 _iter_sse 是同一套思路:

import httpx
with httpx.stream(
    "POST",
    "http://127.0.0.1:8000/v1/chat/completions",
    json={
        "model": "auto",
        "messages": [{"role": "user", "content": "讲个笑话"}],
        "stream": True,
    },
) as r:
    for line in r.iter_lines():
        if not line.startswith("data:"):
            continue
        payload = line[5:].strip()
        if not payload or payload == "[DONE]":
            continue
        print(payload[:120])

5. 中断与取消

为什么需要:生成长、按 token 计费,用户看完开头觉得不对,应该能随时点 “停止”,让上游立刻停止生成,而不是默默等完、白付整段 token 的钱。

当前项目状态:已实现(首页有”⏹ 停止生成"按钮,网关在客户端断开时关闭 上游流)。三个层面缺一不可,这也是面试里常问的"cancel 是怎么贯穿全链路的":

  1. 前端AbortController + fetch 的 signal,点停止时 ctrl.abort(), 浏览器立即断开连接:
const ctrl = new AbortController();
fetch("/v1/chat/completions", {
  method: "POST",
  headers: { "Content-Type": "application/json" },
  body: JSON.stringify(body),
  signal: ctrl.signal,          // 传入 signal
});
// 用户点"停止":
ctrl.abort();
  1. 网关_stream_response 改成 async 生成器,每取一个分片都用 asyncio.to_thread 在 worker 线程跑同步 provider,事件循环保持可取消; 客户端断开时生成器收到 GeneratorExit/CancelledError,在异常处理里 关闭上游连接stream.close()):
def generate():
    try:
        for delta in provider.stream(messages, ...):
            yield sse_chunk(request_id, model_label, {"content": delta})
    except GeneratorExit:
        # 客户端断开:结束生成器,让 provider 的 finally/with 关闭上游流
        raise
    finally:
        yield SSE_DONE
  1. 上游:provider 的 streamclient.streamwith 块管理连接, 生成器被 close 时 with 退出、HTTP 连接关闭,模型端收到连接关闭才会真正停止 生成、停止计费。

一句话总结:取消要贯穿三层才有效——前端 abort → 网关 detect → 上游 close。 已知边界:如果取消时上游恰好正在产出一个分片,可能再多一个分片(或等读取 超时)连接才关闭;严格的即时中断需要把 provider 也做成可取消的 async 客户端。


6. 流式输出异常处理

6.1 网关:流中途出错,往流里写错误块

非流式出错可以整体返回 4xx/5xx;流式不行——响应头已经发出、状态码已是 200, 所以错误只能作为流中的事件发给客户端。本仓库的做法(app.py):

except ProviderError as exc:
    usage_store.record(
        UsageRecord(
            request_id=request_id, provider=provider_name, model=provider.model,
            status="error", error=str(exc),
        )
    )
    yield f"data: {error_payload(str(exc), 'provider_error', exc.status_code)}\n\n"
finally:
    yield SSE_DONE

6.2 重试:只能在"首字节之前"

_request_json 对 429/5xx 做指数退避重试(0.5s → 1s → 2s…),但这只对 整体请求有效:一旦第一个分片已经吐给客户端,就不能重试了——重试会重复 已经看到的内容。所以流式重试的边界是:

  • 首字节前:可以重试(_request_json 的逻辑);
  • 首字节后:只能报错,由客户端决定是否重新发起(客户端要保证请求幂等)。

6.3 跨模型回退:只能在"首字节前"切换

非流式可以在失败后整段换一个模型重试(客户端无感知);流式不行——用户已经 看到了一半内容,这时候切换模型会导致前后文风断裂。所以流式回退的正确边界是: 首字节前失败可切换,一旦流出过内容就只抛错FallbackProvider.stream 按这个边界实现:

def stream(self, messages, ...) -> Iterator[str]:
    """流式 + 回退:只在"首字节前"失败时切换下一个 provider。"""
    providers = []
    if self._primary_available():
        providers.append(self._primary)
    providers.extend(self._fallbacks)
    for provider in providers:
        stream = provider.stream(messages, ...)
        yielded = False
        try:
            for delta in stream:
                yielded = True
                yield delta
            return
        except ProviderError as exc:
            if classify_provider_error(exc) not in FALLBACKABLE_KINDS:
                raise
            if yielded:
                raise  # 已流出内容:不能切换
            continue  # 首字节前失败:试下一个 provider

配合第 3.1 节的 _stream_sse,“首字节前多试几次"有两层:同一上游 429/5xx 先指数退避重试;重试耗尽再由 FallbackProvider 换下一个模型。测试在 tests/test_streaming.py::FallbackStreamTest

6.4 消费端兜底

前端对 reader.read() 的异常(网络抖动、连接被重置)要 catch 并提示; 服务端 _iter_sse 对非 JSON 行静默跳过(except json.JSONDecodeError: continue), 避免个别脏块搞挂整条流。


7. 长文本输出的状态保存

场景:生成一篇长文(比如整本书的导读)时,可能因为 ① 超出 context window、 ② 网络/会话断开,导致任务中途失败。非流式可以整个重试,长文本重试的成本 太高,需要断点续写:把"生成"做成可恢复的 job。

三个要点:

  1. 持久化中间状态job_id + 消息历史 + 已生成部分 + 断点偏移,每生成 一段就追加保存(类似本仓库 usage.py 的 JSONL 落盘思路):
class GenerationJob:
    job_id: str
    model: str
    messages: list          # 完整对话历史
    partial: str            # 已生成文本
    cursor: int             # 断点:下一次续写的偏移
    status: str             # running / cancelled / done
    created_at: datetime
# 落盘:jobs.jsonl 每行一个 job 快照;断线后按 job_id 恢复
  1. 超出 context window:讲义 01 讲过三种手段——分块、摘要、缓存。 生成侧对应的是"先规划章节 → 每章独立生成 → 各章结果持久化”,而不是 一个超长 prompt 一次生成完。

  2. 恢复语义:重连后客户端请求 GET /v1/generations/{job_id},拿到 partial,服务端从 cursor 继续;要保证"续写请求"不会重复生成已有部分。

当前项目状态:已实现 v1generations.py + POST/GET /v1/generations + 首页"长文本断点"面板):

  • GenerationJob 记录 模型 / 消息 / 已生成 partial / 断点字符数 / 状态;
  • 后台线程逐块消费 provider.stream 并追加落盘(JSONL,ANNA_GENERATIONS_LOG 可选,最多 1s 一次节流);
  • 取消后 partial 保持可见,POST /v1/generations/{id}/continue 从断点继续;
  • 任务状态 running / done / cancelled / error,断线后 GET 可查。

v1 的 continue 是"保留已生成内容、重新跑一遍并追加";生产环境应升级为 “按章节规划 + 从断点续写下一节”(本节的思路),避免重复生成。面试时能说出 “把长生成做成有断点的 job,且 continue 的重复生成问题要怎么演进"就是加分项。


小结

  • 流式 = 把"总耗时"换挡成"首字节耗时”,用户感知变快,总耗时不变;
  • SSE 是文本行协议:data: JSON + \n\n 分隔、[DONE] 结束,消费端要处理 多字节字符截断;
  • StreamingResponse 管"网关 → 客户端",client.stream 管"上游 → 网关", 两段都流式才是真流式(本仓库两段都已实现,改前"首字节≈总耗时"的问题 已消除);
  • 前端用 fetch + ReadableStream 逐块消费,EventSource 不适合 POST 聊天;
  • 取消要三层贯通:前端 abort → 网关 detect → 上游 close,否则 token 白付;
  • 流式错误:首字节前可重试/可换模型,首字节后只能报错/客户端重发;
  • 长文本要持久化断点,把"生成"做成可恢复的 job(已实现 v1)。

动手练习

  1. 启动网关,用 curl -N 与不带 -N 各跑一次,观察输出时序差异;
  2. 首页"流式 vs 非流式"面板勾选/取消"放慢流式显示",对比首字节与总耗时;
  3. pacing_ms=500 放慢到极限,观察前端逐片渲染与分片表;
  4. ANNA_RATE_LIMIT_RPM 调小触发 429,观察流式请求的错误处理;
  5. 点首页「模拟流式错误」,观察"前半段 + SSE 错误块 + [DONE]“的完整链路;
  6. 首页「长文本断点」:开始生成 → 取消 → 查看 partial → 继续,观察断点续写;
  7. (验证)对比面板的"首字节 vs 总耗时”:真流式下首字节应远小于总耗时。

附录:讲义内容 × 项目实现对照

快速核对"讲义里讲的,仓库里有没有,首页能不能演示"。

讲义条目项目实现位置状态首页示例
流式 vs 非流式对比(首字节 vs 总耗时)首页对比面板 static/index.html✅ 已实现#chat「对比演示(同一消息 × 两种模式)」:流式栏显示首字节/总耗时/分片数与分片表
SSE 协议组装(data: + [DONE])server/openai_format.py sse_chunk / SSE_DONE✅ 已实现流式对话逐字显示 + 分片表(delta 列表);引导区有 curl -N 示例
SSE 解析(容错跳过脏行)providers/base.py _iter_sse✅ 已实现(有测试)无直接展示(浏览器端解析逻辑,内部行为)
FastAPI StreamingResponse 转发server/app.py _stream_response✅ 已实现(有测试)流式对话 / 对比面板能观察到分片逐块到达
上游低延迟转发(client.stream + 首字节前重试)providers/base.py _stream_sse,三家 adapter 均接入✅ 已实现(有测试)对比面板可验证"首字节 < 总耗时"
演示放慢分片(pacing_ms)网关 + server/schemas.py + 首页✅ 已实现#chat「放慢流式显示(80ms/分片)」勾选框
前端 ReadableStream 逐块消费static/index.html fetch + getReader✅ 已实现整个对话测试面板就是活例子
打字机回放(createTypewriter)static/index.html✅ 已实现流式回复逐字显示 +「正在逐片回放…」提示
EventSource—(GET-only,不适合 POST 聊天)📘 通用说明无(页面用 fetch,讲义 4.1 说明原因)
CLI 消费(curl -N / httpx.stream)—(scripts/ 无对应 CLI)📘 讲义示例引导区有 curl -N 快速上手示例
中断与取消(前端 abort → 网关 → 上游 close)static/index.html 停止按钮 + server/app.py async 生成器✅ 已实现(有测试)#chat「⏹ 停止生成」按钮:点停止后立即停、提示"已停止生成"
流式中途错误处理(SSE 错误块 + 用量记录)server/app.py except 分支 + providers/demo.py✅ 已实现(有测试)#chat「模拟流式错误」按钮:看到"前半段 + 错误块 + [DONE]";用量面板记录 error
流式中途重试 / fallback 切换providers/fallback.py + base.py _stream_sse(首字节前可切换)✅ 已实现(有测试)无专门按钮(可把主模型名写错,观察自动切 fallback)
长文本断点持久化 / 续写generations.py + /v1/generations + 首页✅ 已实现(有测试)#generations 面板:开始生成 → 取消 → 查看 partial → 继续

图例:✅ 已实现并有测试 / ⬜ 未实现(讲义中已注明或属后续阶段)/ 📘 通用知识或示例说明。