Spring Boot实现LLM流式交互:SSE与SseEmitter实战指南
1. 项目概述为什么要在Spring Boot里搞流式交互最近在搞大模型LLM应用落地的朋友估计都遇到过同一个头疼的问题用户问了个稍微复杂点的问题后台吭哧吭哧算了十几秒前端页面就跟卡死了一样用户只能对着一个空白的输入框干瞪眼心里直犯嘀咕“这玩意儿是不是挂了”。这种糟糕的体验在传统的“请求-响应”同步模式下几乎无解。用户提交问题服务器调用LLM API等模型全部生成完毕再把一整段文本塞回给前端——这个过程里用户就是被动的等待者。流式交互Streaming Interaction就是为了干掉这个“等待黑洞”而生的。它的核心思想很简单别等模型全想好了再说想到哪就说到哪。就像两个人聊天对方说一句你听一句边听边理解而不是等对方写完一篇小作文再念给你听。在技术实现上这就意味着服务器需要有能力把LLM生成的一个个“词元”Token或一小段文本像流水一样持续不断地推送给客户端通常是浏览器。那么为什么是Spring Boot作为一个成熟的Java企业级开发框架Spring Boot以其约定大于配置、快速构建微服务的能力著称。当AI能力需要与企业现有的Java技术栈比如用户系统、订单系统、数据库深度集成时用Spring Boot来构建LLM应用的后端就成了一个非常自然且稳健的选择。它提供了强大的依赖管理、自动配置和丰富的生态组件让我们能更专注于业务逻辑而不是底层通信的复杂性。在这个项目里我们的目标就是拆解在Spring Boot中实现LLM流式交互的每一块“积木”从原理到代码让你不仅能搭起来更能明白为什么这么搭。2. 核心原理拆解从HTTP到SSE的演进之路要理解流式交互得先看看我们平时用的HTTP是怎么“说话”的。经典的HTTP 1.1协议遵循的是“一问一答”模式。客户端发一个Request服务器处理完回一个完整的Response然后连接就关闭了。这就像你打电话订餐说完菜单和地址对方说“好的45分钟后送到”电话就挂了。在这45分钟里你完全不知道后厨做到哪一步了。对于LLM这种生成过程较长的任务这种模式显然不行。于是我们需要一种能让服务器主动、持续向客户端发送数据的机制。这里主要有几种技术路径2.1 WebSocket全双工通信的利剑WebSocket在浏览器和服务器之间建立一个持久化的、全双工双方可以同时收发的通道。它功能强大适合需要高频、双向实时交互的场景比如在线游戏、协同编辑。但对于LLM流式输出这种典型的“服务器单向推送客户端主要接收”的场景用WebSocket有点“杀鸡用牛刀”。它引入了更复杂的连接管理、心跳维持和协议处理增加了不必要的复杂度。2.2 Server-Sent Events (SSE)为单向流式而生SSE是HTML5规范的一部分它允许服务器通过一个持久的HTTP连接主动向客户端推送事件流。它的特点非常契合我们的需求单向性服务器到客户端的单向推送。这正是LLM流式输出的核心模式。基于HTTP本质上还是HTTP协议可以利用现有的基础设施如负载均衡器、防火墙规则兼容性更好。简单轻量协议简单浏览器端有原生的EventSource对象支持开发成本低。自动重连EventSource内置了连接断开后的重连机制。SSE的数据格式有特定规范一个标准的响应看起来是这样的HTTP/1.1 200 OK Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive data: {token: 你好} data: {token: } data: {token: 今天} event: complete data: {status: done}每一段数据以data:开头两个换行符(\n\n)标识一个事件的结束。还可以通过event:字段定义自定义事件类型。2.3 Spring Boot的选择SseEmitterSpring Framework从4.2版本开始为SSE提供了优雅的抽象——SseEmitter。它本质上是一个响应式编程模型的产物将HTTP响应的输出流包装成一个Emitter允许我们在不同的线程中异步地向这个流中写入数据。SseEmitter帮我们处理了底层的HTTP连接管理、超时控制以及SSE格式的封装让我们可以像操作一个普通的“发送器”一样专注于业务数据的生产与推送。所以技术选型的逻辑链条很清晰为了实现LLM的流式输出服务器持续推送- 选用更适合单向推送的SSE而非WebSocket - 在Spring Boot中使用SseEmitter作为实现SSE的核心工具。3. 架构设计与核心组件一个健壮的Spring Boot LLM流式交互后端不能只是一个简单的“转发器”。它需要妥善处理并发、异步、资源管理和异常。下面是一个典型的架构分层3.1 控制器层Controller连接的发起与维护这一层由RestController中的接口负责。它的核心职责是创建并返回SseEmitter对象在客户端如浏览器调用连接接口时即时创建一个SseEmitter实例通常我们会设置一个合理的超时时间例如new SseEmitter(30_000L)表示30秒。注册生命周期回调这是关键。需要为SseEmitter注册onCompletion完成和onTimeout超时回调。在这些回调里我们必须进行关键的资源清理工作比如将当前连接从全局的管理器中移除。如果不做清理会导致内存泄漏。存储连接上下文将新创建的SseEmitter与一个唯一的会话ID如UUID关联并存储到一个全局的并发安全的容器中如ConcurrentHashMapString, SseEmitter。这样后续的异步任务才能找到正确的连接来推送数据。触发异步处理连接建立后控制器不应阻塞。它应立即将具体的业务处理如调用LLM API提交给一个异步执行器Async方法或TaskExecutor并将SseEmitter和请求参数传递给这个异步任务。3.2 异步服务层Async Service业务逻辑的核心这是真正“干活”的地方通常由Service组件承载并且方法被Async注解标记。它的工作流是准备LLM调用组装请求参数调用LLM服务提供商如OpenAI、通义千问、DeepSeek等的API。关键点在于必须调用其支持流式输出的接口。这些API通常会返回一个流式响应对象如Spring的ResponseEntityFluxString或OpenAI Java SDK中的Stream。消费流式响应遍历或订阅这个流式响应。每收到一个数据块Chunk就将其封装成前端约定好的格式通常是JSON。通过SseEmitter发送调用SseEmitter.send()方法将封装好的数据块发送出去。这里必须做好异常捕获因为连接可能在任何时候被客户端关闭。发送完成事件当流式响应结束时LLM生成完毕发送一个特殊的事件如event: complete通知前端可以关闭连接或更新UI状态。异常处理与资源释放在任何步骤发生错误如LLM API调用失败、网络中断、JSON解析错误都需要捕获异常并尝试发送一个错误事件给前端最后在finally块中确保调用SseEmitter.complete()或completeWithError()来显式关闭连接触发控制器层的清理回调。3.3 连接管理层全局管理器这是一个隐形的但至关重要的层。我们需要一个中心化的地方来管理所有活跃的SseEmitter连接。它的功能包括注册连接在控制器创建SseEmitter时存入。查找连接供异步服务层根据会话ID查找对应的SseEmitter。移除连接在连接完成、超时或出错时从管理器中移除防止内存泄漏。广播能力如果需要实现类似群聊的功能管理器还可以提供向所有连接广播消息的方法。3.4 配置层异步与线程池Spring Boot的异步能力默认使用一个简单的线程池。但在生产环境中我们必须自定义线程池以更好地控制资源。Configuration EnableAsync public class AsyncConfig { Bean(llmStreamingTaskExecutor) public TaskExecutor taskExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); // 核心线程数即使空闲也保留的线程数根据服务器资源和预期并发量设置 executor.setCorePoolSize(10); // 最大线程数队列满后能创建的最大线程数 executor.setMaxPoolSize(50); // 队列容量核心线程忙时新任务进入队列等待 executor.setQueueCapacity(100); // 线程名前缀便于日志排查 executor.setThreadNamePrefix(llm-stream-); // 拒绝策略当线程池和队列都满时如何拒绝新任务。CallerRunsPolicy表示由调用者线程执行 executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; } }然后在异步方法上指定使用这个执行器Async(llmStreamingTaskExecutor)。合理的线程池配置是系统稳定性的基石能避免因突发流量导致的线程资源耗尽。4. 完整实现步骤与代码剖析我们以一个调用 OpenAI 兼容 API 的流式聊天接口为例将上述架构落地。4.1 第一步定义连接管理器Component public class SseConnectionManager { private final MapString, SseEmitter emitterMap new ConcurrentHashMap(); public SseEmitter createEmitter(String connectionId, Long timeout) { // 设置超时时间0表示永不超时不推荐 SseEmitter emitter timeout ! null ? new SseEmitter(timeout) : new SseEmitter(); this.emitterMap.put(connectionId, emitter); // 设置完成回调用于资源清理 emitter.onCompletion(() - { log.info(SSE连接完成: {}, connectionId); this.emitterMap.remove(connectionId); }); // 设置超时回调 emitter.onTimeout(() - { log.warn(SSE连接超时: {}, connectionId); emitter.complete(); this.emitterMap.remove(connectionId); }); // 设置错误回调 emitter.onError((ex) - { log.error(SSE连接错误: {}, connectionId, ex); this.emitterMap.remove(connectionId); }); return emitter; } public SseEmitter getEmitter(String connectionId) { return emitterMap.get(connectionId); } public void removeEmitter(String connectionId) { emitterMap.remove(connectionId); } }注意onCompletion和onTimeout回调是互斥的只会触发一个。确保在回调中移除连接是防止内存泄漏的关键。4.2 第二步实现控制器RestController RequestMapping(/api/chat) Slf4j public class StreamChatController { Autowired private SseConnectionManager connectionManager; Autowired private LlmStreamingService llmStreamingService; GetMapping(path /connect, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter connect(RequestParam String sessionId) { // 生成或使用传入的会话ID作为连接标识 String connectionId StringUtils.hasText(sessionId) ? sessionId : UUID.randomUUID().toString(); log.info(建立SSE连接连接ID: {}, connectionId); // 创建Emitter设置超时时间例如5分钟应对长文本生成 SseEmitter emitter connectionManager.createEmitter(connectionId, 300_000L); // 可以立即发送一个连接成功的事件 try { SseEmitter.SseEventBuilder event SseEmitter.event() .name(connect) // 自定义事件名 .data(Map.of(connectionId, connectionId, message, 连接已建立)); emitter.send(event); } catch (IOException e) { log.error(初始消息发送失败, e); emitter.completeWithError(e); } return emitter; // 此时连接已建立Spring会保持这个HTTP连接 } PostMapping(/stream) public ResponseEntity? sendMessage(RequestBody ChatRequest request, RequestParam String connectionId) { SseEmitter emitter connectionManager.getEmitter(connectionId); if (emitter null) { return ResponseEntity.status(HttpStatus.GONE).body(连接不存在或已关闭); } // 异步处理消息不阻塞当前请求线程 llmStreamingService.streamLlmResponse(request, connectionId, emitter); return ResponseEntity.accepted().body(Map.of(status, processing, connectionId, connectionId)); } }这里拆成了两个接口/connect用于建立SSE长连接/stream用于接收用户消息并触发异步处理。这种分离更符合RESTful风格也便于前端管理连接和发送请求。4.3 第三步实现异步流式服务这是最核心的部分我们以使用Spring的WebClient调用OpenAI API为例。Service Slf4j public class LlmStreamingService { Autowired private SseConnectionManager connectionManager; Value(${llm.api.base-url}) private String apiBaseUrl; Value(${llm.api.key}) private String apiKey; Async(llmStreamingTaskExecutor) // 指定自定义线程池 public void streamLlmResponse(ChatRequest request, String connectionId, SseEmitter emitter) { SseEmitter localEmitter emitter; // 可能从管理器重新获取这里简化 try { WebClient client WebClient.builder() .baseUrl(apiBaseUrl) .defaultHeader(Authorization, Bearer apiKey) .build(); // 构建流式请求体 MapString, Object requestBody Map.of( model, gpt-3.5-turbo, messages, request.getMessages(), stream, true // 关键开启流式 ); FluxString responseFlux client.post() .uri(/v1/chat/completions) .contentType(MediaType.APPLICATION_JSON) .bodyValue(requestBody) .retrieve() .bodyToFlux(String.class); // 以Flux流的形式接收响应 // 订阅并处理流 responseFlux.doOnNext(dataChunk - { // OpenAI流式响应格式以data: 开头的行最后是data: [DONE] if (dataChunk.startsWith(data: )) { String jsonData dataChunk.substring(6).trim(); if ([DONE].equals(jsonData)) { sendSseEvent(localEmitter, complete, Map.of(status, done)); return; } try { // 解析JSON提取生成的文本delta JsonNode node new ObjectMapper().readTree(jsonData); JsonNode choice node.path(choices).get(0); JsonNode delta choice.path(delta); String content delta.path(content).asText(null); if (content ! null !content.isEmpty()) { // 封装并发送给前端 MapString, Object eventData Map.of(token, content); sendSseEvent(localEmitter, message, eventData); } } catch (JsonProcessingException e) { log.warn(解析SSE数据块失败: {}, jsonData, e); } } }).doOnComplete(() - { log.info(流式响应处理完成 for {}, connectionId); // 发送完成事件确保前端收到结束信号 sendSseEvent(localEmitter, complete, Map.of(status, done)); }).doOnError(error - { log.error(处理LLM流时发生错误 for {}, connectionId, error); sendSseEvent(localEmitter, error, Map.of(message, 模型响应流异常)); localEmitter.completeWithError(error); }).subscribe(); // 订阅以启动流的消费 } catch (Exception e) { log.error(流式处理任务执行失败 for {}, connectionId, e); sendSseEvent(localEmitter, error, Map.of(message, 服务内部错误)); localEmitter.completeWithError(e); } } private void sendSseEvent(SseEmitter emitter, String eventName, Object data) { if (emitter null) return; try { SseEmitter.SseEventBuilder event SseEmitter.event() .name(eventName) .data(data, MediaType.APPLICATION_JSON); emitter.send(event); } catch (IOException e) { // 发送失败通常意味着客户端已断开连接 log.debug(向客户端发送SSE事件失败连接可能已关闭, e); // 这里可以选择不处理因为连接管理器的onError/onCompletion回调会负责清理 } } }实操心得在doOnNext中处理每个数据块时一定要做好异常捕获和空值判断。LLM API返回的JSON结构可能不稳定某个delta里可能没有content字段例如在发送角色信息时。忽略解析错误避免因单个数据块问题导致整个流中断。4.4 第四步前端简单示例前端使用EventSourceAPI进行连接和监听。class SSEClient { constructor(apiBaseUrl) { this.apiBaseUrl apiBaseUrl; this.eventSource null; this.connectionId null; } connect() { // 建立连接 const url ${this.apiBaseUrl}/api/chat/connect; this.eventSource new EventSource(url); this.eventSource.addEventListener(connect, (e) { const data JSON.parse(e.data); this.connectionId data.connectionId; console.log(连接成功ID:, this.connectionId); }); this.eventSource.addEventListener(message, (e) { const data JSON.parse(e.data); // 将收到的token实时追加到UI上 document.getElementById(output).innerText data.token; }); this.eventSource.addEventListener(complete, (e) { console.log(流式响应结束); // 可以更新UI状态如将发送按钮置为可用 }); this.eventSource.addEventListener(error, (e) { console.error(SSE错误:, e); this.close(); }); this.eventSource.onerror (err) { // 处理连接错误 console.error(EventSource连接错误:, err); this.close(); }; } sendMessage(message) { if (!this.connectionId) { alert(请先建立连接); return; } const url ${this.apiBaseUrl}/api/chat/stream?connectionId${this.connectionId}; fetch(url, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ messages: [{ role: user, content: message }] }) }).then(response { if (!response.ok) { throw new Error(发送消息失败); } console.log(消息已发送处理中...); }); } close() { if (this.eventSource) { this.eventSource.close(); this.eventSource null; this.connectionId null; } } }5. 生产环境关键考量与优化策略把Demo跑起来只是第一步要上线稳定运行还有一堆“坑”等着填。5.1 连接管理与资源泄漏防御这是流式服务最核心的稳定性问题。除了在SseEmitter回调中清理还必须有一个兜底机制。定期健康检查与清理启动一个定时任务遍历连接管理器中的所有SseEmitter检查其状态。对于创建时间过久如超过1小时的连接强制将其complete()并移除。这可以清理掉因为客户端异常关闭如直接关闭浏览器标签而未能触发正常回调的“僵尸连接”。使用WeakReference不推荐。SseEmitter与HTTP响应流强绑定需要显式生命周期管理弱引用会导致不可控的清理时机。5.2 背压Backpressure处理当LLM生成速度远快于网络发送速度或者前端处理能力不足时数据会在服务器端堆积可能导致内存溢出OOM。SseEmitter.send()方法是阻塞的如果客户端接收慢它会一直等待直到发送成功或超时。更高级的做法是结合响应式编程的背压机制。我们可以使用Flux和SseEmitter的结合但需要更精细的控制。一种实践是使用一个具有边界的队列如BlockingQueue作为缓冲区异步任务将数据放入队列另一个线程从队列取出并发送。当队列满时可以采取丢弃最新数据或等待的策略。5.3 超时与重连策略服务器超时SseEmitter的超时时间不宜过短避免长文本生成中断也不宜过长避免资源被无效连接长期占用。可以设置为5-10分钟并在前端配合心跳机制。客户端重连EventSource有自动重连但重连后connectionId会变需要重新建立会话状态。更复杂的方案是前端在连接断开后携带之前的sessionId主动调用/connect接口进行“续连”后端需要能恢复之前的上下文如果有的话。5.4 上下文管理与会话状态在多次流式交互中多轮对话需要维护会话上下文。简单的做法是将每轮对话的用户消息和AI的流式回复都追加到一个内存或Redis中的列表里。当新的请求到来时携带sessionId后端取出历史记录组装成完整的消息列表再发给LLM。注意LLM通常有上下文长度限制需要实现一个“滑动窗口”或总结机制来管理过长的历史。5.5 监控与可观测性需要监控的关键指标包括活跃连接数反映当前系统负载。连接创建/关闭速率帮助发现异常流量。消息发送延迟从LLM返回Token到成功推送给客户端的耗时。错误率连接错误、发送失败、LLM API调用失败的比例。 这些指标可以通过Spring Boot Actuator、Micrometer集成到Prometheus和Grafana中实现可视化监控。6. 常见问题排查与实战技巧在实际开发和运维中你会遇到各种各样的问题。下面是一些典型场景和解决思路。6.1 前端收不到数据或连接立即关闭检查响应头确保服务器响应的Content-Type是text/event-stream并且Cache-Control设置为no-cache。Spring Boot的SseEmitter通常会自动设置这些。检查网络代理Nginx等反向代理默认可能缓冲整个响应。需要在代理配置中为SSE路径禁用代理缓冲location /api/chat/ { proxy_pass http://backend; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; proxy_buffering off; # 关键 proxy_cache off; }检查CORS如果前端与后端域名不同需要正确配置CORS允许EventSource请求。注意EventSource不支持自定义Header进行鉴权通常需要将会话信息放在URL参数中。6.2 流式输出中断或不完整服务器端超时检查SseEmitter的超时设置是否足够长。对于生成长篇内容可能需要延长。客户端EventSource自动重连当网络波动导致连接断开EventSource会尝试重连。如果后端没有为同一个会话提供“续连”能力重连后会得到一个新的空连接导致输出中断。需要实现会话恢复逻辑。服务器端资源耗尽检查线程池是否被打满或者队列是否溢出。观察监控指标调整线程池参数。6.3 内存占用过高未释放的SseEmitter这是最常见的原因。务必确保所有路径正常完成、超时、异常都能触发连接管理器的清理逻辑。使用jmap或VisualVM等工具定期检查SseEmitter实例的数量是否与活跃连接数匹配。大消息缓冲区如果单次推送的数据块很大或者背压导致数据在内存队列中大量堆积也会引起内存问题。优化数据块大小实现背压控制。6.4 与特定LLM API的兼容性问题不同厂商的流式API返回格式可能有细微差别。OpenAI格式如上述代码所示以data:为前缀[DONE]结尾。其他API可能是纯JSON流每行一个JSON对象或者使用不同的分隔符。在doOnNext逻辑中需要根据实际的API文档进行解析适配。务必在服务层做好抽象将不同API的流式响应适配成统一的内部数据格式这样业务逻辑代码就不需要关心具体是调用的哪家模型。6.5 异步上下文传递在异步线程中ThreadLocal存储的信息如Spring Security的认证信息、MDC日志跟踪ID会丢失。如果需要可以使用DelegatingSecurityContextAsyncTaskExecutor或手动使用TaskDecorator来包装任务传递必要的上下文。踩过这些坑之后我的体会是构建一个健壮的流式服务三分在功能实现七分在异常处理、资源管理和监控运维。它不是一个简单的“一发一收”接口而是一个有状态的、长生命周期的数据管道。每一个环节的稳健性都直接影响到终端用户能否获得流畅、稳定的AI交互体验。从SseEmitter开始但绝不能止步于此。