尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Spring Boot流式响应实战:SseEmitter与WebFlux实现LLM逐字输出

Spring Boot流式响应实战:SseEmitter与WebFlux实现LLM逐字输出 1. 项目概述为什么要在Spring Boot里搞流式响应最近在折腾大语言模型LLM应用落地的朋友估计都遇到过同一个头疼的问题用户问了个问题前端页面就转圈圈等个十几二十秒模型才慢悠悠地把一整段答案吐出来。这个等待过程用户体验极差感觉像是在和一台老旧的打字机对话。尤其是在做聊天机器人、智能客服或者代码补全这类交互性强的应用时这种“阻塞式”的响应方式简直是灾难。这背后的原因很简单LLM生成文本本质上是“一个词一个词”往外蹦的。传统的HTTP请求-响应模型要求服务器必须生成完整的答案才能一次性返回给客户端。这就好比你要等厨师把整道菜做完才端上桌而不是先给你上前菜和汤。流式响应Streaming Response就是为了解决这个问题而生。它允许服务器将LLM生成的每一个“词元”Token实时地、分批地推送给客户端实现打字机式的逐字输出效果。在Spring Boot生态里实现这种流式响应主要有两条技术路线走得比较顺也是社区里公认的最佳实践一个是基于Servlet API的SseEmitter另一个是拥抱响应式编程范式的Spring WebFlux。选哪个这不是一个简单的二选一而是关乎你的技术栈、团队熟悉度和对并发模型的理解。接下来我就结合最近在项目里的实战把这两种方案的里里外外、坑坑洼洼都给你捋清楚。2. 技术选型深度解析SseEmitter vs. WebFlux在动手写代码之前我们得先弄明白手里这两把“武器”到底有什么区别适合什么样的战场。盲目选型后面可能就是无穷无尽的调试和重构。2.1 SseEmitter传统同步架构的轻量级升级SseEmitter是Spring MVC提供的一个用于服务器发送事件Server-Sent Events, SSE的利器。它的核心思想非常简单在传统的Servlet同步阻塞模型上开一个“长连接”通道允许服务器通过这个通道主动、多次地向客户端发送数据。它的工作模式是这样的客户端通常是浏览器通过EventSource API发起一个GET请求到特定端点。服务器端Controller方法返回一个SseEmitter对象。服务器在另一个线程比如通过Async或线程池中调用LLM服务每获得一个词元就通过emitter.send()方法发送出去。客户端通过监听onmessage事件实时接收并渲染这些数据块。当LLM生成完毕或发生错误时服务器调用emitter.complete()或emitter.completeWithError()来结束流。它的优势非常明显学习成本极低如果你已经很熟悉Spring MVC那套RestController、GetMapping的写法那么使用SseEmitter几乎不需要学习新概念。它完美地集成在现有的Spring MVC框架内。对现有项目侵入性小你不需要改变项目的整体架构只需要在需要流式输出的接口上进行改造即可。这对于在已有大型单体应用中快速增加AI功能模块的场景非常友好。客户端支持广泛SSE是一个W3C标准现代浏览器原生支持前端用起来很简单。对于非浏览器客户端也有成熟的库可以处理SSE流。但它也有天生的局限基于Servlet的阻塞IO模型虽然SseEmitter本身实现了异步发送但其底层依然是Tomcat等Servlet容器的阻塞IO线程模型。每个SSE连接都会占用一个Servlet容器的工作线程如Tomcat的http-nio线程。在高并发流式请求的场景下大量连接可能快速耗尽线程池导致服务器无法处理其他普通请求。单向通信SSE是严格的服务器向客户端的单向通信。如果你想在流式输出过程中允许用户中途打断比如发送一个“停止生成”的指令就需要借助额外的WebSocket或普通的HTTP接口来实现架构上会变得复杂。实操心得如果你的应用并发压力不大比如内部工具、后台管理系统或者团队对响应式编程完全陌生那么SseEmitter是快速上手的首选。它能用最小的改动带来立竿见影的体验提升。2.2 Spring WebFlux为高并发流式而生的响应式方案Spring WebFlux是Spring 5引入的、全新的非阻塞、响应式Web框架。它的核心是Project Reactor库提供了Flux和Mono这两种代表异步数据流的类型。当你的Controller方法返回一个FluxString时Spring WebFlux会自动将其处理为流式HTTP响应通常是text/event-stream格式。它的工作模式是响应式的客户端发起请求。Controller方法直接返回一个FluxString流。这个流的源头Publisher是你的LLM服务它应该以非阻塞的方式生成数据。WebFlux框架底层通常是Netty负责订阅这个Flux并将每个产生的元素实时写入HTTP响应体。整个过程在事件循环Event Loop中完成不会阻塞任何工作线程。它的优势在于高并发和资源效率真正的非阻塞IO从网络IO到业务逻辑全程无阻塞。一个事件循环线程可以处理成千上万的并发连接特别适合LLM这种长耗时、高并发的流式场景。资源利用率远高于线程池模型。背压Backpressure支持这是响应式编程的精髓之一。如果客户端处理速度慢比如网络差或前端渲染卡顿Flux可以感知并通知生产者LLM服务放慢生成速度避免服务器内存被积压的数据撑爆。这在生产环境中是至关重要的稳定性保障。统一的流处理模型Flux不仅可以用于HTTP响应还可以轻松地与响应式的数据库驱动如R2DBC、消息中间件如Kafka Reactive集成构建全链路的非阻塞应用。当然它的挑战也不小编程范式转变从命令式的、同步的思维切换到响应式的、函数式的思维有较高的学习曲线。错误处理、调试的难度也相应增加。库生态兼容性很多常用的Java库特别是涉及阻塞IO的如某些JDBC驱动、同步的HTTP客户端不能在响应式链中直接使用否则会破坏非阻塞性。你需要寻找其响应式版本或进行额外的封装。实操心得如果你的应用从零开始且明确需要面对海量用户同时进行流式对话的场景比如公开的AI聊天服务那么投入时间学习并使用WebFlux是值得的。它为未来的可伸缩性打下了基础。但对于已有庞大同步代码库的项目全面迁移到WebFlux的代价可能过高。选型决策速查表特性维度SseEmitter(Spring MVC)Spring WebFlux编程模型命令式、同步/异步混合声明式、响应式、非阻塞并发模型线程池一连接一线程事件循环多连接一线程资源消耗线程资源敏感连接数受限于线程池内存和CPU资源敏感可支持极高并发连接学习成本低基于传统Spring MVC高需要理解响应式编程范式项目改造侵入性小局部改造即可侵入性大通常需整体架构调整适用场景内部工具、并发量中等、快速迭代高并发公共服务、全新项目、需要全链路非阻塞客户端交互单向SSE单向HTTP Stream或双向WebSocket3. 基于SseEmitter的实现详解与避坑指南理论说再多不如一行代码。我们先来看看如何用SseEmitter实现一个可靠的LLM流式接口。我会从一个最简单的例子开始然后逐步加入超时控制、错误处理和连接管理这些生产级必备特性。3.1 基础实现从零搭建一个流式聊天接口首先确保你的Spring Boot项目版本在2.x以上3.x当然更好。我们假设你已经有了一个能进行流式生成的LLM服务比如调用了OpenAI的API或者本地部署的Ollama。这里我以一个模拟的LLM服务为例。1. 创建Controllerimport 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 SseChatController { // 模拟一个耗时的LLM流式生成服务 private final SimulatedLLMService llmService; // 使用一个线程池来处理异步任务避免阻塞Servlet容器线程 private final ExecutorService asyncExecutor Executors.newCachedThreadPool(); public SseChatController(SimulatedLLMService llmService) { this.llmService llmService; } GetMapping(value /stream, produces text/event-stream) public SseEmitter streamChat(RequestParam String message) { // 创建一个SseEmitter这里设置超时时间为30分钟因为LLM生成可能很慢 SseEmitter emitter new SseEmitter(30 * 60 * 1000L); // 将耗时的LLM调用提交到线程池执行 asyncExecutor.submit(() - { try { // 调用LLM服务传入一个回调函数用于发送每一个生成的片段 llmService.generateStream(message, chunk - { try { // 发送数据格式为SSE规定的 data: {content}\n\n emitter.send(SseEmitter.event() .data(chunk) // 内容 .id(UUID.randomUUID().toString()) // 可选事件ID用于断线重连 .comment(a chunk of text) // 可选注释 ); } catch (IOException e) { // 发送失败通常意味着客户端已断开连接 // 这里可以选择中断LLM生成避免浪费资源 throw new RuntimeException(Client disconnected, e); } }); // LLM生成完成发送结束事件 emitter.complete(); } catch (Exception e) { // 发生任何错误通知客户端 emitter.completeWithError(e); } }); // 设置连接结束时的回调用于资源清理 emitter.onCompletion(() - System.out.println(SSE connection completed.)); emitter.onTimeout(() - System.out.println(SSE connection timed out.)); emitter.onError((ex) - System.out.println(SSE connection error: ex.getMessage())); return emitter; } }2. 模拟的LLM服务import java.util.function.Consumer; Component public class SimulatedLLMService { public void generateStream(String prompt, ConsumerString chunkConsumer) { String simulatedResponse 这是一个模拟的流式响应它会逐词返回。; String[] words simulatedResponse.split(); // 按字分割模拟Token for (String word : words) { try { // 模拟每个词元生成需要一些时间 Thread.sleep(100); // 将生成的词元通过回调函数发送出去 chunkConsumer.accept(word); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } // 模拟生成结束信号 chunkConsumer.accept([DONE]); } }3. 前端页面简单示例!DOCTYPE html html body input typetext idinput placeholder输入你的问题 button onclickstartStream()发送/button div idoutput stylewhite-space: pre-wrap; border:1px solid #ccc; min-height:200px;/div script function startStream() { const message document.getElementById(input).value; const outputDiv document.getElementById(output); outputDiv.innerHTML ; // 清空之前的内容 // 使用EventSource连接SSE端点 const eventSource new EventSource(/api/chat/stream?message${encodeURIComponent(message)}); eventSource.onmessage function(event) { const data event.data; if (data [DONE]) { eventSource.close(); console.log(Stream finished.); return; } // 逐字追加到页面上 outputDiv.innerHTML data; // 自动滚动到底部 outputDiv.scrollTop outputDiv.scrollHeight; }; eventSource.onerror function(error) { console.error(EventSource failed:, error); eventSource.close(); outputDiv.innerHTML \n\n[连接出错]; }; } /script /body /html这样一个最基础的流式聊天功能就完成了。前端会看到文字一个一个地“打”出来。3.2 生产级加固超时、错误与连接管理上面的基础版有很多问题直接上生产环境肯定会崩。我们需要重点解决以下几个痛点痛点一线程池管理不当导致资源泄露我们上面用了Executors.newCachedThreadPool()它创建的线程池没有大小限制可能会创建大量线程。在Spring Boot中更推荐使用ThreadPoolTaskExecutor进行配置化管理。解决方案配置一个专用的异步任务执行器。Configuration EnableAsync // 启用异步支持 public class AsyncConfig { Bean(sseTaskExecutor) public Executor taskExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); // 核心线程数 executor.setMaxPoolSize(20); // 最大线程数 executor.setQueueCapacity(100); // 队列容量 executor.setThreadNamePrefix(sse-async-); executor.initialize(); return executor; } } // 在Controller中注入并使用 RestController public class SseChatController { private final AsyncTaskExecutor asyncTaskExecutor; // 注入 GetMapping(/stream) public SseEmitter streamChat(RequestParam String message) { SseEmitter emitter new SseEmitter(30 * 60 * 1000L); // 使用Async注解的方法并指定执行器 CompletableFuture.runAsync(() - { // ... LLM调用逻辑 }, asyncTaskExecutor).exceptionally(ex - { emitter.completeWithError(ex); return null; }); return emitter; } }痛点二客户端异常断开服务器仍在生成浪费资源这是SseEmitter最经典的坑。当用户关闭浏览器标签页时连接断开但服务器端的线程还在傻傻地调用LLM API既浪费token又消耗服务器资源。解决方案在SseEmitter的回调和LLM生成循环中增加中断检查。public SseEmitter streamChat(RequestParam String message) { SseEmitter emitter new SseEmitter(30 * 60 * 1000L); // 使用一个原子布尔值标记连接状态 AtomicBoolean isConnectionAlive new AtomicBoolean(true); emitter.onCompletion(() - { isConnectionAlive.set(false); System.out.println(连接正常结束); }); emitter.onTimeout(() - { isConnectionAlive.set(false); System.out.println(连接超时); }); emitter.onError((ex) - { isConnectionAlive.set(false); System.out.println(连接出错: ex.getMessage()); }); asyncTaskExecutor.execute(() - { try { llmService.generateStream(message, chunk - { // 每次发送前检查连接是否还活着 if (!isConnectionAlive.get()) { throw new RuntimeException(连接已断开停止生成); } try { emitter.send(chunk); } catch (IOException e) { // 发送失败也意味着连接可能已断 isConnectionAlive.set(false); throw new RuntimeException(发送失败连接可能已断开, e); } }); if (isConnectionAlive.get()) { emitter.complete(); } } catch (Exception e) { if (isConnectionAlive.get()) { emitter.completeWithError(e); } } }); return emitter; }同时在你的SimulatedLLMService的生成循环里也需要能响应中断public void generateStream(String prompt, ConsumerString chunkConsumer) throws InterruptedException { String[] words prompt.split(); for (String word : words) { // 检查当前线程是否被中断由Future.cancel触发 if (Thread.currentThread().isInterrupted()) { throw new InterruptedException(生成任务被中断); } Thread.sleep(100); chunkConsumer.accept(word); } }痛点三超时时间设置不当SseEmitter的默认超时时间很短可能是30秒。对于生成长篇大论的LLM来说完全不够。但设置得过长如几小时又可能导致僵尸连接占用资源。一个折中的方案是设置一个合理的超时如10分钟并配合心跳机制。解决方案增加服务器端心跳保持连接活跃并检测僵尸连接。asyncTaskExecutor.execute(() - { try { // 启动一个心跳任务 ScheduledExecutorService scheduler Executors.newSingleThreadScheduledExecutor(); ScheduledFuture? heartbeatFuture scheduler.scheduleAtFixedRate(() - { if (isConnectionAlive.get()) { try { // 发送一个注释类型的心跳避免影响前端数据解析 emitter.send(SseEmitter.event().comment(heartbeat)); } catch (IOException e) { isConnectionAlive.set(false); scheduler.shutdown(); } } else { scheduler.shutdown(); } }, 0, 15, TimeUnit.SECONDS); // 每15秒发送一次心跳 // ... LLM生成逻辑 // LLM生成完成后取消心跳任务 heartbeatFuture.cancel(true); scheduler.shutdown(); } catch (Exception e) { // ... 错误处理 } });避坑总结使用SseEmitter核心就是管理好生命周期和资源。牢记三点1. 用可控的线程池执行异步任务2. 必须监听并处理onCompletion、onTimeout、onError及时释放资源和停止后台任务3. 对于超长任务考虑用心跳保活。做好这三点SseEmitter方案在中等并发下可以非常稳定。4. 基于Spring WebFlux的实现详解与核心技巧如果你决定拥抱响应式那么WebFlux会给你带来不一样的编程体验。它更优雅也更挑战思维习惯。我们来实现一个功能对等的流式聊天接口。4.1 基础实现使用Flux构建响应式流首先你需要将Spring Boot项目转换为WebFlux应用。如果你是新项目可以直接选择Spring Reactive Web依赖。如果是现有项目需要添加依赖并做相应调整。1. 添加依赖Mavendependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-webflux/artifactId /dependency2. 创建响应式Controller关键点在于你的LLM服务也需要是响应式的即返回FluxString。我们改造一下模拟的LLM服务。import org.springframework.web.bind.annotation.*; import reactor.core.publisher.Flux; import java.time.Duration; RestController RequestMapping(/api/reactive-chat) public class ReactiveChatController { private final ReactiveLLMService reactiveLLMService; public ReactiveChatController(ReactiveLLMService reactiveLLMService) { this.reactiveLLMService reactiveLLMService; } GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxString streamChat(RequestParam String message) { // 直接返回一个Flux流 return reactiveLLMService.generateStream(message) // 可以方便地添加背压策略、超时、错误处理等操作符 .timeout(Duration.ofMinutes(10)) // 设置10分钟超时 .onErrorResume(e - { // 错误处理例如返回一个错误信息然后结束流 return Flux.just([ERROR: e.getMessage() ]); }) // 在流结束时发送一个结束信号这是一个好习惯 .concatWithValues([DONE]); } }3. 创建响应式的LLM服务这里我们模拟一个异步生成词元的Flux。在实际项目中这可能对应着调用一个支持响应式的HTTP客户端如WebClient去请求外部LLM API。import org.springframework.stereotype.Service; import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; Service public class ReactiveLLMService { public FluxString generateStream(String prompt) { // 将同步阻塞的生成过程包装在Flux.create或Flux.generate中并切换到弹性调度器执行 return Flux.Stringcreate(sink - { // 模拟生成过程 String simulatedResponse 这是WebFlux实现的流式响应。; String[] words simulatedResponse.split(); for (String word : words) { try { Thread.sleep(100); // 模拟延迟 } catch (InterruptedException e) { sink.error(e); return; } sink.next(word); // 发射一个元素 } sink.complete(); // 生成完毕结束流 }) .subscribeOn(Schedulers.boundedElastic()); // 在弹性线程池上执行阻塞操作 // 注意如果LLM调用本身是异步非阻塞的如WebClient则不需要切换调度器。 } }前端代码可以复用之前的EventSource因为WebFlux的TEXT_EVENT_STREAM_VALUE输出也是标准的SSE格式。4.2 高级特性背压、熔断与链路追踪WebFlux的强大之处在于其丰富的操作符和与响应式生态的集成。下面看几个生产环境有用的技巧。技巧一实现背压控制假设你的LLM生成速度极快但客户端网络很慢数据会在服务器端堆积。我们可以使用onBackpressureBuffer操作符来定义一个缓冲策略。public FluxString streamChat(String message) { return reactiveLLMService.generateStream(message) .onBackpressureBuffer(50, // 缓冲区大小 BufferOverflowStrategy.DROP_OLDEST) // 缓冲区满时丢弃最老的数据 .doOnNext(chunk - log.debug(Emitting chunk: {}, chunk)) .doOnError(e - log.error(Error in stream, e)); }更高级的做法是如果你的LLM服务支持比如某些API可以暂停/继续你可以在Flux.create中利用sink.request(n)来动态控制上游的生成速度。技巧二与Resilience4j集成实现熔断在微服务架构中调用外部LLM API可能失败。我们可以使用Resilience4j为这个流式端点添加熔断器。import io.github.resilience4j.circuitbreaker.CircuitBreaker; import io.github.resilience4j.reactor.circuitbreaker.operator.CircuitBreakerOperator; Service public class ReactiveChatController { private final ReactiveLLMService llmService; private final CircuitBreaker circuitBreaker; public ReactiveChatController(ReactiveLLMService llmService, CircuitBreakerRegistry registry) { this.llmService llmService; this.circuitBreaker registry.circuitBreaker(llmService); } GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxString streamChat(RequestParam String message) { return llmService.generateStream(message) .transform(CircuitBreakerOperator.of(circuitBreaker)) // 应用熔断器 .onErrorResume(e - Flux.just([服务暂时不可用请稍后重试])); } }技巧三响应式链路追踪在分布式系统中追踪一个流式请求的完整生命周期很有挑战性。你可以利用reactor.core.publisher.Hooks或与Micrometer、Brave等集成。import reactor.core.publisher.SignalType; import reactor.core.publisher.Hooks; PostConstruct public void init() { // 为每个Flux/Mono的操作符添加跟踪信息方便调试 Hooks.onOperatorDebug(); } // 在流中手动添加跟踪点 public FluxString generateStreamWithTrace(String prompt, String traceId) { return Flux.deferContextual(ctx - { log.info(TraceId: {}, Starting stream generation for prompt: {}, traceId, prompt); return generateStream(prompt); }) .doOnEach(signal - { if (signal.isOnNext()) { log.debug(TraceId: {}, Emitted: {}, traceId, signal.get()); } else if (signal.isOnComplete()) { log.info(TraceId: {}, Stream completed, traceId); } else if (signal.isOnError()) { log.error(TraceId: {}, Stream error, traceId, signal.getThrowable()); } }); }核心技巧使用WebFlux时时刻记住“一切都是流”。将你的LLM调用、数据库查询、外部服务调用都建模为Flux或Mono。利用丰富的操作符map,filter,flatMap,timeout,retry,onBackpressureBuffer来组合和控制系统行为。调试时善用Hooks.onOperatorDebug()和log()操作符来观察数据流动。5. 实战问题排查与性能优化无论选择哪种方案上线后都会遇到各种稀奇古怪的问题。这里我整理了一份从实战中总结出来的问题排查清单和优化建议。5.1 常见问题速查表问题现象可能原因排查步骤与解决方案前端收不到数据或连接立即关闭1. 响应头Content-Type不正确。2.SseEmitter超时时间太短。3. 服务器端抛出未捕获的异常。1. 检查Controller方法produces属性是否为text/event-stream。2. 检查SseEmitter构造器传入的超时时间单位毫秒。3. 查看服务器日志确保异步任务中的异常被try-catch并调用emitter.completeWithError()。流式输出卡顿数据堆积一段时间后突然全部涌出1.SseEmitter服务器端缓冲区满。Tomcat等容器对响应体有缓冲区。2.通用网络延迟或浏览器渲染阻塞。3.WebFlux没有正确处理背压生产者速度远大于消费者。1. 尝试在发送数据后调用emitter.flush()强制刷新缓冲区但频繁flush影响性能。2. 检查前端代码确保onmessage事件处理函数执行很快没有同步阻塞操作。3. 在WebFlux中使用onBackpressureBuffer定义合理缓冲策略或优化LLM生成速度。高并发下服务器线程数飙升响应变慢SseEmitter典型问题Servlet容器线程池被大量SSE长连接占满。1. 监控线程池使用情况如Tomcat的threads-busy。2. 优化SseEmitter的异步任务执行器使用有界队列和合理的拒绝策略。3.考虑限流对/stream接口进行并发连接数限制。4.评估迁移至WebFlux。客户端断开后服务器后台任务仍在运行没有正确监听和处理SseEmitter的生命周期事件。1. 务必实现onCompletion,onTimeout,onError回调。2. 在这些回调中设置标志位如AtomicBoolean并在LLM生成循环中检查该标志位及时中断。Nginx/Apache等反向代理后流式失效代理服务器默认会缓冲整个响应后再转发给客户端。1.Nginx在location配置中添加proxy_buffering off;和proxy_cache off;。2.Apache设置ProxyRequests Off并考虑使用mod_proxy的disablereuseOn等参数或升级到支持HTTP/1.1分块传输的版本。流式响应被浏览器或CDN缓存缓存服务器或浏览器对text/event-stream内容类型进行了缓存。在响应头中添加Cache-Control: no-cache, no-transform和X-Accel-Buffering: no针对Nginx。WebFlux应用内存持续增长1. 存在内存泄漏如未取消的订阅。2. 背压失控缓冲区堆积大量未发送数据。1. 使用Profiler工具如VisualVM, YourKit分析内存堆转储查看Flux相关对象。2. 检查代码中是否有Flux.create创建的sink没有在适当时候调用complete/error。3. 强化背压控制为Flux添加onBackpressureBuffer或onBackpressureDrop策略并设置合理的缓冲区大小。5.2 性能优化建议连接复用与HTTP/2如果客户端和服务器都支持优先启用HTTP/2。HTTP/2的多路复用特性可以显著减少多个流式连接带来的开销。在Spring Boot中Tomcat和Netty都支持HTTP/2需要相应配置。数据格式优化压缩对于文本数据启用GZIP压缩可以节省大量带宽。但注意流式响应下压缩是逐块进行的可能增加延迟。需要权衡。在Spring Boot中可以通过server.compression.enabledtrue全局开启或使用ResponseBody装饰器针对特定端点处理。协议除了SSE也可以考虑使用WebSocket进行双向、低延迟的流式通信尤其适合需要频繁交互的场景。Spring Boot对WebSocket也有很好的支持。LLM服务调用优化连接池如果通过HTTP调用外部LLM API使用带有连接池的客户端如Apache HttpClient、OkHttp或WebClient并合理配置参数最大连接数、超时时间。超时设置为LLM API调用设置连接超时、读取超时和写入超时避免一个慢请求拖死整个线程池或事件循环。异步客户端在WebFlux方案中务必使用非阻塞的客户端如WebClient来调用外部服务否则会阻塞事件循环线程。监控与告警关键指标监控活跃SSE连接数、WebFlux的event-loop线程使用率、LLM API调用延迟和错误率。设置告警当活跃连接数接近线程池上限SseEmitter或异常断开率飙升时触发告警。链路追踪为每个流式请求分配唯一的Trace ID并贯穿整个调用链前端-网关-业务服务-LLM服务便于定位问题。6. 扩展思考超越基础流式实现了基础的流式输出后我们可以思考一些更进阶的场景这些往往能体现出一个AI应用的成熟度。场景一支持中途打断Stop Generation用户看到一半觉得答案不对想停止生成。这在SseEmitter的单向SSE中比较难实现通常需要额外提供一个REST API让前端发送一个“停止”请求服务器端根据Session或Request ID找到对应的生成任务并中断。而在WebFlux中可以更优雅地结合WebSocket实现双向通信服务器端可以监听客户端的“停止”消息。场景二流式输出格式化如Markdown、代码高亮LLM直接返回的是纯文本。我们可以在服务器端或前端进行格式化。更高效的做法是让LLM返回结构化的数据片段如JSON包含文本类型和内容服务器端流式传输这些片段前端根据类型实时渲染为Markdown、代码块等。这需要定义一套简单的协议。场景三与RAG检索增强生成结合在流式生成的过程中先显示检索到的参考文档片段再开始生成答案。这需要将“检索”和“生成”两个阶段都流式化。可以设计为先流式返回检索结果[RETRIEVED_DOC] ...然后返回一个分隔符[GENERATION_START]再流式返回生成的内容。这能给用户更强的掌控感和信任感。场景四多模态流式输出不仅是文本还有图片、音频。这通常需要更复杂的协议。一种方案是使用服务器发送事件SSE传输不同的event类型如event: text、event: image_url。前端根据事件类型进行不同的渲染处理。实现这些进阶功能无论是用SseEmitter还是WebFlux核心思想都是定义清晰的数据协议和管理好复杂的客户端状态。从简单的文本流开始逐步迭代是更稳妥的做法。
返回列表