揭秘DeepSeek流式输出的核心技术:从Langchain到SSE的完整实现指南(Python+FastAPI实战)
1. 为什么流式输出是AI对话的“灵魂”?
你有没有在跟DeepSeek、ChatGPT这类AI聊天时,被那种“一个字一个字蹦出来”的实时回复体验所吸引?那种感觉就像对面真的有人在思考、在打字,而不是等上好几秒,突然给你扔过来一大段冰冷的文字。这种体验上的巨大差异,背后就是**流式输出(Streaming Output)**技术。
简单来说,传统AI对话是“打包发货”:你把问题扔给服务器,服务器吭哧吭哧让大模型生成完整答案,然后一次性把整个“包裹”发回给你。这个过程,用户面对的是一个空白的输入框,在等待中可能会焦虑,甚至怀疑是不是网络断了。而流式输出则是“流水线发货”:模型每生成一个词、一个片段,就立刻通过管道“流”到你的界面上。你看到的是动态的、渐进的思考过程,交互感、实时感和参与感都直接拉满。
在我做过的多个AI对话和RAG(检索增强生成)项目里,一旦接入了流式输出,用户的平均对话轮次和停留时长都会有明显提升。这不仅仅是“酷”,而是实实在在提升了产品的可用性和用户粘性。对于需要低延迟交互的场景,比如智能客服、实时翻译、代码辅助编程,流式输出几乎是必选项。
那么,这种“打字机”效果是怎么实现的呢?核心链路其实可以概括为三步:后端异步生成数据块 -> 通过SSE协议流式传输 -> 前端实时接收并渲染。听起来简单,但里面有不少坑,比如怎么处理Langchain里复杂的Chain输出,怎么确保SSE连接稳定,前端又该怎么优雅地拼接这些碎片化的文本。别担心,接下来我会结合一个真实的Python + FastAPI + Vue项目,带你从原理到代码,一步步拆解,保证你能自己动手实现一套。
2. 核心基石:理解Python的生成器与yield
在深入代码之前,我们必须先搞懂两个最核心的Python概念:**生成器(Generator)**和 yield关键字。这是理解后端如何“流”出数据的关键。
2.1 迭代器:一切可遍历对象的基石
想象一下,你有一个装满水果的篮子(列表)。你伸手进去,一次拿一个水果出来,直到拿空。这个“伸手去拿”的动作,背后就是迭代器(Iterator)在起作用。列表、元组、字典、字符串都是可迭代对象(Iterable),但不是迭代器本身。你可以通过iter()函数把一个可迭代对象变成迭代器。
迭代器有一个核心方法:next()。每次调用next(iterator),它就返回下一个元素,没有更多元素时,就会抛出StopIteration异常。for循环的本质,就是自动帮你创建迭代器并不断调用next()。
my_list = [‘苹果‘, ‘香蕉‘, ‘橙子‘]
my_iterator = iter(my_list) # 把列表变成迭代器
print(next(my_iterator)) # 输出:苹果
print(next(my_iterator)) # 输出:香蕉
print(next(my_iterator)) # 输出:橙子
print(next(my_iterator)) # 抛出 StopIteration
2.2 生成器:用函数“懒”生成序列的利器
现在,假设我们不想提前准备好整个水果篮子,而是希望根据需求,现场“变”出水果。这就是生成器函数的用武之地。
生成器函数看起来和普通函数一样,但它的魔法在于使用yield而不是return。当函数执行到yield语句时,它会“暂停”,并把yield后面的值返回给调用者。神奇的是,函数的所有局部变量状态都会被冻结保存。下次再请求下一个值时,函数会从上次yield之后的地方继续执行。
def fruit_generator():
print("开始生成苹果")
yield "苹果"
print("继续生成香蕉")
yield "香蕉"
print("最后生成橙子")
yield "橙子"
print("生成器结束")
# 调用生成器函数,不会立即执行函数体,而是返回一个生成器对象
gen = fruit_generator()
# 第一次调用next,执行到第一个yield,返回“苹果”,并暂停
fruit = next(gen)
print(f"拿到了:{fruit}") # 输出:开始生成苹果 \n 拿到了:苹果
# 第二次调用next,从上次暂停处继续,执行到第二个yield,返回“香蕉”
fruit = next(gen)
print(f"拿到了:{fruit}") # 输出:继续生成香蕉 \n 拿到了:香蕉
你可以用for循环来优雅地遍历生成器:
for fruit in fruit_generator():
print(f"循环中拿到:{fruit}")
这个循环会依次打印出“开始生成苹果”、“循环中拿到:苹果”、“继续生成香蕉”…… 生成器在for循环结束时自动处理了StopIteration。
yield vs return 核心区别:
return:彻底结束函数,返回一个值,函数状态被销毁。yield:临时“交出”一个值,函数状态被完整保存,等待下次“唤醒”。
在我们的流式输出场景里,yield就是那个“流水线传送带”。大模型每生成一个词(chunk),我们就通过yield把它“放”到传送带上,立刻送出去,而不是等所有词都造好了再一起打包。
3. 后端引擎:Langchain的astream()如何驱动数据流
知道了yield怎么“送”数据,接下来要看数据从哪里来。在基于Langchain构建的AI应用中,答案就是 Runnable.astream() 方法。
3.1 从“invoke”到“astream”的转变
在非流式场景下,我们调用Langchain的Chain(链)通常是这样的:
# 传统一次性调用
result = my_chain.invoke({"input": "你好,世界"})
print(result["output"]) # 等待完整结果生成后,一次性打印
invoke方法是同步的,它会阻塞直到整个Chain执行完毕,返回最终结果。
而流式调用的核心,就是把invoke换成astream(异步流):
# 流式调用
async for chunk in my_chain.astream({"input": "你好,世界"}):
print(chunk, end="|", flush=True) # 收到一个chunk就打印一个
astream()返回的是一个异步迭代器(AsyncIterator)。async for循环会异步地、逐个地从迭代器中取出数据块(chunk)。大模型每生成一个token(或一小段文本),这个chunk就会立刻通过迭代器出来,而不是等到全部生成完毕。
3.2 实战:处理复杂的Chain输出结构
在实际项目中,我们的Chain结构可能很复杂,比如一个典型的RAG链:先检索文档,再把检索结果和问题一起喂给LLM。astream()出来的chunk可能包含多种类型的数据,我们需要从中提取出我们想要的文本内容。
假设我们的链最终输出一个字典,包含answer和context。但流式过程中,LLM本身输出的可能是AIMessageChunk对象。下面是一个我在项目中使用的、健壮性很强的chunk处理函数片段:
from langchain_core.messages import BaseMessage
async def stream_chat(chain, user_input):
"""处理流式输出的核心生成器函数"""
async for chunk in chain.astream({"input": user_input}):
content_piece = ""
# 情况1:chunk是Langchain的消息对象(如AIMessageChunk)
if isinstance(chunk, BaseMessage):
if hasattr(chunk, 'content'):
content_piece = chunk.content
# 情况2:chunk是字典,比如我们自定义的RAG链输出
elif isinstance(chunk, dict):
# 假设我们的链输出格式为 {'answer': '...', 'context': [...]}
if 'answer' in chunk:
answer_part = chunk['answer']
# answer部分可能已经是字符串,也可能是消息对象
if isinstance(answer_part, str):
content_piece = answer_part
elif isinstance(answer_part, BaseMessage) and hasattr(answer_part, 'content'):
content_piece = answer_part.content
# 你也可以处理其他键,比如流式返回检索到的上下文片段
elif 'context' in chunk:
# 这里可以yield一个上下文提示,给前端显示
yield {"type": "context", "data": f"检索到相关文档..."}
continue # 跳过本次循环,不输出文本内容
# 情况3:其他未知类型,可以记录日志或跳过
else:
logging.debug(f"收到未知类型的chunk: {type(chunk)}")
continue
# 只有提取到有效文本内容时才yield
if content_piece:
yield {"type": "chunk", "data": content_piece}
这段代码的关键在于类型判断和容错。因为Langchain的流式输出可能因Chain的构成不同而变化,我们必须确保无论收到什么,都能安全地提取出文本。yield出去的也不是纯文本,而是一个结构化的字典,用type字段来区分消息类型(如普通文本块chunk、上下文信息context、错误error),这样前端处理起来就非常清晰。
4. 传输协议:用FastAPI的StreamingResponse封装SSE
数据块通过yield生产出来了,怎么高效、标准地“流”到浏览器呢?这里就要请出Server-Sent Events (SSE) 协议和FastAPI的 StreamingResponse。
4.1 为什么是SSE,而不是WebSocket?
你可能听说过WebSocket,它支持全双工通信(客户端和服务器可以互相随时发消息)。但对于AI对话这种典型的“客户端提问,服务器流式推送回答”的场景,SSE是更简单、更轻量的选择。
- SSE是单向的:服务器向客户端推送。这完美匹配了AI回答的推送模式。
- 基于HTTP:SSE就是普通的HTTP连接,只是
Content-Type设置为text/event-stream。这意味着它更容易处理身份验证、CORS,也更利于在现有HTTP基础设施上部署和调试。 - 自动重连:浏览器端的EventSource API内置了断线重连机制。
- 简单易懂:数据格式就是简单的
data: <内容>\n\n。
4.2 构建SSE响应生成器
我们的目标是把前面stream_chat函数yield出的结构化字典,转换成SSE协议要求的格式,并通过FastAPI的路由返回。下面是一个完整的StreamingResponse应用示例:
from fastapi import FastAPI, APIRouter
from fastapi.responses import StreamingResponse
import json
import logging
app = FastAPI()
chat_router = APIRouter(prefix="/chat", tags=["Chat"])
@chat_router.post("/stream")
async def chat_stream_endpoint(question: str):
"""流式聊天接口"""
# 1. 这里可以加入你的业务逻辑,比如验证用户、获取配置等
# 2. 关键:返回StreamingResponse,并指定媒体类型为 text/event-stream
return StreamingResponse(
stream_response_generator(question), # 传入异步生成器
media_type="text/event-stream" # 必须指定!
)
async def stream_response_generator(question: str):
"""SSE响应生成器:将业务数据转换为SSE格式"""
logging.info(f"开始为问题‘{question}‘生成流式响应")
try:
# 模拟或调用你真正的流式处理链,这里用异步生成器模拟
# 假设 get_streaming_chain() 返回一个类似3.2节中的异步生成器
async for chunk_dict in get_streaming_chain(question):
# chunk_dict 格式如: {"type": "chunk", "data": "你好"}
# 转换为SSE格式: `data: <json_data>\n\n`
# 注意:必须用双换行符\n\n标识一个事件的结束
event_data = f"data: {json.dumps(chunk_dict)}\n\n"
yield event_data
logging.debug(f"已发送chunk: {chunk_dict}")
except Exception as e:
logging.error(f"流式响应生成过程中发生错误: {e}", exc_info=True)
# 即使出错,也要以SSE格式通知前端
error_payload = json.dumps({
"type": "error",
"data": f"流处理发生严重错误: {str(e)}"
})
yield f"data: {error_payload}\n\n"
finally:
logging.info(f"流式响应结束")
代码解读:
StreamingResponse:这是FastAPI提供的专门用于流式响应的类。它接受一个异步生成器作为content参数。media_type="text/event-stream":这是灵魂!告诉浏览器这是一个SSE流,浏览器会以特定方式处理。stream_response_generator:这是核心的生成器函数。它async for循环遍历业务层产生的数据块(chunk_dict)。- SSE格式封装:对于每个
chunk_dict,我们用json.dumps将其转换为JSON字符串,然后按照SSE格式拼接成data: <json_string>\n\n。这个\n\n(两个换行符)是SSE协议规定的事件边界,绝对不能少。 - 错误处理:即使在生成器内部发生异常,我们也应该捕获它,并构造一个类型为
error的SSE事件发送给前端,让前端能优雅地告知用户,而不是连接突然中断。
这样,一个标准的、支持错误处理的SSE流后端接口就搭建好了。你可以用curl命令测试一下:
curl -N -X POST http://你的服务器地址/chat/stream \
-H "Content-Type: application/json" \
-d '{"question": "你好"}'
你会看到数据以data: {...}\n\n的形式一行行实时返回。
5. 前端接收:使用fetch-event-source实现实时渲染
后端的数据流已经准备就绪,前端如何接住并实时显示呢?由于原生的EventSource API不支持发送POST请求体(我们的问题数据需要放在body里),我们选择微软开源的 @microsoft/fetch-event-source 这个库,它基于fetch API,功能更强大。
5.1 安装与基础配置
首先,在你的Vue(或React)项目中安装它:
pnpm add @microsoft/fetch-event-source
# 或
npm install @microsoft/fetch-event-source
# 或
yarn add @microsoft/fetch-event-source
5.2 核心实现:连接、接收与状态管理
下面是一个在Vue 3组合式API中实现的完整sendMessage函数,它包含了流式请求的发起、数据处理、状态控制和错误处理。
<script setup>
import { ref } from 'vue'
import { fetchEventSource } from '@microsoft/fetch-event-source'
import { ElMessage, ElNotification } from 'element-plus' // 示例UI库,可按需替换
const inputMessage = ref('')
const messages = ref([]) // 消息列表
const isStreaming = ref(false) // 是否正在流式传输
const abortController = ref(null) // 用于取消请求
const API_BASE = import.meta.env.VITE_API_BASE_URL // 你的API基础地址
const sendMessage = async () => {
// 1. 前置检查
if (!inputMessage.value.trim()) return
if (isStreaming.value) {
ElMessage.warning('正在等待AI回复,请稍候')
return
}
// 2. 取消上一个可能存在的请求
if (abortController.value) {
abortController.value.abort()
}
abortController.value = new AbortController()
// 3. 准备用户消息和AI占位符
const userMessage = {
id: Date.now(),
role: 'user',
content: inputMessage.value.trim()
}
messages.value.push(userMessage)
const aiMessageId = Date.now() + 1
const aiMessage = {
id: aiMessageId,
role: 'assistant',
content: '' // 初始为空,后续用chunk拼接
}
messages.value.push(aiMessage)
// 4. 清空输入框,更新状态
const currentQuestion = inputMessage.value
inputMessage.value = ''
isStreaming.value = true
// 5. 准备请求体
const payload = {
question: currentQuestion,
// 其他参数,如session_id, model等
session_id: getCurrentSessionId(),
model: 'deepseek-chat'
}
let accumulatedContent = '' // 用于累积AI回复内容
try {
await fetchEventSource(`${API_BASE}/chat/stream`, {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Accept': 'text/event-stream' // 关键头,声明接受SSE流
},
body: JSON.stringify(payload),
signal: abortController.value.signal, // 关联取消信号
// 5.1 连接建立时
onopen: async (response) => {
if (response.ok && response.headers.get('content-type')?.includes('text/event-stream')) {
console.log('SSE连接已建立')
} else {
// 连接失败,抛出错误会被onerror捕获
throw new Error(`连接失败: ${response.status} ${response.statusText}`)
}
},
// 5.2 收到消息时(核心逻辑)
onmessage: (event) => {
// event.data 就是后端发来的 `data: {...}\n\n` 中 `{...}` 的部分
try {
const parsedData = JSON.parse(event.data)
const aiMsgIndex = messages.value.findIndex(msg => msg.id === aiMessageId)
if (aiMsgIndex === -1) return // 消息占位符意外丢失
switch (parsedData.type) {
case 'context':
// 处理上下文信息,例如用通知显示
ElNotification({
title: '检索到上下文',
message: parsedData.data,
type: 'info'
})
break
case 'chunk':
// 核心:拼接文本块
accumulatedContent += parsedData.data
// 更新Vue响应式数据,触发视图更新
messages.value[aiMsgIndex].content = accumulatedContent
// 可选:自动滚动到底部
scrollToBottom()
break
case 'error':
// 处理后端报告的流内错误
console.error('流处理错误:', parsedData.data)
messages.value[aiMsgIndex].content += `\n\n**系统错误:** ${parsedData.data}`
ElMessage.error(`AI回复出错: ${parsedData.data}`)
// 发生错误,主动中止连接
if (abortController.value) {
abortController.value.abort()
}
break
}
} catch (e) {
console.error('解析SSE数据失败:', e, '原始数据:', event.data)
ElMessage.error('收到无效的数据格式')
if (abortController.value) abortController.value.abort()
}
},
// 5.3 连接关闭时
onclose: () => {
console.log('SSE连接已关闭')
isStreaming.value = false
abortController.value = null
scrollToBottom()
},
// 5.4 发生错误时
onerror: (err) => {
console.error('SSE错误:', err)
isStreaming.value = false
// 特殊处理用户手动取消的情况
if (err.name === 'AbortError') {
console.log('流已被用户中止')
// 如果AI消息还是空的,可以移除占位符
const aiMsgIndex = messages.value.findIndex(msg => msg.id === aiMessageId)
if (aiMsgIndex !== -1 && !messages.value[aiMsgIndex].content) {
messages.value.splice(aiMsgIndex, 1)
}
return // 对于AbortError,不显示错误提示,直接返回
}
// 处理其他错误(网络错误、服务器错误等)
ElMessage.error(`连接错误: ${err.message || '未知错误'}`)
const aiMsgIndex = messages.value.findIndex(msg => msg.id === aiMessageId)
if (aiMsgIndex !== -1) {
messages.value[aiMsgIndex].content += `\n\n**网络错误:** ${err.message || '未知错误'}`
}
// 对于非AbortError,需要抛出,以阻止fetchEventSource的默认重试逻辑(如果需要的话)
throw err
}
})
} catch (err) {
// 捕获在初始化请求或onerror中抛出的、未被处理的错误
if (err.name !== 'AbortError') {
console.error('发起SSE请求失败:', err)
ElMessage.error('请求发送失败')
isStreaming.value = false
abortController.value = null
}
}
}
// 一个停止流式传输的函数,供“停止生成”按钮调用
const stopStreaming = () => {
if (abortController.value) {
abortController.value.abort()
isStreaming.value = false
ElMessage.info('已停止生成')
}
}
</script>
<template>
<div class="chat-container">
<div class="messages">
<div v-for="msg in messages" :key="msg.id" :class="`message ${msg.role}`">
{{ msg.content }}
</div>
</div>
<div class="input-area">
<input v-model="inputMessage" @keyup.enter="sendMessage" :disabled="isStreaming" />
<button @click="sendMessage" :disabled="isStreaming">发送</button>
<button @click="stopStreaming" v-if="isStreaming">停止生成</button>
</div>
</div>
</template>
5.3 关键点解析与避坑指南
Accept: 'text/event-stream':这个请求头至关重要,它告诉服务器客户端期望SSE流。有些后端框架会根据这个头来调整行为。signal:关联AbortController的signal,这是实现“停止生成”功能的关键。用户点击停止时,调用abortController.abort(),会触发AbortError,从而优雅地中断连接。onmessage:这里是处理数据的核心。我们根据后端返回的type字段来区分处理逻辑。对于chunk类型,我们将data字段的内容累加到当前AI消息上,并立即更新Vue的响应式数据,界面就会实时渲染。- 错误处理分层:
onopen中检查响应是否成功。onmessage中处理数据解析错误和后端返回的业务错误(type: 'error')。onerror中处理网络错误、连接错误等。要特别注意区分AbortError(用户取消)和其他错误。
- 状态管理:
isStreaming这个状态非常有用,可以用于控制按钮的禁用状态、输入框的只读状态,防止用户在流式传输过程中重复提交。 - 性能与体验:在
onmessage中频繁更新DOM(通过Vue响应式系统)可能会影响性能。如果响应速度极快(比如本地模型),可以考虑使用防抖(debounce)来合并短时间内的多次更新,减少UI重绘次数。但对于网络请求,通常直接更新即可。
6. 进阶优化与实战踩坑经验
把基础链路跑通只是第一步。在实际生产环境中,你会遇到各种边界情况和性能问题。这里分享几个我踩过的坑和优化方案。
6.1 连接稳定性与超时处理
SSE连接是长连接,可能会因为网络波动、代理服务器、负载均衡器超时等原因中断。fetch-event-source库内置了重试逻辑,但其默认策略可能不满足需求。
优化方案:自定义重试策略。你可以在fetchEventSource的配置中提供一个onerror回调,并在其中决定是否重试。
await fetchEventSource(url, {
// ... 其他配置
onerror: (err) => {
// 如果是网络错误或服务器5xx错误,可以重试
if (err instanceof TypeError || // 网络错误
(err.message && err.message.includes('500'))) {
console.log('遇到可重试错误,将在2秒后重试...')
// 返回一个数字,表示重试延迟(毫秒)
// 返回undefined或不返回任何值,则停止重试
return 2000
}
// 其他错误(如4xx客户端错误、AbortError)则直接抛出,停止重试
throw err
}
})
同时,后端(Nginx等代理)也需要配置合适的超时时间(如proxy_read_timeout 300s;),避免过早关闭连接。
6.2 流式输出中的上下文与思考过程展示
在复杂的RAG或Agent应用中,你可能不仅想展示最终的答案文本,还想让用户看到模型的“思考过程”,比如“我正在检索文档...”、“我找到了三条相关信息...”。这可以通过我们之前定义的type: 'context'事件来实现。
在后端的stream_response_generator中,你可以在流开始前、检索到文档时、调用工具时,yield不同类型的事件:
async def stream_response_generator(question):
# 1. 通知前端开始思考
yield {"type": "status", "data": "思考中..."}
# 2. 执行检索
docs = retrieve_documents(question)
yield {"type": "context", "data": f"已检索到 {len(docs)} 条相关文档"}
# 3. 开始流式生成答案
async for chunk in llm_chain.astream({"question": question, "context": docs}):
if has_valid_content(chunk):
yield {"type": "chunk", "data": extract_content(chunk)}
# 4. 生成结束
yield {"type": "status", "data": "生成完成"}
前端在onmessage中根据不同的type更新不同的UI组件,比如在消息上方显示一个状态栏,或者用特殊样式展示检索到的文档摘要,极大地增强了交互的透明度和用户体验。
6.3 性能监控与调试技巧
流式应用调试起来比普通API麻烦一些,因为数据是持续的。这里有几个实用技巧:
- 后端日志:在
stream_response_generator的每个yield前后打上日志,记录时间戳和chunk内容的前几个字符。这能帮你确认数据是否按预期产生和发送。 - 前端网络面板:在Chrome DevTools的Network标签页,找到你的
/chat/stream请求,查看“Response”标签页。如果连接正确,你会看到数据在实时追加。这是检查SSE流是否正常工作的最直接方法。 - 模拟慢速网络:在DevTools的Network条件中,选择“Slow 3G”来测试前端在弱网环境下拼接和渲染数据的能力,以及超时重试逻辑是否健壮。
- 压力测试:使用工具模拟大量并发流式请求,观察服务器的内存和CPU使用情况。确保你的异步生成器能正确释放资源,避免内存泄漏。
6.4 结合LangChain新特性的实践
LangChain社区活跃,不断有新特性推出。例如,LangGraph更适合构建有状态的、多步骤的Agent应用。其流式输出也基于类似的异步生成器模式。
在LangGraph中,你可以通过监听特定节点的输出来实现更精细的流式控制:
from langgraph.graph import StateGraph, MessagesState
from langgraph.checkpoint import MemorySaver
graph_builder = StateGraph(MessagesState)
# ... 构建你的图
graph = graph_builder.compile()
# 流式运行
async for event in graph.astream({"messages": [("user", "你好")]}):
# event 是一个元组 (node_name, node_output)
node_name, output = event
if node_name == "llm_node":
# 处理LLM产生的流式chunk
if hasattr(output, 'content'):
for chunk in output.content:
yield {"type": "chunk", "data": chunk}
elif node_name == "tool_node":
# 通知前端工具调用
yield {"type": "action", "data": f"调用了工具: {output}"}
这让你能将Agent复杂的推理步骤也“流式”地展示给用户,体验更上一层楼。
流式输出从技术上看,是生成器、异步迭代、SSE协议和前端事件处理的结合。但从产品角度看,它关乎的是用户的耐心、注意力和对智能的感知。实现它并不需要多高深的技术,但需要对整个数据流有清晰的认识和细致的处理。当你看到自己实现的聊天界面里,文字像有了生命一样逐字浮现时,那种成就感就是对我们开发者最好的回报。
更多推荐


所有评论(0)