1. 项目概述当Java应用遇上AI流式输出最近在折腾几个AI应用项目发现一个挺普遍的需求用户输入一个问题后台大模型吭哧吭哧推理半天最后才一次性吐出所有答案。这个等待过程对用户来说简直是煎熬尤其是答案比较长的时候屏幕一片空白用户心里直打鼓——是卡死了吗还是我网断了这种体验非常糟糕。于是“流式输出”就成了提升AI应用交互体验的刚需。简单说就是让AI像真人聊天一样一个字一个字、或者一个词一个词地“流”出来用户能实时看到生成过程感知到进度体验立马就上去了。在Java生态里尤其是Spring Boot框架下实现AI的流式输出Server-Sent Events是目前最主流、也最优雅的方案之一很多人也简称它为SSE。它和WebSocket不同SSE是服务器向浏览器单向推送的技术特别适合这种“服务器主动推送连续数据流”的场景比如新闻更新、股票行情当然还有我们今天的主题——AI生成内容。相比于WebSocket的双向通信SSE更轻量实现起来也更简单天然就是为这种“只读”的数据流设计的。所以今天我们就来深挖一下在Java特别是Spring Boot应用里如何从原理到代码稳稳当当地实现AI的流式输出。我会结合我最近在项目里的实战把SSE协议怎么工作、Spring Boot里怎么集成、怎么对接大模型的流式接口、以及那些容易踩的坑都掰开揉碎了讲清楚。无论你是想给自己项目增加点“科技感”还是面试时被问到这块能对答如流这篇文章都能给你整明白。2. 核心原理SSE协议与流式传输机制要玩转流式输出首先得搞清楚SSE是怎么一回事。别被名字吓到它的原理其实非常直观。2.1 SSE协议的工作机制你可以把SSE想象成服务器打开了一个通向客户端比如浏览器的单向“水管”。一旦连接建立服务器就可以随时通过这根“水管”推送消息而客户端则持续监听收到消息就立刻处理。这个连接基于普通的HTTP/HTTPS所以兼容性极好。SSE通信有自己的一套简单格式。服务器推送的每条消息由若干个字段组成用换行符分隔最后以两个换行符\n\n表示一条消息结束。核心字段有三个data:消息的数据内容。如果一条消息很长可以分成多行每行前面都加上data:。event:自定义事件类型。客户端可以根据不同的事件类型来触发不同的处理逻辑。这是个可选字段非常有用。id:消息ID。主要用于连接意外中断后重连客户端可以通过Last-Event-ID头告诉服务器“我从哪条消息之后开始断的”从而实现断点续传。一个典型的SSE响应流看起来是这样的HTTP/1.1 200 OK Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive event: start data: 连接已建立开始生成内容... id: 1 data: 这是第一段 data: 连续的内容。 event: chunk data: {token: 这是, delta: 一段} id: 2 event: chunk data: {token: 流式, delta: 输出} id: 3 event: done data: 内容生成完毕。 id: 4注意响应头Content-Type: text/event-stream这是告诉客户端“接下来是SSE流”的关键。连接会一直保持直到服务器主动关闭或网络异常。2.2 为何选择SSE而非WebSocket这里有个常见的选型困惑流式输出为啥常用SSE而不是听起来更强大的WebSocket协议复杂度WebSocket是一个独立的、全双工的协议需要一次额外的“握手”升级连接。SSE则直接跑在HTTP之上实现简单无需额外协议。通信模型AI内容生成绝大多数时候是服务器单向推送给客户端。WebSocket的双工能力客户端也能随时主动发消息在这里是冗余的反而增加了实现的复杂性和资源开销。自动重连SSE客户端内置了自动重连机制。连接断开后它会自动尝试重新连接并可以携带上次最后的消息ID。WebSocket需要自己实现这套逻辑。浏览器兼容性与开发便利性现代浏览器都原生支持SSE前端用EventSource对象几行代码就能搞定。虽然WebSocket支持也很广但SSE的API对这类只读数据流场景更友好。所以对于AI流式输出这种典型的“服务器推送、客户端接收”模式SSE是更轻量、更专注、也更合适的技术选型。当然如果你的应用场景同时需要前端频繁地向服务器发送复杂指令比如一个实时协作的AI绘图应用那WebSocket可能更合适。2.3 AI模型侧的流式支持光有SSE通道还不够关键是AI大模型本身要支持流式响应。目前主流的模型API如OpenAI的Chat Completion、国内各大模型的类似接口都提供了“流式”调用的选项。通常你在调用API时设置一个参数如stream: true那么API返回的就不是一个完整的JSON响应体而是一个HTTP流HTTP Stream。这个流里服务器会持续返回一系列以data:开头的SSE格式数据块每个数据块是一个独立的JSON对象包含模型最新生成的一小段内容通常是一个token或一个词。最后一个数据块会有特殊的结束标志。我们的Java后端服务角色就是作为这个“AI流”和“前端SSE流”之间的中转站和适配器它要以流式方式调用AI API然后实时地将收到的一个个数据块通过自己建立的SSE连接推送给前端。3. Spring Boot中实现SSE流式输出的实战理论清楚了我们动手在Spring Boot里搭一个。我会用一个模拟AI生成故事的场景来演示。3.1 环境准备与基础依赖首先创建一个标准的Spring Boot项目这里用3.x版本。除了基本的Web依赖我们不需要为SSE引入任何特殊库因为Spring MVC已经原生支持了。pom.xml关键依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency是的就这么简单。SSE的核心是SseEmitter类它已经在spring-web模块里了。3.2 构建SSE控制器与连接管理我们来创建一个控制器用于处理前端的SSE连接请求并管理这些连接。import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.*; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.io.IOException; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; RestController RequestMapping(/api/sse) public class SseController { // 用于存储每个用户的SseEmitter键可以是用户ID或会话ID private final MapString, SseEmitter emitterMap new ConcurrentHashMap(); /** * 前端建立SSE连接 * param clientId 客户端标识可以从请求头或参数传入 * return SseEmitter */ GetMapping(path /connect, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter connect(RequestParam String clientId) { // 设置连接超时时间0表示永不超时生产环境建议设置一个合理值如30秒或1分钟 SseEmitter emitter new SseEmitter(0L); emitterMap.put(clientId, emitter); // 设置连接完成和超时/错误时的回调用于清理资源 emitter.onCompletion(() - { System.out.println(SSE连接完成: clientId); emitterMap.remove(clientId); }); emitter.onTimeout(() - { System.out.println(SSE连接超时: clientId); emitter.complete(); emitterMap.remove(clientId); }); emitter.onError((ex) - { System.out.println(SSE连接错误: clientId , error: ex.getMessage()); emitter.completeWithError(ex); emitterMap.remove(clientId); }); // 发送一个初始事件告知客户端连接成功 try { emitter.send(SseEmitter.event() .name(connect) // 事件类型 .data(SSE连接已建立客户端ID: clientId) .id(String.valueOf(System.currentTimeMillis()))); // 用时间戳作为初始ID } catch (IOException e) { emitter.completeWithError(e); } return emitter; } /** * 关闭指定客户端的SSE连接 */ PostMapping(/disconnect) public String disconnect(RequestParam String clientId) { SseEmitter emitter emitterMap.remove(clientId); if (emitter ! null) { emitter.complete(); return 连接已关闭: clientId; } return 未找到连接: clientId; } // 提供一个方法让其他服务如AI调用服务能获取到Emitter并发送消息 public SseEmitter getEmitter(String clientId) { return emitterMap.get(clientId); } }关键点解析produces MediaType.TEXT_EVENT_STREAM_VALUE这个注解是关键它确保响应的Content-Type是text/event-stream。SseEmitter这是Spring提供的用于发送SSE事件的类。创建时可以指定超时时间。连接管理使用一个ConcurrentHashMap来存储活跃的连接是常见做法。务必在onCompletion、onTimeout、onError回调中从Map里移除对应的emitter防止内存泄漏。这是第一个大坑忘了清理Map会越来越大最终导致OutOfMemoryError。emitter.send()用于发送事件。你可以构建不同name事件类型和data数据内容的事件。数据可以是String也可以是Spring能序列化成JSON的对象。3.3 模拟AI流式生成与推送服务现在我们创建一个服务它模拟调用一个流式AI接口并将生成的内容块实时推送给对应的SSE客户端。import com.fasterxml.jackson.databind.ObjectMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import java.io.IOException; import java.time.LocalTime; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; Service Slf4j public class AiStreamService { private final SseController sseController; private final ObjectMapper objectMapper; // Jackson用于序列化JSON private final ExecutorService executorService Executors.newCachedThreadPool(); public AiStreamService(SseController sseController, ObjectMapper objectMapper) { this.sseController sseController; this.objectMapper objectMapper; } /** * 处理AI生成请求并流式推送结果 * param clientId 客户端ID * param prompt 用户输入的提示词 */ public void streamAiResponse(String clientId, String prompt) { executorService.submit(() - { SseEmitter emitter sseController.getEmitter(clientId); if (emitter null) { log.warn(客户端 {} 的SSE连接不存在或已关闭, clientId); return; } try { // 1. 发送开始事件 emitter.send(SseEmitter.event() .name(start) .data(开始处理您的请求: \ prompt \) .id(generateMessageId())); // 2. 模拟调用流式AI API并处理数据块 // 这里用一个字符串数组模拟AI流式返回的多个片段 String simulatedResponse 在一个阳光明媚的下午Java程序员小李正在调试一个复杂的并发问题。; String[] tokens simulatedResponse.split(); // 简单按字拆分实际按token for (int i 0; i tokens.length; i) { // 模拟网络延迟和AI思考时间 Thread.sleep(50 (int)(Math.random() * 50)); // 构建要推送的数据对象 AiResponseChunk chunk new AiResponseChunk(tokens[i], i, LocalTime.now().toString()); // 将对象序列化为JSON字符串发送 String chunkData objectMapper.writeValueAsString(chunk); // 发送数据块事件 emitter.send(SseEmitter.event() .name(chunk) // 事件名前端根据这个来区分 .data(chunkData) .id(generateMessageId())); log.debug(向客户端 {} 发送第 {} 个数据块: {}, clientId, i, tokens[i]); } // 3. 发送结束事件 emitter.send(SseEmitter.event() .name(end) .data(内容生成完成。) .id(generateMessageId())); log.info(客户端 {} 的AI流式请求处理完毕, clientId); } catch (IOException e) { log.error(向客户端 {} 发送SSE消息失败: {}, clientId, e.getMessage()); emitter.completeWithError(e); } catch (InterruptedException e) { Thread.currentThread().interrupt(); log.error(处理线程被中断: {}, clientId); } catch (Exception e) { log.error(处理AI流式响应时发生未知错误: {}, clientId, e); try { emitter.send(SseEmitter.event() .name(error) .data(生成过程发生错误: e.getMessage())); } catch (IOException ex) { // 忽略发送错误时的异常 } emitter.completeWithError(e); } }); } private String generateMessageId() { return String.valueOf(System.currentTimeMillis()); } // 定义一个内部类表示AI返回的每个数据块 lombok.Data lombok.AllArgsConstructor private static class AiResponseChunk { private String token; // 当前生成的词或token private int index; // 序号 private String timestamp; // 时间戳 } }关键点解析异步处理AI生成可能是耗时的所以必须使用异步线程这里用了ExecutorService来处理避免阻塞SSE连接线程。Spring MVC的请求线程池是有限的阻塞会导致应用无法处理新请求。异常处理这是重中之重。网络不稳定、AI服务超时、客户端提前关闭连接等都可能导致IOException。必须捕获这些异常并调用emitter.completeWithError(e)或发送一个error事件通知前端同时做好资源清理在控制器的回调里已经做了。数据格式我们定义了一个AiResponseChunk对象来结构化每个数据块并用Jackson序列化成JSON再发送。这样前端解析起来非常方便。实际项目中这个对象的结构应该和你对接的真实AI API返回的流式数据块结构保持一致。事件命名使用不同name的事件start,chunk,end,error可以让前端清晰地知道当前处于哪个阶段从而进行不同的UI更新如显示开始提示、逐字显示、显示结束标志、显示错误信息。3.4 前端如何消费SSE流后端准备好了前端怎么接呢非常简单使用浏览器原生的EventSourceAPI。!DOCTYPE html html head titleJava AI 流式输出演示/title /head body h2AI故事生成器/h2 input typetext idpromptInput placeholder输入一个主题例如Java程序员的一天 size50/ button onclickstartStream()开始生成/button button onclickcloseConnection()停止接收/button br/br/ div idoutput stylewhite-space: pre-wrap; border:1px solid #ccc; padding:10px; min-height:200px; 生成的内容将显示在这里... /div script let eventSource null; const clientId user_ Math.random().toString(36).substr(2, 9); // 生成一个随机客户端ID // 1. 首先建立SSE连接 function connectSSE() { // 连接端点带上clientId const url http://localhost:8080/api/sse/connect?clientId${clientId}; eventSource new EventSource(url); // 监听通用消息未指定event类型 eventSource.onmessage function(event) { console.log(收到消息:, event.data); // 通常用于处理没有指定事件名的数据但建议都用addEventListener }; // 监听特定事件 eventSource.addEventListener(connect, function(event) { document.getElementById(output).innerHTML span stylecolor:green;${event.data}/spanbr; console.log(连接成功:, event.data); }); eventSource.addEventListener(start, function(event) { appendOutput(span stylecolor:blue;[开始] ${event.data}/spanbr); }); eventSource.addEventListener(chunk, function(event) { // 解析后端发送的JSON数据块 try { const chunk JSON.parse(event.data); // 将token追加到显示区域实现逐字输出效果 appendOutput(span${chunk.token}/span); } catch (e) { console.error(解析数据块失败:, e, event.data); } }); eventSource.addEventListener(end, function(event) { appendOutput(brspan stylecolor:green;[结束] ${event.data}/spanbr); // 生成结束后可以自动或手动关闭连接 // closeConnection(); }); eventSource.addEventListener(error, function(event) { console.error(SSE错误:, event); appendOutput(brspan stylecolor:red;[错误] 连接出现异常/spanbr); // 发生错误时EventSource会自动尝试重连 }); // 监听连接打开事件 eventSource.onopen function() { console.log(SSE连接已打开); }; } // 2. 触发AI生成 function startStream() { const prompt document.getElementById(promptInput).value; if (!prompt) { alert(请输入提示词); return; } // 先确保连接已建立 if (!eventSource || eventSource.readyState EventSource.CLOSED) { connectSSE(); // 简单起见这里设置一个延迟再发送请求确保连接先建立。 // 更好的做法是等connect事件收到后再触发。 setTimeout(() sendGenerateRequest(prompt), 500); } else { sendGenerateRequest(prompt); } } function sendGenerateRequest(prompt) { // 这里应该调用后端另一个接口触发AI生成。为了简化我们直接调用服务端方法。 // 实际项目中你需要发一个POST请求到后端例如/api/ai/generate // fetch(/api/ai/generate?clientId${clientId}prompt${encodeURIComponent(prompt)}, {method: POST}); console.log(触发生成clientId: ${clientId}, prompt: ${prompt}); // 由于我们做了简化这里假设请求已发出。实际开发中需要这个HTTP调用。 appendOutput(br请求已发送: ${prompt}br); // 模拟在实际代码中startStream函数里的setTimeout应该替换为真实的fetch调用。 } // 3. 关闭连接 function closeConnection() { if (eventSource) { eventSource.close(); eventSource null; appendOutput(brspan stylecolor:gray;[连接已手动关闭]/spanbr); // 也可以通知后端清理资源 fetch(/api/sse/disconnect?clientId${clientId}, {method: POST}); } } function appendOutput(html) { document.getElementById(output).innerHTML html; // 滚动到底部 const outputDiv document.getElementById(output); outputDiv.scrollTop outputDiv.scrollHeight; } // 页面加载时自动连接可选 window.onload connectSSE; /script /body /html前端代码的核心就是new EventSource(url)。通过为不同的事件类型connect,start,chunk等添加监听器我们就能在内容生成的每个阶段更新UI实现真正的流式体验。4. 对接真实AI大模型流式API上面的例子是模拟的。真实项目中我们需要在AiStreamService里替换掉那个模拟循环改为真实调用AI服务的流式接口。这里以使用Spring AI如果项目用了或简单的RestTemplate/WebClient调用OpenAI兼容API为例。关键思路我们需要以流式的方式消费AI API返回的HTTP响应体而不是等全部完成再解析。4.1 使用WebClient进行流式消费Spring 5引入了反应式编程的WebClient它非常适合处理流式响应。首先添加依赖如果还没引入dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-webflux/artifactId /dependency然后修改服务类中的核心方法import org.springframework.web.reactive.function.client.WebClient; import reactor.core.publisher.Flux; import org.springframework.http.HttpHeaders; import org.springframework.http.MediaType; Service Slf4j public class RealAiStreamService { // ... 其他注入 ... private final WebClient webClient; public RealAiStreamService(WebClient.Builder webClientBuilder) { this.webClient webClientBuilder.baseUrl(https://api.your-ai-provider.com/v1).build(); } public void streamAiResponseFromRealApi(String clientId, String prompt) { executorService.submit(() - { SseEmitter emitter sseController.getEmitter(clientId); if (emitter null) { return; } try { emitter.send(SseEmitter.event().name(start).data(开始连接AI模型...)); // 构建请求体注意设置 stream: true MapString, Object requestBody Map.of( model, gpt-3.5-turbo, messages, List.of(Map.of(role, user, content, prompt)), stream, true ); // 使用WebClient发起流式请求 FluxString responseFlux webClient.post() .uri(/chat/completions) .header(HttpHeaders.AUTHORIZATION, Bearer YOUR_API_KEY) .contentType(MediaType.APPLICATION_JSON) .bodyValue(requestBody) .accept(MediaType.TEXT_EVENT_STREAM) // 关键接受text/event-stream .retrieve() .bodyToFlux(String.class); // 将响应体转换为字符串流 // 订阅这个流每收到一块数据就处理一块 responseFlux.subscribe( dataLine - { // 成功收到一行数据 // AI API返回的流通常是多个data: {...}\n\n格式的数据行 if (dataLine.startsWith(data: )) { String jsonData dataLine.substring(6).trim(); if ([DONE].equals(jsonData)) { // 流结束标志 emitter.send(SseEmitter.event().name(end).data(生成完毕)); emitter.complete(); return; } try { // 解析JSON提取delta content JsonNode node objectMapper.readTree(jsonData); JsonNode choices node.path(choices); if (choices.isArray() choices.size() 0) { JsonNode delta choices.get(0).path(delta); String content delta.path(content).asText(null); if (content ! null !content.isEmpty()) { // 将内容块通过SSE推送给前端 AiResponseChunk chunk new AiResponseChunk(content, chunkIndex.getAndIncrement(), LocalTime.now().toString()); emitter.send(SseEmitter.event().name(chunk).data(objectMapper.writeValueAsString(chunk))); } } } catch (Exception e) { log.warn(解析AI流数据行失败: {}, line: {}, e.getMessage(), dataLine); } } }, error - { // 发生错误 log.error(消费AI API流失败, error); try { emitter.send(SseEmitter.event().name(error).data(AI服务调用失败: error.getMessage())); } catch (IOException ignored) {} emitter.completeWithError(error); }, () - { // 流正常结束 log.info(AI流消费完成); // 如果前面没遇到[DONE]这里可以发送结束事件 } ); } catch (Exception e) { log.error(发起AI流式请求失败, e); try { emitter.send(SseEmitter.event().name(error).data(请求初始化失败)); } catch (IOException ignored) {} emitter.completeWithError(e); } }); } }这里的关键是accept(MediaType.TEXT_EVENT_STREAM)和bodyToFlux(String.class)。这告诉WebClient我们期望一个流式响应并且我们要以Flux一个反应式流的形式来消费它实现边收边处理。4.2 使用Spring AI简化操作如果你的项目使用Spring AI整个过程会变得更加声明式和简洁。Spring AI抽象了不同模型供应商的API提供了统一的流式调用接口。import org.springframework.ai.chat.ChatClient; import org.springframework.ai.chat.prompt.Prompt; import org.springframework.ai.chat.messages.UserMessage; import org.springframework.ai.chat.ChatResponse; import reactor.core.publisher.Flux; Service public class SpringAiStreamService { private final ChatClient chatClient; // 注入ChatClient具体实现由配置的模型决定如OpenAI public void streamWithSpringAi(String clientId, String prompt) { // ... 获取emitter逻辑同上 ... // 构建Prompt Prompt springPrompt new Prompt(new UserMessage(prompt)); // 流式调用返回的是一个FluxChatResponse FluxChatResponse responseFlux chatClient.stream(springPrompt); // 订阅并处理流 responseFlux.subscribe( chatResponse - { // 每个ChatResponse包含一个Generation里面有生成的Content String content chatResponse.getResult().getOutput().getContent(); if (content ! null) { // 推送内容块给前端 // ... 发送SSE事件 ... } }, error - { /* 错误处理 */ }, () - { /* 完成处理 */ } ); } }使用Spring AI后我们不再需要关心不同AI供应商的API细节和HTTP流解析框架帮我们处理了这些复杂性让代码更专注于业务逻辑。5. 生产环境中的关键考量与避坑指南把Demo跑起来只是第一步要上线稳定运行以下几个坑必须得填平。5.1 连接管理与资源泄漏这是SSE应用最常见的痛点。每个SseEmitter都会持有一定的资源如响应输出流。如果不妥善管理会导致内存泄漏emitterMap不断增长最终OutOfMemoryError。文件描述符耗尽大量未关闭的连接占用系统资源。解决方案设置合理超时创建SseEmitter时不要用0L无限等待。根据业务场景设置比如SseEmitter emitter new SseEmitter(30_000L); // 30秒。超时后会自动触发onTimeout回调进行清理。强制清理除了依赖超时和回调还应该提供一个主动清理的机制。例如在用户退出登录、关闭浏览器标签时前端发送一个断开连接的请求后端调用emitter.complete()。也可以设置一个后台定时任务定期扫描emitterMap清理掉长时间没有通信比如最后一次发送时间超过5分钟的emitter。使用WeakReference或缓存框架对于更复杂的场景可以考虑使用WeakHashMap或者集成类似Caffeine这样的缓存并设置基于时间和基于引用的过期策略让GC能帮忙回收失效的连接。5.2 背压与流量控制当AI生成速度极快或者网络较慢时可能会出现服务器推送数据的速度远大于客户端接收/处理速度的情况。这会导致数据在服务器端堆积最终可能引发内存问题。解决方案SseEmitter的发送队列SseEmitter.send()方法是非阻塞的它会将事件放入一个内部队列。如果客户端处理慢这个队列会变长。你需要监控这个情况。使用反应式编程模型这是更根本的解决方案。考虑使用Spring WebFlux将整个链路从消费AI API流到推送SSE流都构建在Reactor的Flux之上。Reactor内置了背压机制下游客户端可以向上游AI服务请求特定数量的数据从而实现流量控制。简单限流在服务端代码中可以在每次send之后加入一个小的延迟或者检查emitter的队列状态虽然SseEmitter没有直接暴露此API来控制推送频率。5.3 异常处理与连接恢复网络是不稳定的。SSE连接可能因为网络抖动、代理超时、服务器重启等原因中断。前端EventSource对象在连接断开后默认会自动尝试重新连接。重连时浏览器会自动在请求头中带上上次收到的最后一个事件的IDLast-Event-ID。后端我们需要处理这个Last-Event-ID头以便从断点处恢复数据推送。这要求我们在发送每个事件时都设置一个有序且唯一的id并在服务端有能力根据这个ID定位到中断时的数据位置例如记录每个用户会话的生成进度。对于AI生成这种一次性过程断点续传实现较复杂通常更简单的做法是在重连后如果发现之前的生成任务未完成则重新开始一个新的生成任务并通知前端“连接已恢复重新生成”。对于聊天对话场景则需要更精细的状态管理。实操建议至少要做到友好的错误提示。在onError回调中不仅记录日志还要尝试发送一个最终的error事件给前端让用户知道发生了什么而不是默默地断开。5.4 认证与授权集成在生产环境中SSE端点不能对所有人开放。你需要集成Spring Security等安全框架。挑战标准的SSEEventSource在建立连接时只能发送标准的HTTP头如Authorization: Bearer token。它不支持在连接建立后动态修改头部也不支持像WebSocket那样在协议升级时进行复杂握手。解决方案Token作为查询参数在连接URL中携带认证token如/sse/connect?clientIdxxxtokenyyy。注意这存在token在日志、浏览器历史中泄露的风险需确保使用HTTPS并谨慎记录日志。Cookie如果应用使用基于Cookie的会话EventSource会自动携带Cookie可以利用Session进行认证。先认证后连接前端先调用一个普通的REST接口进行登录认证获取一个短期有效的“连接令牌”connection token。然后使用这个令牌作为参数或自定义Header需服务器支持CORS预检来建立SSE连接。后端验证此令牌的有效性。使用Spring Security的事件流支持Spring Security 5 对Server-Sent Events有更好的支持可以结合其安全上下文。在控制器中你可以像保护普通RestController端点一样使用PreAuthorize等注解。只需注意认证逻辑会在连接建立时执行一次。5.5 性能与可扩展性当并发用户数很高时每个SSE连接都会占用一个服务器线程如果使用Tomcat等传统Servlet容器。虽然这些线程大部分时间处于休眠等待发送事件状态但数量过多仍会影响服务器性能。优化方向使用WebFlux和Netty将应用迁移到Spring WebFlux基于Netty。Netty使用事件循环和非阻塞IO可以用少量线程处理大量并发连接非常适合SSE这种长连接、低频率数据推送的场景。网关聚合在微服务架构下可以考虑使用API网关如Spring Cloud Gateway来统一处理SSE连接后端服务只负责生成事件数据通过消息队列如Kafka, RabbitMQ推送给网关由网关分发给对应的客户端。这样可以解耦和水平扩展。连接数监控务必监控应用中的活跃SSE连接数设置告警阈值以便提前发现容量问题。6. 常见问题排查与调试技巧在实际开发中你肯定会遇到各种奇怪的问题。这里记录几个我踩过的坑和解决方法。问题1前端收不到任何消息连接状态一直是CONNECTING(0)。检查后端响应头确保控制器方法的produces MediaType.TEXT_EVENT_STREAM_VALUE生效并且响应头确实是Content-Type: text/event-stream;charsetUTF-8。可以用Postman或curl测试接口。检查CORS如果前端和后端域名不同浏览器会因同源策略阻止SSE连接。必须在后端配置CORS允许前端域名并暴露必要的头如Cache-Control。Spring Boot中可以在配置类或控制器上添加CrossOrigin注解。Configuration public class WebConfig implements WebMvcConfigurer { Override public void addCorsMappings(CorsRegistry registry) { registry.addMapping(/api/sse/**) .allowedOrigins(http://your-frontend-domain.com) .allowedMethods(GET) // SSE通常是GET请求 .allowCredentials(true) // 如果需要携带cookie .exposedHeaders(Cache-Control); // 可能需要暴露 } }问题2连接建立后很快就自动断开了前端触发onerror。检查超时设置后端创建的SseEmitter超时时间太短。生产环境网络延迟或AI生成较慢时容易超时。适当调大超时时间但不要设为0。检查心跳有些代理服务器如Nginx或负载均衡器对长时间空闲的连接有默认超时例如60秒。解决方案是服务器端定期发送“心跳”事件比如一个注释行:\n\n或一个空的data事件保持连接活跃。// 可以启动一个定时任务每隔一段时间向所有活跃的emitter发送心跳 Scheduled(fixedDelay 30000) // 每30秒一次 public void sendHeartbeat() { emitterMap.forEach((clientId, emitter) - { try { emitter.send(SseEmitter.event().comment(heartbeat)); } catch (IOException e) { // 发送失败连接可能已失效可以考虑移除 log.debug(发送心跳到客户端 {} 失败, clientId); } }); }检查服务器/代理配置确保Nginx等代理服务器配置了较长的proxy_read_timeout并且没有缓冲响应proxy_buffering off;对于SSE至关重要否则数据会被缓存直到连接关闭。问题3前端能收到消息但EventSource的onmessage触发了addEventListener却没反应。事件名匹配检查后端发送事件时指定的name如.name(chunk)是否和前端的addEventListener(chunk, ...)中的事件名完全一致大小写敏感。默认onmessage如果没有指定event字段或者event字段为空数据会触发onmessage回调。如果指定了event则只会触发对应addEventListener的回调而不会触发onmessage。建议统一使用addEventListener来监听特定事件。问题4高并发下服务器出现大量IOException: Broken pipe或Connection reset by peer。客户端主动断开这是正常现象用户关闭了浏览器标签。关键是要在后端的onError或onCompletion回调中做好资源清理避免内存泄漏。网络不稳定在不可靠的网络环境下这是常态。确保你的异常处理代码足够健壮不会因为一个连接异常而影响整个服务。监控与日志对这些错误进行适当的监控和日志记录但不必视为严重错误。可以区分日志级别大量的Broken pipe可以记录为DEBUG或WARN而非ERROR。调试技巧使用浏览器开发者工具在Network标签页查看SSE连接类型为eventsource可以看到请求头、响应头以及实时流入的数据流。使用curl命令行测试curl -N http://localhost:8080/api/sse/connect?clientIdtest。-N参数表示不缓冲直接输出原始流非常适合调试SSE接口。服务端日志增强在SseEmitter的所有回调onCompletion,onTimeout,onError以及每次send操作前后添加详细的日志记录客户端ID和关键状态便于追踪连接生命周期。流式输出为AI应用带来了质的体验提升而SSE是在Java Web生态中实现这一特性的利器。从理解协议原理到在Spring Boot中一步步构建控制器和服务再到处理生产环境中的各种疑难杂症整个过程需要我们对HTTP、并发编程、资源管理都有清晰的认识。希望这篇从原理到实战的长文能帮你绕过我踩过的那些坑顺利地在你的Java AI应用中实现流畅的“打字机”效果。记住良好的异常处理和资源管理是这类长连接应用稳定的基石多测试、多监控才能让功能平稳上线。