拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于SSE实现AI对话流式响应:Spring Boot后端与前端EventSource实战

1. 项目概述从“一问一答”到“流式对话”的体验升级最近在做一个全栈项目核心功能是让一个AI“章鱼哥”来解答用户的各种问题。最开始我把它做成了一个典型的“后端接口-前端展示”模式用户在前端输入框里提问点击提交然后页面转圈圈等个一两秒后端处理完一次性把完整的答案吐回来前端再整个渲染出来。功能是实现了但体验上总觉得差点意思尤其是答案比较长的时候用户得干等着看着那个加载动画心里没底。这就像你去餐厅点餐服务员记下菜单后消失在后厨过了十分钟才端上来一盘完整的菜中间你完全不知道后厨是在切菜、炒菜还是厨师在发呆。所以这次实战的目标很明确把这种“后端回答”的模式升级为“流式对话界面”。简单说就是让答案像水流一样一个字一个字、或者一个词一个词地“流”到前端页面上用户可以实时看到答案的生成过程。这种体验在现在的AI应用里非常普遍它能极大地降低用户的等待焦虑提升交互的沉浸感。这不仅仅是加个动画那么简单它涉及到前后端通信方式的根本性改变从传统的“请求-响应”变成了“请求-流式响应”是一个挺典型的全栈优化场景。这个项目我称之为“Vibe Coding”实战意思是不只关注功能实现更关注整个开发过程的“感觉”和最终产品的“体验”。我们需要打通从前端的事件监听、网络请求到后端的流式数据生成与推送再到前端的渐进式渲染这整条链路。下面我就把从零开始构建这个流式对话界面的完整思路、技术选型、踩坑记录和优化心得分享出来。2. 技术架构与核心思路拆解要实现流式对话核心在于改变数据交换的协议和方式。传统的RESTful API一次请求返回一个完整的JSON对象而流式响应需要服务器能够持续不断地向客户端发送数据片段。2.1 为什么选择 Server-Sent Events (SSE)面对流式数据推送我们有几个备选方案WebSocket、Server-Sent Events (SSE) 和长轮询。WebSocket功能最强大支持全双工通信客户端和服务器可以随时互相发送消息。但它也最重需要建立独立的、持久的连接协议也相对复杂。对于我们的场景——主要是服务器向客户端单向推送文本流——有点杀鸡用牛刀。长轮询 (Long Polling)一种模拟实时性的“黑客”手段。客户端发起请求服务器持有这个请求直到有数据或超时。收到响应后客户端立即发起下一个请求。这种方式实现简单但效率低下会产生大量HTTP请求并且延迟不可控。Server-Sent Events (SSE)这正是为我们这种场景量身定制的。它基于普通的HTTP协议允许服务器主动向客户端推送数据。连接建立后服务器可以持续发送事件流而客户端通过EventSourceAPI来监听。它的优点是协议简单、轻量天然支持自动重连并且与HTTP生态如认证、代理兼容性好。注意SSE是单向的服务器-客户端。如果你的对话场景需要客户端在流式接收过程中频繁中断或发送指令例如“停止生成”那么单纯的SSE会有点吃力可能需要结合额外的API调用。不过对于基本的问答流式输出SSE是完美选择。因此我的架构决策很清晰后端使用SSE协议输出流式文本前端使用EventSource或兼容的库来接收并实时渲染。2.2 整体数据流设计整个流程可以分解为以下几个关键步骤前端触发用户在界面输入问题点击“发送”。建立连接前端创建一个指向特定后端端点的EventSource连接。后端处理后端接收到请求开始调用AI模型或任何文本生成器。模型不是一次性生成全部文本而是以“token”可以理解为词或字为单位逐步生成。每生成一个或一小批token后端就将其封装成SSE格式data: {“content”: “生成的词”}\n\n并立即写入响应流。这个过程持续进行直到整个答案生成完毕。前端流式渲染前端通过EventSource的onmessage事件监听器实时收到这些数据块。每收到一块就将其追加到页面上的对话气泡或答案区域中。连接关闭当后端生成完毕发送一个特殊的事件如[DONE]或直接关闭连接前端随之更新UI状态如隐藏加载指示器。这个设计的关键在于后端响应体的flush操作。我们必须确保每个数据块都能立即从服务器缓冲区发送到客户端而不是等所有数据都生成完再一次性发送。3. 后端实现Spring Boot下的SSE端点我用的后端框架是Spring Boot它提供了对异步处理和响应式编程的良好支持非常适合实现SSE。3.1 创建SSE控制器首先定义一个REST控制器其端点返回SseEmitter对象。SseEmitter是Spring对SSE的抽象封装。import org.springframework.web.bind.annotation.*; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.io.IOException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; RestController RequestMapping(/api/chat) public class StreamChatController { // 建议使用线程池来管理异步任务避免为每个请求创建新线程 private final ExecutorService nonBlockingService Executors.newCachedThreadPool(); GetMapping(path /stream, produces text/event-stream) public SseEmitter streamAnswer(RequestParam String question) { // 设置连接超时时间0表示永不超时对于长流可以设置一个较长的时间如30秒 SseEmitter emitter new SseEmitter(0L); // 提交一个异步任务来处理流式生成 nonBlockingService.execute(() - { try { // 模拟或实际调用AI模型生成流式内容 simulateStreamingGeneration(question, emitter); // 生成完成后发送完成事件或直接完成 emitter.send(SseEmitter.event().name(complete).data([DONE])); emitter.complete(); } catch (IOException e) { // 发生IO异常通常是客户端断开连接直接完成并清理 emitter.completeWithError(e); } catch (Exception e) { // 其他业务异常 emitter.completeWithError(e); } }); // 设置连接结束时的回调用于资源清理 emitter.onCompletion(() - System.out.println(SSE连接完成)); emitter.onTimeout(() - System.out.println(SSE连接超时)); emitter.onError((ex) - System.out.println(SSE连接错误: ex.getMessage())); return emitter; } private void simulateStreamingGeneration(String question, SseEmitter emitter) throws IOException, InterruptedException { // 这里模拟一个AI模型逐步生成答案的过程 String fullAnswer 你好我是章鱼哥这是一个模拟的流式回答。我将逐词输出这句话。; String[] tokens fullAnswer.split(); // 按字拆分实际场景可能是按模型返回的token for (String token : tokens) { // 关键每次生成一点就发送一点 // SSE格式要求data: 内容\n\nSpring的SseEmitter帮我们处理了格式 emitter.send(SseEmitter.event().data(token)); // 模拟模型生成每个token需要的时间 Thread.sleep(50); } } }实操心得一连接管理与超时SseEmitter的构造函数可以传入超时时间毫秒。对于流式对话这个时间应该设置得足够长或者设为0不超时。因为一次生成可能持续数十秒。但设为0有风险如果客户端异常断开服务器端连接可能不会及时释放。更稳健的做法是设置一个合理的超时如5分钟并在前端实现断线重连逻辑。实操心得二异常处理与资源清理流式连接生命周期较长异常处理至关重要。onCompletion、onTimeout、onError这三个回调必须设置好用于记录日志和清理与该连接关联的资源如中断模型生成任务。否则可能导致内存泄漏或后台任务空转。3.2 集成真实的AI模型流式调用上面的例子是模拟的。真实场景中你需要调用支持流式输出的AI模型API例如OpenAI的Chat Completions API设置stream: true或本地部署的类似Ollama、vLLM等支持流式的模型服务。以调用OpenAI风格接口为例伪代码private void callAIStreaming(String question, SseEmitter emitter) throws IOException { // 假设有一个HttpClient配置为流式读取响应 HttpClient client HttpClient.newHttpClient(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://your-ai-server/v1/chat/completions)) .header(Content-Type, application/json) .POST(HttpRequest.BodyPublishers.ofString( {\model\:\gpt-3.5-turbo\,\messages\:[{\role\:\user\,\content\:\ question \}],\stream\:true} )) .build(); client.sendAsync(request, HttpResponse.BodyHandlers.ofLines()) .thenAccept(response - { // 逐行读取服务器返回的SSE流 response.body().forEach(line - { if (line.startsWith(data: )) { String data line.substring(6).trim(); if ([DONE].equals(data)) { emitter.complete(); } else if (!data.isEmpty()) { // 解析JSON提取delta content // 假设返回格式为data: {choices:[{delta:{content:Hi}}]} try { JsonNode node objectMapper.readTree(data); String content node.path(choices).get(0).path(delta).path(content).asText(); if (content ! null !content.isEmpty()) { emitter.send(SseEmitter.event().data(content)); } } catch (Exception e) { // 忽略单次解析错误继续处理后续流 } } } }); }); }这里的关键是后端作为“中间人”需要将从AI模型服务收到的流式数据几乎无延迟地转发给前端客户端。这要求后端处理响应体时也必须采用流式方式不能等全部读完再转发。4. 前端实现使用EventSource接收与渲染前端相对直接主要使用EventSourceAPI。但为了更好的兼容性和功能如自定义头部、错误重试我推荐使用一个轻量级的封装库如eventsource-polyfill或直接使用Fetch API模拟SSE。这里先用原生EventSource演示。4.1 建立连接与基本监听class StreamChatUI { constructor() { this.eventSource null; this.answerContainer document.getElementById(answer-container); this.questionInput document.getElementById(question-input); this.sendButton document.getElementById(send-button); this.isStreaming false; this.bindEvents(); } bindEvents() { this.sendButton.addEventListener(click, () this.startStreaming()); // 也可以监听输入框的Enter键 this.questionInput.addEventListener(keypress, (e) { if (e.key Enter !e.shiftKey) { e.preventDefault(); this.startStreaming(); } }); } startStreaming() { if (this.isStreaming) { console.warn(正在流式输出中请等待完成或取消。); return; } const question this.questionInput.value.trim(); if (!question) return; // 清空上一次的答案显示加载状态 this.answerContainer.innerHTML ; this.showLoadingIndicator(); // 构建SSE连接URL const url /api/chat/stream?question${encodeURIComponent(question)}; // 创建EventSource连接 this.eventSource new EventSource(url); this.isStreaming true; // 监听消息事件 this.eventSource.onmessage (event) { // 收到的数据就是后端emitter.send()里的data const chunk event.data; // 如果后端发送了特殊的结束标记 if (chunk [DONE]) { this.finishStreaming(); return; } // 将收到的数据块追加到答案容器中 this.appendAnswerChunk(chunk); }; // 监听自定义事件如果后端用了emitter.send(SseEmitter.event().name(complete)...) this.eventSource.addEventListener(complete, (event) { console.log(流式传输完成); this.finishStreaming(); }); // 监听错误事件 this.eventSource.onerror (error) { console.error(EventSource 错误:, error); // 错误发生时尝试关闭连接并更新UI this.finishStreaming(); this.showError(连接出现异常请重试。); }; } appendAnswerChunk(chunk) { // 简单的追加实际中你可能需要处理Markdown、代码高亮等 this.answerContainer.textContent chunk; // 滚动到底部确保用户能看到最新内容 this.answerContainer.scrollTop this.answerContainer.scrollHeight; } finishStreaming() { this.isStreaming false; if (this.eventSource) { this.eventSource.close(); this.eventSource null; } this.hideLoadingIndicator(); } showLoadingIndicator() { /* UI逻辑 */ } hideLoadingIndicator() { /* UI逻辑 */ } showError(msg) { /* UI逻辑 */ } } // 初始化 new StreamChatUI();4.2 处理复杂场景与用户体验优化原生EventSource有几个限制1) 不支持POST请求体只能GET参数放URL2) 不支持自定义HTTP头部如认证Token。对于生产环境这通常是不可接受的。解决方案使用Fetch API模拟SSE我们可以用Fetch API来发起请求然后手动读取流式响应体这给了我们最大的灵活性。async function startStreamingWithFetch(question) { const controller new AbortController(); const signal controller.signal; try { const response await fetch(/api/chat/stream, { method: POST, // 可以使用POST了 headers: { Content-Type: application/json, Authorization: Bearer ${yourToken} // 可以自定义头部了 }, body: JSON.stringify({ question: question }), signal: signal // 用于后续取消请求 }); if (!response.ok || !response.body) { throw new Error(HTTP error! status: ${response.status}); } const reader response.body.getReader(); const decoder new TextDecoder(); while (true) { const { done, value } await reader.read(); if (done) { console.log(Stream finished); break; } // 解码并处理数据块 const chunk decoder.decode(value, { stream: true }); // SSE数据流是以\n\n分隔的多个事件需要按行解析 processSSEChunk(chunk); } } catch (error) { if (error.name AbortError) { console.log(请求被用户取消); } else { console.error(流式请求失败:, error); showError(获取回答失败); } } finally { hideLoadingIndicator(); } } function processSSEChunk(rawChunk) { // 简单的行解析器 const lines rawChunk.split(\n); for (const line of lines) { if (line.startsWith(data: )) { const data line.substring(6).trim(); if (data [DONE]) { // 结束处理 return; } if (data) { try { // 假设后端直接发送文本内容或解析JSON const parsed JSON.parse(data); appendAnswerChunk(parsed.content || parsed); } catch { // 如果不是JSON直接当作文本处理 appendAnswerChunk(data); } } } } }实操心得三前端的“停止生成”功能这是一个提升用户体验的关键功能。当答案生成到一半用户可能觉得不对想重问。使用Fetch API后我们可以通过AbortController轻松实现。let currentAbortController null; function askQuestion(question) { // 如果已有请求在进行先中止它 if (currentAbortController) { currentAbortController.abort(); } currentAbortController new AbortController(); startStreamingWithFetch(question, currentAbortController.signal); } // 停止按钮的点击事件 stopButton.addEventListener(click, () { if (currentAbortController) { currentAbortController.abort(); currentAbortController null; updateUIForStopped(); } });实操心得四渲染优化与打字机效果直接textContent chunk在内容很多时可能导致性能问题。更好的做法是使用文档片段DocumentFragment进行批量更新或者使用requestAnimationFrame来调度渲染。如果想实现更优雅的“打字机”效果逐字出现可以在收到完整句子或词后用CSS动画或setInterval控制每个字符的显示延迟。5. 部署、调试与性能考量将这套流式系统部署到生产环境还需要考虑几个实际问题。5.1 连接数与服务器资源每个SSE连接都是一个长期的HTTP连接会占用一个线程或一个非阻塞的IO通道。对于Tomcat这样的传统Servlet容器大量并发长连接可能消耗大量线程资源。可以考虑调整服务器配置增加最大线程数、连接超时时间。使用响应式Web框架如Spring WebFlux它基于Netty专为高并发、长连接场景设计资源利用率更高。网关/代理配置确保Nginx等反向代理配置了合适的proxy_buffering off;和长超时时间以免代理层缓冲或切断数据流。5.2 网络与稳定性心跳机制为了防止中间网络设备代理、防火墙因长时间无数据而断开连接服务器可以定期发送注释行以:开头的行作为心跳。Spring Boot的SseEmitter可以发送SseEmitter.event().comment(“keepalive”)。自动重连EventSource原生支持自动重连。如果连接意外断开它会尝试重新连接。但重连后如何恢复上下文例如继续之前的回答是个复杂问题通常需要客户端携带一个会话ID。5.3 监控与调试浏览器开发者工具在Network标签页你可以看到类型为eventsource的请求点击可以查看流式传输的事件列表非常直观。服务器日志记录SSE连接的生命周期创建、完成、错误、超时便于排查连接泄漏问题。流量观察观察服务器带宽和连接数确保流式传输没有成为性能瓶颈。6. 常见问题与排查实录在实际开发中我遇到了不少坑这里记录下最典型的几个及其解决方法。问题1前端收不到数据或者收到的是完整的一大段不是流式的。排查首先检查浏览器Network面板看对应请求的响应类型是否是text/event-stream以及数据是否是一段一段陆续到达的。如果是一次性到达问题出在后端。解决确保后端代码中在每次emitter.send(data)后都执行了response.flushBuffer()Spring的SseEmitter通常内部处理了。如果使用了RestController避免返回值被包装成统一的JSON格式。检查是否有拦截器或过滤器缓冲了响应。问题2连接几秒后就自动断开了。排查可能是服务器或代理的超时设置太短。解决Spring Boot检查SseEmitter的超时设置。对于Tomcat可能需要调整server.tomcat.connection-timeout和server.tomcat.keep-alive-timeout。Nginx在location配置中增加proxy_read_timeout 300s;设置一个很长的超时和proxy_buffering off;。前端检查EventSource的重连逻辑是否正常触发。问题3前端显示乱码或数据解析错误。排查检查SSE格式是否正确。每个事件必须以data:开头以两个换行符\n\n结束。数据本身最好是纯文本或JSON字符串。解决后端确保发送格式正确。前端使用Fetch API读取流时注意TextDecoder的使用并正确处理数据块拼接因为一个chunk可能只包含事件的一部分。问题4在开发环境正常部署到生产服务器后流式失效。排查这通常是环境差异导致。生产环境可能有额外的网关、负载均衡器、防火墙或安全策略。解决逐层检查。从应用服务器日志开始看连接是否建立成功。然后检查网关/负载均衡器如Nginx, API Gateway的配置确保它们支持并正确转发了text/event-stream。最后检查网络安全组或防火墙规则是否允许长连接。从传统的“后端回答”升级到“流式对话界面”虽然增加了前后端的一些复杂度但对用户体验的提升是质的飞跃。这个过程让我深刻体会到全栈开发不仅仅是前后端功能的简单拼接更是需要对数据流动的协议、性能、异常情况有全局的掌控。选择SSE作为技术方案在满足需求的同时保持了架构的简洁。前端的Fetch API方案提供了比原生EventSource更强的灵活性适合大多数需要认证和自定义请求的生产场景。最后一个小技巧在开发初期可以先用一个简单的、返回模拟流式数据的后端接口让前端先跑通整个接收和渲染流程。然后再去集成复杂的AI模型调用这样能有效隔离问题提高调试效率。当看到第一个词从屏幕上“流”出来的时候那种成就感就是“Vibe Coding”最好的诠释。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门