1. 概述
先从一个普通聊天接口开始。客户端提交问题后,服务器调用模型,等模型生成完整答案,再一次性返回 JSON:
1客户端发出请求2↓3服务器等待模型生成完整答案4↓5服务器返回响应头和完整响应体6↓7客户端开始显示答案
假设模型用 20 秒生成答案,那么用户在前 20 秒里什么也看不到。即使服务器一直在工作,页面看起来也像是没有响应。
流式响应改变了返回答案的方式。服务器拿到第一小段内容后就立即发送,后续内容生成一段、发送一段:
1客户端发出请求2↓3服务器返回响应头4↓5第 1 段内容 → 第 2 段内容 → 第 3 段内容 → 响应结束
这里要区分两个时间:
- 首块时间:从发出请求到收到第一段内容所需的时间;
- 完整时间:从发出请求到收到全部内容所需的时间。
流式响应主要缩短的是用户感知到的首块时间,不一定能缩短模型生成完整答案的时间。它让用户更早看到结果,也允许页面边接收边渲染,并为停止生成、进度提示和工具调用状态提供了基础。
不过,流式响应也带来了新的问题:客户端如何知道一段内容从哪里开始、在哪里结束?流已经开始后发生错误怎么办?用户关闭页面后,服务器如何停止模型?网络断开后能不能继续?这正是后面要重点学习 SSE 的知识内容。
2. 响应、数据块与事件
阅读流式代码时,经常会看到响应、数据块和事件这三个词。它们不是同一个概念:
| 名称 | 所在层级 | 含义 |
|---|---|---|
| HTTP 响应 | HTTP 层 | 一次请求对应的一份响应,由响应头和响应体组成 |
| 网络数据块 | 传输层 | 响应体在网络传输时被拆成的若干字节片段 |
| SSE 事件 | 应用协议层 | 按 SSE 格式组织的一条完整业务消息 |
一次流式请求仍然只有一个 HTTP 响应。服务器不会为每个模型片段重新创建一次响应,而是在同一个响应体中持续写入数据。
网络数据块的边界也没有业务含义。一次 yield 可能被底层拆成多个网络块,多次 yield 也可能合并后才到达客户端。一个 UTF-8 汉字、一个 JSON 对象,甚至 SSE 中的一行文本,都可能跨越两个网络块。因此,客户端不能对每次 reader.read() 得到的内容直接执行 JSON.parse(),而要先按照协议重新拼出完整事件。
3. StreamingResponse
FastAPI 的 StreamingResponse 可以把普通迭代器或异步迭代器产生的内容持续写入响应体。我们先看一个不涉及模型和 SSE 的最小示例:
01# src/fastapi/api/stream/text_stream.py02import asyncio03from collections.abc import AsyncIterator0405from fastapi import APIRouter06from fastapi.responses import StreamingResponse0708router = APIRouter()091011class TextStreamingResponse(StreamingResponse):12media_type = "text/plain"131415# 请求路径:GET /api/v1/stream/text/stream16@router.get(17"/text/stream",18response_class=TextStreamingResponse,19)20async def stream_text() -> AsyncIterator[str]:21for chunk in ("FastAPI", " 可以", " 流式", " 输出"):22await asyncio.sleep(0.2)23yield chunk
这段代码的执行过程如下:
- 客户端请求
/api/v1/stream/text/stream - FastAPI 创建
TextStreamingResponse,响应类型是text/plain - 路径函数运行到第一个
yield,把FastAPI交给响应层 - 函数暂停在当前位置,响应层把内容写给客户端
- 下一次迭代从暂停处继续,等待 0.2 秒后产生下一段内容
- 循环结束,生成器结束,FastAPI 关闭响应体
yield 的作用是产生一段值并暂停函数,await 的作用是暂停当前协程并把执行机会交还事件循环。二者解决的问题不同。真实模型的 astream() 内部通常包含异步网络读取,所以自然会到达 await;如果自定义异步生成器只有长时间的 CPU 计算,即使不断 yield,取消信号也可能无法及时执行
启动本地服务之后,我们可以通过下面这个案例来感受一下这个接口的调用过程
浏览器拿到的是连续的 text/plain 字节流,并不知道 Python 端一共执行了几次 yield。页面最后显示为一段完整文本,也不代表浏览器恰好执行了四次 reader.read();服务器生成边界和网络读取边界仍然是两回事
TextStreamingResponse 把媒体类型固定为 text/plain。因为它属于文本类型,响应层还会补充 UTF-8 字符集。我们也可以显式创建响应对象:
01# src/fastapi/api/stream/text_stream.py02from collections.abc import AsyncIterable0304from fastapi.responses import StreamingResponse0506async def generate_text() -> AsyncIterable[str]:07yield "第一段"08yield "第二段"091011async def stream_text() -> StreamingResponse:12return StreamingResponse(13generate_text(),14media_type="text/plain; charset=utf-8",15)1617# GET /api/v1/stream/text/response18@router.get(19"/text/response",20summary="显式创建纯文本流响应",21)22async def stream_text_response() -> StreamingResponse:23"""在运行时显式传入迭代器、媒体类型和字符集。"""2425return StreamingResponse(26generate_text(),27media_type="text/plain; charset=utf-8",28)
显式返回 StreamingResponse 适合在运行时决定响应头、状态码或媒体类型。直接 yield 更简洁,也能让 FastAPI 根据返回类型进行文档生成和数据处理。两种形式最终都要把一个可迭代对象交给响应层。
StreamingResponse 同时支持同步生成器和异步生成器,但选择时要看数据来源:
| 数据来源 | 推荐形式 | 原因 |
|---|---|---|
| 异步模型 SDK、异步 HTTP 客户端 | async def + async for | 等待网络时不会阻塞事件循环 |
| 普通文件对象、同步迭代库 | def + yield | FastAPI 可以在线程池中迭代 |
| CPU 密集计算 | 进程池或任务系统 | 改成异步函数不会自动消除 CPU 阻塞 |
不要在异步生成器里直接调用耗时的同步 SDK。它会占住事件循环,使同一进程中的其他请求、心跳和取消处理一起停下来。如果供应商只提供同步流式接口,应使用同步生成器,让框架在线程池中迭代,或者把阻塞调用明确放到线程中执行。
3. 接入 LangChain
LangChain 的 astream() 会返回异步迭代器,可以直接接入 StreamingResponse:
01# src/fastapi/api/stream/langchain_stream.py02from collections.abc import AsyncIterator03from typing import Annotated0405from fastapi import APIRouter, Depends06from fastapi.responses import StreamingResponse0708from .common import ChatRequest, DemoChatService, get_chat_service0910router = APIRouter()1112# POST /api/v1/stream/langchain13@router.post("/langchain")14async def stream_chat(15payload: ChatRequest,16service: Annotated[DemoChatService, Depends(get_chat_service)],17) -> StreamingResponse:18async def generate() -> AsyncIterator[str]:19async for chunk in service.astream(payload.question):20yield chunk2122return StreamingResponse(23generate(),24media_type="text/plain; charset=utf-8",25)
真实 LangChain 模型的 chunk 通常是消息块,chunk.content 才是模型产生的内容。一个消息块不保证恰好对应模型分词器中的一个 token,也不保证是一句话,所以后文虽然沿用常见的 token 事件名,客户端仍应把它理解为文本增量。当前项目的 DemoChatService.astream() 已经把消息块转换为字符串,路由可以直接转发。
不同模型集成返回的 content 结构可能不同。有的返回字符串,有的返回内容块列表,有的还会混入工具调用参数。真实项目要先确认当前模型的消息类型,再决定如何提取文本,不能无条件执行 str(chunk.content),否则结构化工具信息也可能被错误地显示给用户。
还要注意,调用了 astream() 不等于整条链一定能够逐步产出。链中的模型、解析器和自定义节点都要支持流式处理;只要中间某个步骤必须等待完整结果,客户端就可能在那一步暂时收不到新内容。
下面这个案例中,客户端的输入会作为 {"question":"..."} 发送给 POST /api/v1/stream/langchain。接口仍然返回纯文本,所以客户端只负责逐块解码和追加,不会出现 event、id 或 data 等 SSE 字段
你可以把这个案例与后面的 /chat 放在一起比较:两者使用同一个 DemoChatService.astream(),但 /langchain 只返回正文,/chat 还会返回开始、结束、错误和事件 ID 等协议信息
关于 langchain 更详细的使用,我们会在后续章节中深入学习
4. 对接 deepseek
前面的案例没有请求真实的 LLM 模型。这个案例,我们使用 deepseek 的 SSE 接口来对接流式响应
所有的配置,我们都配置在项目根目录的 .env 文件中:
1DEEPSEEK_MODEL=deepseek-chat2DEEPSEEK_BASE_URL=https://api.deepseek.com3DEEPSEEK_API_KEY=your-api-key
下面是案例和对应的代码实现
001# src/fastapi/api/stream/deepseek_sse.py002"""使用 .env 中的 DeepSeek 配置调用真实模型的 SSE 接口。"""003004import asyncio005import logging006from collections.abc import AsyncIterator007from typing import Annotated008from uuid import uuid4009010from fastapi import APIRouter, Depends, HTTPException011from fastapi.sse import EventSourceResponse, ServerSentEvent012from langchain_openai import ChatOpenAI013014from app.core.config import Settings, get_settings015from app.model import create_model016017from .common import ConversationRequest018019logger = logging.getLogger(__name__)020router = APIRouter()021022023def get_deepseek_model(024settings: Annotated[Settings, Depends(get_settings)],025) -> ChatOpenAI:026"""在流启动前创建模型,缺少密钥时返回普通 HTTP 错误。"""027028try:029return create_model(settings)030except RuntimeError as exc:031raise HTTPException(032status_code=503,033detail="DeepSeek 模型配置不完整",034) from exc035036037# POST /api/v1/stream/deepseek038@router.post(039"/deepseek",040response_class=EventSourceResponse,041summary="DeepSeek 流式聊天",042)043async def stream_deepseek(044payload: ConversationRequest,045model: Annotated[ChatOpenAI, Depends(get_deepseek_model)],046) -> AsyncIterator[ServerSentEvent]:047"""调用 DeepSeek,并把 LangChain 消息块转换成 SSE 生命周期事件。"""048049message_id = f"msg_{uuid4().hex}"050sequence = 0051usage: dict[str, object] | None = None052053yield ServerSentEvent(054event="message_start",055id=f"{message_id}:{sequence}",056retry=3_000,057data={058"message_id": message_id,059"model": model.model_name,060},061)062063try:064async for chunk in model.astream(payload.to_langchain_messages()):065if chunk.usage_metadata is not None:066usage = dict(chunk.usage_metadata)067068# ``text`` 会从字符串或标准文本内容块中提取可展示文本。069if not chunk.text:070continue071072sequence += 1073yield ServerSentEvent(074event="token",075id=f"{message_id}:{sequence}",076data={"text": chunk.text},077)078except asyncio.CancelledError:079# 客户端断开时让取消信号继续传播,以便终止上游读取。080raise081except Exception:082logger.exception("DeepSeek streaming request failed")083sequence += 1084yield ServerSentEvent(085event="stream_error",086id=f"{message_id}:{sequence}",087data={088"code": "model_unavailable",089"message": "模型暂时不可用,请稍后重试",090},091)092return093094sequence += 1095yield ServerSentEvent(096event="message_end",097id=f"{message_id}:{sequence}",098data={099"finish_reason": "stop",100"usage": usage,101},102)
最终请求路径是 POST /api/v1/stream/deepseek:
01curl -iN -X POST \02http://127.0.0.1:8000/api/v1/stream/deepseek \03-H 'Content-Type: application/json' \04-d '{05"messages": [06{"role":"system","content":"你是 Python 学习助手。"},07{"role":"user","content":"什么是 FastAPI?"},08{"role":"assistant","content":"FastAPI 是一个 Python Web 框架。"},09{"role":"user","content":"它支持异步吗?"}10]11}'
这个案例与前面的单轮 /chat 案例最重要的区别,在于请求体:每次提问都会把当前
UI 组件中已经完成的 user 和 assistant 消息连同新问题一起提交。
可以先发送“请记住:我正在学习 FastAPI”,等回答结束后再问“我刚才说我在学习 什么?”。第二轮请求不是只发送最后一个问题,而是发送类似下面的完整数组:
1{2"messages": [3{"role": "user", "content": "请记住:我正在学习 FastAPI。"},4{"role": "assistant", "content": "好的,我记住了。"},5{"role": "user", "content": "我刚才说我在学习什么?"}6]7}
5. 选择流式协议
纯文本流只能表达正文。如果服务器还要告诉客户端 消息开始了、正在执行工具、已经结束 、或者 生成失败,就需要在文本之上再约定一种格式
常见方案如下:
| 协议 | 数据组织方式 | 适用场景 |
|---|---|---|
| 纯文本 | 只有连续文本,没有事件边界 | 只展示正文 |
| JSON Lines | 每行一个独立 JSON | 程序间批量数据流 |
| SSE | 用空行分隔事件,内置 data、event、id 等字段 | AI 回复、通知、日志 |
| WebSocket | 建立独立的全双工消息通道 | 语音对话、协作编辑、双方频繁主动发消息 |
AI 聊天通常是客户端先提交一次问题,随后服务器持续返回结果。这个过程仍然符合普通 HTTP 的请求与响应模型,SSE 已经能够表达正文、状态、错误和恢复信息。WebSocket 功能更多,但也会增加连接管理、心跳、鉴权、重连和跨进程广播的复杂度,不应该成为默认选择。