公司动态
基于FastAPI与SSE实现Web流式输出:构建极简智能体后端
1. 项目缘起从“等待”到“流淌”的体验革命做后端开发久了你肯定遇到过这种场景用户提交一个需要复杂处理的请求比如生成一篇长文、分析一份文档或者调用一个大语言模型。然后呢前端页面就卡在那里转着圈圈用户盯着进度条心里默数着秒数。后端在吭哧吭哧地处理直到最后“哐当”一下把整个结果一股脑地吐给前端。这个过程我们称之为“阻塞式响应”或“同步响应”。对于用户来说体验是割裂的对于服务器来说连接被长时间占用资源利用率也谈不上高效。“流式输出”就是为了解决这个问题而生的。它的核心思想很简单别等所有事情都做完了再给用户看而是像拧开水龙头一样让数据一点一点地“流”出来。用户能立刻看到部分结果感知到进度体验的流畅度和即时反馈感会得到质的提升。尤其是在AI应用爆发的今天让大模型一个字一个字地“吐出”回答而不是等上十几秒才看到完整段落已经成为高质量交互的标配。这次我们要聊的就是在Web环境下实现这种“流淌”体验。标题里的“极简智能体”点明了场景——一个轻量级的、具备一定自主处理能力的服务端程序而“Web流式输出”则是我们实现与前端流畅对话的技术手段。这不仅仅是技术实现更是一种产品思维和用户体验的优化。我们将聚焦于如何用最简洁、最现代的技术栈构建一个支持流式输出的Web后端并确保它稳定、高效。2. 技术选型为什么是SSE与FastAPI实现Web流式输出主流有几种技术路径WebSocket、Server-Sent Events (SSE) 和普通的HTTP分块传输。我们的目标是“极简”这意味着我们需要在功能、复杂度和浏览器兼容性之间找到一个优雅的平衡点。WebSocket功能最强大支持全双工通信客户端和服务器可以随时互相发送消息适合聊天室、实时游戏等场景。但它也最重需要单独的协议ws://或wss://服务端和客户端都要处理连接建立、维护和断开引入了额外的复杂度。HTTP分块传输 (Chunked Transfer Encoding)是HTTP/1.1的标准特性服务器可以把响应体分成多个“块”依次发送。但它本质上还是一个“请求-响应”模型一次请求对应一个流式响应响应结束连接就关闭。它缺乏标准的事件机制前端需要自己解析流处理起来不够直观。Server-Sent Events (SSE)恰恰是介于两者之间的“甜点”。它基于普通的HTTP协议使用长连接但只支持服务器向客户端的单向数据推送这正是我们流式输出的主要方向。它的协议非常简单数据格式有标准规范并且绝大多数现代浏览器都原生支持。对于“服务器推送文本流”这个特定场景SSE在简洁性和功能性上达到了完美统一。因此SSE是我们的不二之选。它极简无需额外协议它标准前端有EventSource对象原生支持它够用完美契合智能体逐字吐出结果的需求。后端框架我们选择FastAPI。原因同样围绕着“极简”和“现代”异步原生支持流式输出本质上是I/O密集型操作异步编程可以让我们在等待AI模型返回下一个词时轻松处理其他请求极大提升并发能力。FastAPI基于Starlette对异步的支持是刻在基因里的。简洁直观使用Python的async/await语法和Pydantic模型定义API清晰得如同写说明文。实现一个流式端点代码量很少。性能优异基于ASGI标准性能表现非常出色足以应对高并发流式请求。生态友好与httpx,aiohttp等异步HTTP客户端以及openai等AI库的异步版本搭配使用浑然天成。所以我们的技术栈锚定为FastAPI (后端) SSE (协议) AsyncOpenAI (AI调用)。这个组合能以最小的开发成本构建出体验优秀的流式智能体接口。3. 核心实现从零搭建FastAPI流式端点理论说清楚了我们直接上代码。一个最核心的、支持SSE流式输出的FastAPI端点是如何构建的。3.1 项目初始化与依赖安装首先创建一个新的项目目录并初始化虚拟环境这是保持环境清洁的好习惯。mkdir simple-agent-streaming cd simple-agent-streaming python -m venv venv # Windows 使用 venv\Scripts\activate source venv/bin/activate接着安装核心依赖。我们需要的包不多体现了“极简”。pip install fastapi uvicorn httpx sse-starlette openaifastapi: Web框架本体。uvicorn: ASGI服务器用于运行FastAPI应用。httpx: 异步HTTP客户端用于异步调用外部API虽然本例直接使用OpenAI库但它是优秀的备选。sse-starlette: 一个让FastAPI/Starlette轻松生成SSE响应的库。虽然FastAPI可以直接用StreamingResponse手动构造SSE但这个库让代码更优雅。openai: OpenAI官方库我们使用其异步客户端。3.2 构建SSE流式响应端点我们来创建一个main.py文件这是应用的核心。from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import asyncio import json from openai import AsyncOpenAI import os from sse_starlette.sse import EventSourceResponse app FastAPI(title极简流式智能体) # 初始化异步OpenAI客户端建议从环境变量读取API Key client AsyncOpenAI(api_keyos.getenv(OPENAI_API_KEY)) app.get(/stream-simple) async def stream_simple(request: Request): 一个最简单的手动构造SSE流的示例。 每秒发送一个数字。 async def event_generator(): for i in range(1, 6): # 检查客户端是否还连接着如果断开则停止生成 if await request.is_disconnected(): break # SSE数据格式data: 内容\n\n # 内容必须是字符串如果要发送JSON需要先序列化 yield fdata: 当前数字是 {i}\n\n await asyncio.sleep(1) # 模拟耗时操作 yield data: [DONE]\n\n # 发送一个结束标记 return StreamingResponse( event_generator(), media_typetext/event-stream, headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no # 针对Nginx代理的重要设置 } )这个/stream-simple端点展示了最原始的SSE流构造。关键点在于我们定义了一个异步生成器函数event_generator它yield出的每一个字符串都必须遵循SSE格式以data:开头以两个换行符\n\n结尾。StreamingResponse接收这个生成器并设置media_typetext/event-stream这是告诉浏览器这是一个SSE流。头部信息很重要Cache-Control: no-cache和Connection: keep-alive是SSE的常见设置。X-Accel-Buffering: no尤其关键如果你用了Nginx做反向代理这个设置可以防止Nginx对响应进行缓冲否则前端可能无法实时收到数据块。3.3 集成AsyncOpenAI实现真正的AI流上面的例子是模拟数据。现在我们集成真正的AI大模型实现智能体的流式对话。创建一个新的端点/chat/stream。from pydantic import BaseModel from typing import AsyncGenerator import logging logging.basicConfig(levellogging.INFO) logger logging.getLogger(__name__) class ChatMessage(BaseModel): role: str content: str class ChatRequest(BaseModel): messages: list[ChatMessage] model: str gpt-3.5-turbo # 默认模型 app.post(/chat/stream) async def chat_stream(chat_request: ChatRequest, request: Request): 接收聊天消息流式返回AI的回复。 async def stream_openai_response() - AsyncGenerator[str, None]: try: # 调用OpenAI的异步流式API stream await client.chat.completions.create( modelchat_request.model, messages[msg.dict() for msg in chat_request.messages], streamTrue, # 关键参数开启流式输出 temperature0.7, max_tokens500, ) async for chunk in stream: if await request.is_disconnected(): logger.info(客户端断开连接停止流式响应。) break # 提取流中的内容 if chunk.choices[0].delta.content is not None: content chunk.choices[0].delta.content # 将内容封装为SSE格式 # 这里我们发送JSON包含内容和可能的其他元数据如是否结束 data json.dumps({content: content, finish_reason: None}) yield fdata: {data}\n\n except Exception as e: logger.error(f流式处理发生错误: {e}) # 发生错误时也发送一个错误事件给前端 error_data json.dumps({error: str(e), content: None, finish_reason: error}) yield fdata: {error_data}\n\n finally: # 流正常结束发送结束标志 done_data json.dumps({content: , finish_reason: done}) yield fdata: {done_data}\n\n return EventSourceResponse( stream_openai_response(), headers{ Cache-Control: no-cache, Connection: keep-alive, X-Accel-Buffering: no } )这个实现是核心中的核心有几个细节需要深究streamTrue这是OpenAI API调用中最关键的参数。没有它你会等到所有token都生成完毕才收到一个完整的响应对象。有了它client.chat.completions.create返回的是一个异步生成器每次yield一个包含部分结果的chunk。数据提取每个chunk的结构需要仔细处理。chunk.choices[0].delta包含了本次流增量delta的内容。delta可能包含content文本内容、role角色等。我们主要关心content。重要提示delta里的content可能为None例如流刚开始时可能只返回role所以必须做判空处理。连接状态检查await request.is_disconnected()是一个非常重要的优化。如果用户在流式输出过程中关闭了浏览器标签页这个检查能让我们及时感知并停止后续的AI生成和网络传输避免浪费宝贵的AI Token和服务器资源。错误处理与结束信号流式处理中网络错误、API限额错误都可能发生。我们用try...except包裹并在异常时也向客户端发送一个格式化的错误事件。无论成功失败最后都发送一个明确的结束事件finish_reason: done或error让前端知道流已关闭可以更新UI状态比如把“正在输入”的动画去掉。使用EventSourceResponse这里我们换用了sse_starlette提供的EventSourceResponse。它内部帮我们处理了SSE格式的封装和心跳维持默认每5秒发送一个: ping注释行以保持连接活性比手动StreamingResponse更省心。3.4 运行与测试应用写好了让我们运行它。uvicorn main:app --reload --host 0.0.0.0 --port 8000打开浏览器访问http://localhost:8000/docs你会看到自动生成的交互式API文档。你可以直接在Docs里测试/chat/stream接口但Swagger UI对SSE的支持有限通常不会实时显示流内容。为了真正测试流式效果我们需要一个简单的前端页面。4. 前端对接使用EventSource消费流式接口流式接口的后端已经就绪前端如何接收并展示这个“数据流”呢HTML5原生提供了EventSource对象它就是为消费SSE流而生的。创建一个templates/index.html文件确保项目根目录有templates文件夹。!DOCTYPE html html head title极简智能体 - 流式对话测试/title style body { font-family: sans-serif; max-width: 800px; margin: 2em auto; padding: 1em; } #chatBox { border: 1px solid #ccc; height: 400px; overflow-y: auto; padding: 1em; margin-bottom: 1em; } .message { margin-bottom: 0.8em; } .user { text-align: right; color: #0066cc; } .assistant { text-align: left; color: #333; } #inputArea { display: flex; } #userInput { flex-grow: 1; padding: 0.5em; } #sendBtn { padding: 0.5em 2em; } #status { font-style: italic; color: #666; margin-top: 0.5em; } /style /head body h2与极简智能体对话 (流式模式)/h2 div idchatBox/div div idinputArea input typetext iduserInput placeholder输入你的问题... / button idsendBtn发送/button /div div idstatus准备就绪。/div script const chatBox document.getElementById(chatBox); const userInput document.getElementById(userInput); const sendBtn document.getElementById(sendBtn); const statusDiv document.getElementById(status); let currentEventSource null; let assistantMessageDiv null; function appendMessage(role, content) { const div document.createElement(div); div.className message ${role}; div.textContent ${role}: ${content}; chatBox.appendChild(div); chatBox.scrollTop chatBox.scrollHeight; // 自动滚动到底部 return div; } function updateStatus(msg) { statusDiv.textContent msg; } sendBtn.addEventListener(click, async () { const question userInput.value.trim(); if (!question) return; // 禁用输入和按钮防止重复发送 userInput.disabled true; sendBtn.disabled true; userInput.value ; // 显示用户消息 appendMessage(user, question); // 创建并显示一个空的助手消息容器用于流式追加内容 assistantMessageDiv appendMessage(assistant, ); assistantMessageDiv.textContent assistant: ; // 重置只保留前缀 // 如果存在之前的连接先关闭 if (currentEventSource) { currentEventSource.close(); } updateStatus(正在思考...); // 构建请求体 const requestBody { messages: [{ role: user, content: question }], model: gpt-3.5-turbo }; // 关键使用 EventSource 连接 SSE 端点 // 注意EventSource 只支持 GET 请求且不能自定义Header如Authorization。 // 对于需要POST或认证的接口这不是最佳选择。更通用的方案是使用 fetch。 // 这里为了演示SSE原理我们假设有一个GET端点实际项目中慎用。 // 更推荐使用 Fetch API 来消费流见下方补充说明。 // 模拟一个GET端点需要在后端实现 /chat/stream-get?qxxx // const eventSource new EventSource(/chat/stream-get?q${encodeURIComponent(question)}); // 更通用的方案使用 Fetch API 处理流 try { const response await fetch(/chat/stream, { method: POST, headers: { Content-Type: application/json, }, body: JSON.stringify(requestBody) }); if (!response.ok || !response.body) { throw new Error(HTTP error! status: ${response.status}); } const reader response.body.getReader(); const decoder new TextDecoder(); let accumulatedText ; updateStatus(正在接收流式回复...); while (true) { const { done, value } await reader.read(); if (done) { updateStatus(回复完成。); break; } // 解码 chunk const chunk decoder.decode(value, { stream: true }); // SSE流中每个事件由 \n\n 分隔但通过fetch读取时可能不会按完整事件到达 // 我们需要一个简单的缓冲区来拼接和处理完整行 accumulatedText chunk; const lines accumulatedText.split(\n); // 保留最后可能不完整的行 accumulatedText lines.pop() || ; for (const line of lines) { if (line.startsWith(data: )) { const dataStr line.slice(6); // 去掉 data: if (dataStr.trim() ) continue; // 忽略空行或心跳ping try { const data JSON.parse(dataStr); if (data.error) { assistantMessageDiv.textContent \n[错误: ${data.error}]; updateStatus(发生错误。); } else if (data.content ! null data.content ! undefined) { // 流式追加内容 assistantMessageDiv.textContent data.content; chatBox.scrollTop chatBox.scrollHeight; } if (data.finish_reason done || data.finish_reason error) { // 流结束可以清理资源 updateStatus(data.finish_reason done ? 回复完成。 : 流异常结束。); // 注意这里不能break while循环因为reader.read()可能还有后续数据如关闭信号 } } catch (e) { console.error(解析SSE数据失败:, e, 原始数据:, dataStr); } } // 忽略以 : 开头的注释行如心跳 : ping } } reader.releaseLock(); } catch (error) { console.error(请求失败:, error); assistantMessageDiv.textContent \n[请求失败: ${error.message}]; updateStatus(请求失败。); } finally { // 重新启用输入 userInput.disabled false; sendBtn.disabled false; userInput.focus(); currentEventSource null; assistantMessageDiv null; } }); // 支持回车发送 userInput.addEventListener(keypress, (e) { if (e.key Enter) { sendBtn.click(); } }); /script /body /html为了让这个页面能被访问我们需要在main.py中加一个路由来提供它。from fastapi.staticfiles import StaticFiles from fastapi.templating import Jinja2Templates from fastapi.responses import HTMLResponse app.mount(/static, StaticFiles(directorystatic), namestatic) templates Jinja2Templates(directorytemplates) app.get(/, response_classHTMLResponse) async def read_root(request: Request): return templates.TemplateResponse(index.html, {request: request})现在访问http://localhost:8000/你就可以看到一个简单的聊天界面。输入问题点击发送就能看到AI的回答一个字一个字地“流”出来了。前端实现的关键抉择EventSourcevsFetch API上面的前端代码注释里提到了一个关键点原生的EventSource对象虽然简单但它只支持GET请求且不能自定义请求头如Authorization。这在需要认证或发送复杂POST body如我们的聊天消息列表的场景下是硬伤。因此在生产环境中更推荐使用Fetch API来消费SSE流正如我们示例中所做。fetch()返回的Response对象的body属性是一个可读流ReadableStream我们可以用getReader()来读取。虽然需要手动解析SSE格式处理data:前缀和\n\n分隔符但这带来了完全的灵活性支持任何HTTP方法、可以添加自定义头、可以发送任意格式的请求体。手动解析SSE流时要注意数据块chunk的边界不一定与SSE事件边界对齐所以需要一个缓冲区accumulatedText来拼接和处理完整的行。这是一个常见的细节处理。5. 进阶优化与生产环境考量一个能跑通的Demo只是第一步。要让这个“极简智能体”真正可靠、可用还需要考虑以下几个进阶问题。5.1 连接管理与超时控制SSE连接是长连接。如果不加管理可能会有大量闲置连接占用服务器资源。客户端超时与重连浏览器端的EventSource对象在连接断开时会自动尝试重连。你可以通过监听onerror事件并检查eventSource.readyState来进行更精细的控制。对于fetch方案你需要自己实现重连逻辑。服务器端超时Starlette/FastAPI本身没有为SSE连接设置特定超时。但你可以通过底层ASGI服务器如Uvicorn的配置或者在前端发送周期性“心跳”并在后端检查最后活动时间的方式来清理死连接。sse-starlette库的EventSourceResponse自带发送: ping注释作为心跳的功能这有助于保持连接活跃并被负载均衡器识别。5.2 错误处理与重试机制流式传输过程中网络波动、服务端错误都可能发生。前端重试当fetch流意外中断非正常结束前端应提示用户并提供“重试”按钮。重试时最好能重新发送上次的请求消息。后端优雅降级如果AI服务调用失败后端不应直接抛出500错误导致连接粗暴断开。应该像我们示例中那样捕获异常并通过SSE流发送一个格式化的错误事件 ({error: ..., finish_reason: error})让前端能优雅地展示错误信息并结束流。限流与配额在API网关或应用层对/chat/stream这类端点实施限流如令牌桶算法防止恶意用户耗尽你的AI API配额。FastAPI可以很方便地集成slowapi等限流库。5.3 性能与并发异步框架虽然擅长处理高并发I/O但在流式场景下仍需注意生成器内存确保你的异步生成器 (event_generator或stream_openai_response) 不会在内存中积累大量数据。我们的代码是“来一块数据就yield一块”内存占用是常数级的这是正确的做法。数据库与外部服务如果在流式生成过程中需要查询数据库或调用其他微服务务必使用异步客户端如asyncpgfor PostgreSQL,aioredisfor Redis,httpxfor HTTP。任何同步的阻塞调用都会卡住整个事件循环拖累所有并发请求。背景任务如果流式响应结束后还有一些清理工作如记录日志、更新统计数据不要放在生成器函数里做因为连接关闭后这些代码可能无法执行。应该使用FastAPI的BackgroundTasks在响应返回后异步执行。5.4 部署与代理配置这是最容易踩坑的地方。当你把应用部署到生产环境前面通常会有一个反向代理如Nginx。Nginx配置要点location /chat/stream { # 关键禁用代理缓冲否则数据会在Nginx处堆积无法实时推送到客户端 proxy_buffering off; # 关键禁用Nginx对后端响应的缓存 proxy_cache off; # 设置较长的超时时间因为SSE是长连接 proxy_read_timeout 300s; proxy_connect_timeout 75s; # 确保能正确传递这些头部 proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; # 将请求代理到你的Uvicorn服务器 proxy_pass http://backend_server; }proxy_buffering off;这一行至关重要。Nginx默认会缓冲后端服务器的响应以达到优化性能的目的。但对于SSE缓冲会导致数据在Nginx内存中堆积直到达到一定大小或超时才会发送给客户端这就完全破坏了“流式”的实时性。X-Accel-Buffering: no这个响应头通常就是用来控制Nginx这个行为的但在配置文件中直接关闭更彻底。6. 踩坑实录那些年我们遇到的流式问题在实际开发和线上运维中我遇到过不少关于流式输出的“坑”。这里分享几个典型的希望能帮你提前避雷。坑一StreamingResponse被中间件或框架提前关闭有一次我在一个全局的异常处理中间件里习惯性地在捕获异常后返回了一个JSONResponse。结果发现当流式端点内部发生异常时前端收到的不是流式错误事件而是一个普通的JSON错误。原因是异常中间件拦截了异常并生成了新的响应覆盖了原本的StreamingResponse。解决方案在FastAPI中对于流式端点要谨慎使用会生成响应的全局异常处理器。更好的做法是在流式生成器函数内部如stream_openai_response进行try...except并在异常块中yield错误事件就像我们示例代码中做的那样。这样错误信息也能通过流的形式发送出去。坑二前端收不到数据或者数据攒在一起才收到这个问题十有八九出在代理缓冲上。无论是开发环境如果你用了某些开发服务器的代理功能还是生产环境Nginx/Apache一定要检查缓冲是否被禁用。除了配置代理服务器确保后端响应头包含了X-Accel-Buffering: no和Cache-Control: no-cache也是一个好习惯。坑三连接数过多服务器资源耗尽SSE是长连接每个活跃用户都会占用一个连接。如果用户不关闭页面连接会一直保持。在用户量大的时候这可能导致服务器文件描述符或内存耗尽。解决方案设置合理的心跳和超时让不活跃的连接自动关闭。考虑连接池化或使用更高效协议对于超大规模、双向通信需求强的场景SSE可能不是最优解可以评估WebSocket。但对于大多数单向推送场景SSE配合良好的连接管理是足够的。监控与告警监控服务器的连接数、内存和CPU使用率设置告警阈值。坑四async for循环中的阻塞操作在async for chunk in stream:循环里如果你不小心执行了一个同步的、耗时的操作比如用requests库发一个HTTP请求或者进行一个复杂的CPU计算整个事件循环就会被阻塞所有其他并发请求都会被卡住。解决方案牢记“异步环境中一切皆需异步”。调用外部服务用httpx.AsyncClient访问数据库用异步驱动如asyncpg,aiomysqlCPU密集型任务考虑放到单独的线程池中执行asyncio.to_thread。坑五前端EventSource无法发送POST请求或自定义Header这是我们之前讨论过的。很多开发者一开始都会想用EventSource因为它简单但很快就会被GET方法和无自定义Header的限制卡住。这时需要果断切换到Fetch API 手动解析SSE流的方案。虽然代码量多一点但换来的是完全的灵活性和可控性。