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)关联,并存储到一个全局的并发安全的容器中(如ConcurrentHashMap<String, SseEmitter>)。这样,后续的异步任务才能找到正确的连接来推送数据。 - 触发异步处理:连接建立后,控制器不应阻塞。它应立即将具体的业务处理(如调用LLM API)提交给一个异步执行器(
@Async方法或TaskExecutor),并将SseEmitter和请求参数传递给这个异步任务。
3.2 异步服务层(Async Service):业务逻辑的核心这是真正“干活”的地方,通常由@Service组件承载,并且方法被@Async注解标记。它的工作流是:
- 准备LLM调用:组装请求参数,调用LLM服务提供商(如OpenAI、通义千问、DeepSeek等)的API。关键点在于,必须调用其支持流式输出的接口。这些API通常会返回一个流式响应对象(如Spring的
ResponseEntity<Flux<String>>或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 Map<String, 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(); // 构建流式请求体 Map<String, Object> requestBody = Map.of( "model", "gpt-3.5-turbo", "messages", request.getMessages(), "stream", true // 关键:开启流式 ); Flux<String> 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()) { // 封装并发送给前端 Map<String, 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开始,但绝不能止步于此。
