创建时间: 2026-08-25最后更新: 2026-08-27

1. 概述

先从一个普通聊天接口开始。客户端提交问题后,服务器调用模型,等模型生成完整答案,再一次性返回 JSON:

code.ts
1
客户端发出请求
2
↓
3
服务器等待模型生成完整答案
4
↓
5
服务器返回响应头和完整响应体
6
↓
7
客户端开始显示答案

假设模型用 20 秒生成答案,那么用户在前 20 秒里什么也看不到。即使服务器一直在工作,页面看起来也像是没有响应。

流式响应改变了返回答案的方式。服务器拿到第一小段内容后就立即发送,后续内容生成一段、发送一段:

code.ts
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 的最小示例:

text_stream.py
01
# src/fastapi/api/stream/text_stream.py
02
import asyncio
03
from collections.abc import AsyncIterator
04
05
from fastapi import APIRouter
06
from fastapi.responses import StreamingResponse
07
08
router = APIRouter()
09
10
11
class TextStreamingResponse(StreamingResponse):
12
media_type = "text/plain"
13
14
15
# 请求路径:GET /api/v1/stream/text/stream
16
@router.get(
17
"/text/stream",
18
response_class=TextStreamingResponse,
19
)
20
async def stream_text() -> AsyncIterator[str]:
21
for chunk in ("FastAPI", " 可以", " 流式", " 输出"):
22
await asyncio.sleep(0.2)
23
yield chunk

这段代码的执行过程如下:

  1. 客户端请求 /api/v1/stream/text/stream
  2. FastAPI 创建 TextStreamingResponse,响应类型是 text/plain
  3. 路径函数运行到第一个 yield,把 FastAPI 交给响应层
  4. 函数暂停在当前位置,响应层把内容写给客户端
  5. 下一次迭代从暂停处继续,等待 0.2 秒后产生下一段内容
  6. 循环结束,生成器结束,FastAPI 关闭响应体

yield 的作用是产生一段值并暂停函数,await 的作用是暂停当前协程并把执行机会交还事件循环。二者解决的问题不同。真实模型的 astream() 内部通常包含异步网络读取,所以自然会到达 await;如果自定义异步生成器只有长时间的 CPU 计算,即使不断 yield,取消信号也可能无法及时执行

启动本地服务之后,我们可以通过下面这个案例来感受一下这个接口的调用过程

直接 yield 纯文本GET /api/v1/stream/text/stream · text/plain
本地服务:心跳检测中

浏览器拿到的是连续的 text/plain 字节流,并不知道 Python 端一共执行了几次 yield。页面最后显示为一段完整文本,也不代表浏览器恰好执行了四次 reader.read();服务器生成边界和网络读取边界仍然是两回事

TextStreamingResponse 把媒体类型固定为 text/plain。因为它属于文本类型,响应层还会补充 UTF-8 字符集。我们也可以显式创建响应对象:

return-streaming-response.py
01
# src/fastapi/api/stream/text_stream.py
02
from collections.abc import AsyncIterable
03
04
from fastapi.responses import StreamingResponse
05
06
async def generate_text() -> AsyncIterable[str]:
07
yield "第一段"
08
yield "第二段"
09
10
11
async def stream_text() -> StreamingResponse:
12
return StreamingResponse(
13
generate_text(),
14
media_type="text/plain; charset=utf-8",
15
)
16
17
# GET /api/v1/stream/text/response
18
@router.get(
19
"/text/response",
20
summary="显式创建纯文本流响应",
21
)
22
async def stream_text_response() -> StreamingResponse:
23
"""在运行时显式传入迭代器、媒体类型和字符集。"""
24
25
return StreamingResponse(
26
generate_text(),
27
media_type="text/plain; charset=utf-8",
28
)
显式创建 StreamingResponseGET /api/v1/stream/text/response · text/plain
本地服务:心跳检测中

显式返回 StreamingResponse 适合在运行时决定响应头、状态码或媒体类型。直接 yield 更简洁,也能让 FastAPI 根据返回类型进行文档生成和数据处理。两种形式最终都要把一个可迭代对象交给响应层。


StreamingResponse 同时支持同步生成器和异步生成器,但选择时要看数据来源:

数据来源推荐形式原因
异步模型 SDK、异步 HTTP 客户端async def + async for等待网络时不会阻塞事件循环
普通文件对象、同步迭代库def + yieldFastAPI 可以在线程池中迭代
CPU 密集计算进程池或任务系统改成异步函数不会自动消除 CPU 阻塞

不要在异步生成器里直接调用耗时的同步 SDK。它会占住事件循环,使同一进程中的其他请求、心跳和取消处理一起停下来。如果供应商只提供同步流式接口,应使用同步生成器,让框架在线程池中迭代,或者把阻塞调用明确放到线程中执行。

3. 接入 LangChain

LangChain 的 astream() 会返回异步迭代器,可以直接接入 StreamingResponse:

