🧠 图解记忆: 流式输出是端到端管道,后端逐块转发、前端增量渲染,并处理取消、断线和背压。
💡 答案要点
Streaming = LLM边生成边返回token,首屏时间从5s→0.3s
为什么需要流式输出
非流式:用户等5秒 → 看到完整回答(体验差)
流式: 0.3秒开始显示第一个字 → 逐字显示 → 5秒显示完(体验好)
指标对比:
TTFT(首Token时间):5000ms → 300ms(-94%)
用户感知等待时间:5s → 0.3s(-94%)
用户满意度:↑ 显著提升后端实现(FastAPI SSE)
展开 Python 代码示例(64 行)
python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
from openai import OpenAI
import asyncio
import json
app = FastAPI()
client = OpenAI()
@app.post("/chat/stream")
async def chat_stream(request: ChatRequest):
"""流式聊天接口"""
async def generate():
try:
# 开启流式模式
stream = client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": request.message}],
stream=True, # 关键参数
max_tokens=2048,
)
# 逐chunk发送
for chunk in stream:
delta = chunk.choices[0].delta
if delta.content:
# SSE格式:data: {json}\n\n
data = json.dumps({
"type": "content",
"content": delta.content
}, ensure_ascii=False)
yield f"data: {data}\n\n"
# 工具调用流式处理
if delta.tool_calls:
for tc in delta.tool_calls:
data = json.dumps({
"type": "tool_call",
"tool": tc.function.name if tc.function.name else "",
"args_chunk": tc.function.arguments if tc.function.arguments else ""
})
yield f"data: {data}\n\n"
# 让出事件循环,避免阻塞
await asyncio.sleep(0)
# 发送结束信号
yield f"data: {json.dumps({'type': 'done'})}\n\n"
except Exception as e:
error_data = json.dumps({"type": "error", "message": str(e)})
yield f"data: {error_data}\n\n"
return StreamingResponse(
generate(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no", # 禁止Nginx缓冲
"Connection": "keep-alive",
}
)前端实现(EventSource / fetch)
展开 Javascript 代码示例(54 行)
javascript
// 方式1:EventSource(简单,仅支持GET)
const es = new EventSource('/chat/stream?message=你好');
es.onmessage = (event) => {
const data = JSON.parse(event.data);
if (data.type === 'content') {
// 追加到显示区域
document.getElementById('output').textContent += data.content;
} else if (data.type === 'done') {
es.close();
} else if (data.type === 'error') {
console.error(data.message);
es.close();
}
};
// 方式2:fetch + ReadableStream(支持POST,更灵活)
async function streamChat(message) {
const response = await fetch('/chat/stream', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({message})
});
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const {done, value} = await reader.read();
if (done) break;
buffer += decoder.decode(value, {stream: true});
// 按SSE格式解析
const lines = buffer.split('\n');
buffer = lines.pop(); // 保留不完整的行
for (const line of lines) {
if (line.startsWith('data: ')) {
const jsonStr = line.slice(6);
if (jsonStr === '[DONE]') return;
try {
const data = JSON.parse(jsonStr);
if (data.type === 'content') {
appendToOutput(data.content);
}
} catch (e) {}
}
}
}
}中间件处理(LangChain流式)
展开 Python 代码示例(43 行)
python
from langchain.callbacks.streaming_stdout import StreamingStdOutCallbackHandler
from langchain.callbacks.base import BaseCallbackHandler
class CustomStreamHandler(BaseCallbackHandler):
"""自定义流式回调"""
def __init__(self, queue: asyncio.Queue):
self.queue = queue
def on_llm_new_token(self, token: str, **kwargs):
"""每生成一个token触发"""
self.queue.put_nowait(token)
def on_llm_end(self, response, **kwargs):
"""生成结束"""
self.queue.put_nowait(None) # 发送结束信号
# 在FastAPI中使用
@app.post("/langchain/stream")
async def langchain_stream(request: ChatRequest):
queue = asyncio.Queue()
handler = CustomStreamHandler(queue)
async def generate():
# 在后台线程运行LangChain(避免阻塞事件循环)
import threading
def run_chain():
llm = ChatOpenAI(streaming=True, callbacks=[handler])
llm.invoke(request.message)
thread = threading.Thread(target=run_chain)
thread.start()
# 从队列读取token并发送
while True:
token = await queue.get()
if token is None:
yield f"data: {json.dumps({'type': 'done'})}\n\n"
break
yield f"data: {json.dumps({'type': 'content', 'content': token})}\n\n"
return StreamingResponse(generate(), media_type="text/event-stream")生产注意事项
python
# 1. Nginx配置(防止缓冲)
nginx_config = """
location /chat/stream {
proxy_pass http://backend;
proxy_buffering off; # 关键:禁用缓冲
proxy_cache off;
proxy_set_header Connection '';
proxy_http_version 1.1;
chunked_transfer_encoding on;
}
"""
# 2. 超时设置(流式响应时间长)
# 普通接口:30s超时
# 流式接口:600s超时(10分钟)
# 3. 断点续传(网络中断后继续)
class StreamResumer:
def __init__(self):
self.cache = {} # session_id → 已生成内容
def resume_stream(self, session_id: str, offset: int):
"""从offset位置继续输出"""
cached = self.cache.get(session_id, "")
if len(cached) > offset:
# 先发送已缓存的部分
return cached[offset:], True
return "", False面试话术:
"流式输出核心是Server-Sent Events(SSE)协议:后端yield逐个token,前端EventSource或fetch ReadableStream实时接收追加。TTFT从5s降到300ms,用户体验大幅提升。生产上3个关键点:1)Nginx必须关闭proxy_buffering,否则还是等全部生成才返回;2)流式接口超时设置要长,普通30s会被截断;3)用asyncio.sleep(0)让出事件循环,避免阻塞其他请求。工具调用也可以流式传递,让用户看到Agent正在思考的过程。"