langchain_stream.py
01
# src/fastapi/api/stream/langchain_stream.py
02
from collections.abc import AsyncIterator
03
from typing import Annotated
04
05
from fastapi import APIRouter, Depends
06
from fastapi.responses import StreamingResponse
07
08
from .common import ChatRequest, DemoChatService, get_chat_service
09
10
router = APIRouter()
11
12
# POST /api/v1/stream/langchain
13
@router.post("/langchain")
14
async def stream_chat(
15
payload: ChatRequest,
16
service: Annotated[DemoChatService, Depends(get_chat_service)],
17
) -> StreamingResponse:
18
async def generate() -> AsyncIterator[str]:
19
async for chunk in service.astream(payload.question):
20
yield chunk
21
22
return StreamingResponse(
23
generate(),
24
media_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 字段

LangChain 风格文本流POST /api/v1/stream/langchain · text/plain
本地服务:心跳检测中

你可以把这个案例与后面的 /chat 放在一起比较:两者使用同一个 DemoChatService.astream(),但 /langchain 只返回正文,/chat 还会返回开始、结束、错误和事件 ID 等协议信息

关于 langchain 更详细的使用,我们会在后续章节中深入学习

4. 对接 deepseek

前面的案例没有请求真实的 LLM 模型。这个案例,我们使用 deepseek 的 SSE 接口来对接流式响应

所有的配置,我们都配置在项目根目录的 .env 文件中:

.env
1
DEEPSEEK_MODEL=deepseek-chat
2
DEEPSEEK_BASE_URL=https://api.deepseek.com
3
DEEPSEEK_API_KEY=your-api-key

下面是案例和对应的代码实现

DeepSeek 多轮 SSE 对话POST /api/v1/stream/deepseek · 完整 messages 历史
本地服务:心跳检测中
deepseek_sse.py
common.py
001
# src/fastapi/api/stream/deepseek_sse.py
002
"""使用 .env 中的 DeepSeek 配置调用真实模型的 SSE 接口。"""
003
004
import asyncio
005
import logging
006
from collections.abc import AsyncIterator
007
from typing import Annotated
008
from uuid import uuid4
009
010
from fastapi import APIRouter, Depends, HTTPException
011
from fastapi.sse import EventSourceResponse, ServerSentEvent
012
from langchain_openai import ChatOpenAI
013
014
from app.core.config import Settings, get_settings
015
from app.model import create_model
016
017
from .common import ConversationRequest
018
019
logger = logging.getLogger(__name__)
020
router = APIRouter()
021
022
023
def get_deepseek_model(
024
settings: Annotated[Settings, Depends(get_settings)],
025
) -> ChatOpenAI:
026
"""在流启动前创建模型,缺少密钥时返回普通 HTTP 错误。"""
027
028
try:
029
return create_model(settings)
030
except RuntimeError as exc:
031
raise HTTPException(
032
status_code=503,
033
detail="DeepSeek 模型配置不完整",
034
) from exc
035
036
037
# POST /api/v1/stream/deepseek
038
@router.post(
039
"/deepseek",
040
response_class=EventSourceResponse,
041
summary="DeepSeek 流式聊天",
042
)
043
async def stream_deepseek(
044
payload: ConversationRequest,
045
model: Annotated[ChatOpenAI, Depends(get_deepseek_model)],
046
) -> AsyncIterator[ServerSentEvent]:
047
"""调用 DeepSeek,并把 LangChain 消息块转换成 SSE 生命周期事件。"""
048
049
message_id = f"msg_{uuid4().hex}"
050
sequence = 0
051
usage: dict[str, object] | None = None
052
053
yield ServerSentEvent(
054
event="message_start",
055
id=f"{message_id}:{sequence}",
056
retry=3_000,
057
data={
058
"message_id": message_id,
059
"model": model.model_name,
060
},
061
)
062
063
try:
064
async for chunk in model.astream(payload.to_langchain_messages()):
065
if chunk.usage_metadata is not None:
066
usage = dict(chunk.usage_metadata)
067
068
# ``text`` 会从字符串或标准文本内容块中提取可展示文本。
069
if not chunk.text:
070
continue
071
072
sequence += 1
073
yield ServerSentEvent(
074
event="token",
075
id=f"{message_id}:{sequence}",
076
data={"text": chunk.text},
077
)
078
except asyncio.CancelledError:
079
# 客户端断开时让取消信号继续传播,以便终止上游读取。
080
raise
081
except Exception:
082
logger.exception("DeepSeek streaming request failed")
083
sequence += 1
084
yield ServerSentEvent(
085
event="stream_error",
086
id=f"{message_id}:{sequence}",
087
data={
088
"code": "model_unavailable",
089
"message": "模型暂时不可用,请稍后重试",
090
},
091
)
092
return
093
094
sequence += 1
095
yield ServerSentEvent(
096
event="message_end",
097
id=f"{message_id}:{sequence}",
098
data={
099
"finish_reason": "stop",
100
"usage": usage,
101
},
102
)

最终请求路径是 POST /api/v1/stream/deepseek:

inspect-deepseek-sse.sh
01
curl -iN -X POST \
02
http://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”,等回答结束后再问“我刚才说我在学习 什么?”。第二轮请求不是只发送最后一个问题,而是发送类似下面的完整数组:

code.ts
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 功能更多,但也会增加连接管理、心跳、鉴权、重连和跨进程广播的复杂度,不应该成为默认选择。

正在验证登录状态
请稍候,验证完成后将继续显示文章内容