拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

LLM流式对话实战:Spring Boot+SSE实现打字机效果与abort中断

LLM流式对话实战:Spring Boot+SSE实现打字机效果与abort中断 1. 为什么“流式”是 LLM 对话产品的分水岭做过对话类产品的朋友大概率都有过这种体验用户点下发送按钮界面转圈转了七八秒然后“啪”一下整段答案全冒出来。功能上没毛病但用起来就是别扭——像对着一个沉默寡言的人说话对方憋半天才一次性回你一大段。而 ChatGPT 那种一个字一个字往外蹦的效果哪怕总耗时一样主观感受也快得多、自然得多。这中间的差别就是流式对话。我这次要聊的就是围绕LLM 流式对话做的一套前后端联接的完整实现。核心目标很明确让大模型的回答像打字机一样实时渲染到页面上而不是等全部生成完再一次性返回。技术栈上后端用Spring Boot承载业务与模型调用传输层用SSEServer-Sent Events做单向流式推送前端负责接收事件流并逐段拼接渲染同时配合abort能力让用户可以随时中断生成。这套东西适合谁如果你正在做对话机器人、AI 助手、知识库问答这类产品或者你手上已经有一个能跑通“一问一答”的接口但想把它升级成真正的流式体验那这篇内容基本可以照着抄。哪怕你之前没接触过 SSE只要会写基本的 Spring Boot 接口和前端请求也能跟下来。我会把选型理由、参数细节、踩过的坑都摊开讲尽量让你少走弯路。先说清楚一个前提流式对话不是“把接口改快一点”这么简单它牵扯到模型调用方式、传输协议、前端渲染节奏、连接生命周期管理四个层面的协同。任何一层没处理好用户看到的要么是卡顿要么是断流要么是中断后状态错乱。下面我按这四个层面拆开讲。2. 整体方案设计与技术选型拆解2.1 为什么是 SSE而不是 WebSocket 或轮询传输层的选择是这套方案里第一个要拍板的事。常见候选有三个轮询、WebSocket、SSE。轮询的问题最直接——客户端每隔几百毫秒问一次“好了没”服务端每次都要重新处理上下文延迟高、资源浪费大而且很难做到“逐字”的细腻节奏。WebSocket 是双向全双工能力最强但对于“服务端单向推、客户端只接收”的对话场景来说属于杀鸡用牛刀握手升级、心跳保活、连接状态管理都要额外写一堆代码。SSE 恰好卡在中间它基于普通 HTTP服务端以text/event-stream的 MIME 类型持续向客户端推送文本事件客户端用EventSource或 fetch 流式读取即可。对于 LLM 对话这种“一问、持续多答”的模式SSE 的语义天然契合。而且它走标准 HTTP穿透性好部署时不需要像 WebSocket 那样单独考虑升级协商。注意SSE 是单向的客户端到服务端仍然走普通 POST 请求。所以典型结构是“POST 发起对话 SSE 接收流式回复”而不是用 SSE 连接去发消息。我最终选的是Spring Boot 的SseEmitter来封装服务端推送。它是 Spring MVC 自带的类不用引入额外依赖配合ResponseBodyEmitter体系能很好地管理超时和完成回调。相比手写HttpServletResponse的PrintWriter循环 flushSseEmitter帮我们处理了事件格式、连接关闭、异常回调这些琐事代码干净很多。2.2 后端调用模型流式接口必须端到端打通很多人第一次做流式会踩一个坑前端用了 SSE但后端调模型时用的还是同步阻塞接口结果就是后端等模型全部生成完再一次性通过 SSE 推出去——用户看到的还是“憋一大段”。流式的关键是从模型调用到前端渲染整条链路都必须是流式的。所以后端调用 LLM 时必须使用模型服务提供的流式接口通常是返回一个可迭代的流对象或回调式 API。Spring Boot 这边拿到流之后每收到一个增量片段delta就通过SseEmitter.send()推给前端。这里要注意线程模型模型调用往往是阻塞式的流读取不能占用 Web 请求线程太久否则并发一上来线程池就爆了。我的做法是把模型调用放到独立的线程池里执行Web 线程只负责建立 SSE 连接并返回 emitter。2.3 前端渲染逐段拼接而不是整段替换前端这边接收到的是一连串事件。每个事件里带着一小段文本增量。渲染逻辑有两种常见写法一种是每次收到增量就innerHTML delta另一种是维护一个完整的字符串变量每次追加后重新赋值给渲染节点。我推荐后者。原因是直接操作 DOM 追加在增量很密集时容易触发频繁重排而且一旦需要支持 Markdown 渲染就必须拿到完整文本重新解析。维护一个fullText变量每次fullText delta后交给 Markdown 渲染器处理逻辑更清晰也方便做中断后的状态保留。2.4 abort 中断用户体验的隐藏加分项用户点了发送模型开始哗哗输出突然发现问错了想重新问——这时候如果没有中断能力只能等它说完。abort 的价值就在这里。实现上分两步前端调用AbortController.abort()取消 fetch 请求后端在SseEmitter的onCompletion或onError回调里感知到连接断开进而停止读取模型流、释放资源。这里有个细节如果后端不主动停止模型调用即使前端断了模型那边还在继续生成、继续消耗资源。所以中断必须是双向的——前端断连接后端断模型流。这一点后面在排查章节会详细讲。3. 核心细节解析与实操要点3.1 SseEmitter 的超时设置不能想当然SseEmitter构造时可以传一个超时时间毫秒。如果不传默认走容器的默认超时Tomcat 下通常是 30 秒左右。问题来了LLM 生成一段长回答超过 30 秒是家常便饭。一旦超时连接被容器关掉前端就会收到中断用户看到回答戛然而止。我的做法是把超时设得足够长比如 5 分钟300000L同时在业务层做兜底如果模型生成确实超过预期主动发送一个结束事件并关闭连接。超时时间不是越长越好太长会导致异常连接迟迟不释放占用资源。5 分钟对绝大多数对话场景够用了。SseEmitter emitter new SseEmitter(300_000L);提示不同容器对 SSE 超时的默认值和行为略有差异。如果你用 Tomcat注意connectionTimeout和异步请求超时是两回事别混淆。3.2 事件格式与前端解析的约定SSE 的报文格式是有规范的每个事件由若干字段行组成字段名和值之间用冒号分隔事件之间用空行分隔。常见字段有event事件类型、data数据、id事件 ID、retry重连间隔。SseEmitter.send()有多个重载。如果只传字符串默认作为data字段发送。如果我想区分“正常增量”和“结束信号”可以自定义事件名emitter.send(SseEmitter.event().name(delta).data(chunk)); emitter.send(SseEmitter.event().name(done).data([DONE]));前端解析时EventSource会自动按事件名分发用addEventListener(delta, ...)和addEventListener(done, ...)分别处理。但如果你用 fetch 流式读取我更推荐这种方式因为可以配合 AbortController就需要自己按\n\n切分事件块再解析每块里的字段。这块解析代码不难但要处理跨 chunk 的半截事件——一个事件块可能被 TCP 分包切成两半必须用缓冲区拼接后再切分。3.3 模型流读取的线程与背压问题前面提到模型调用要放到独立线程池。这里还有个容易被忽略的点背压。如果模型生成速度远快于前端消费速度比如前端在做复杂的 Markdown 渲染事件会在服务端堆积。SseEmitter本身没有内置背压机制堆积过多会占内存。实际场景里LLM 的生成速度通常不会快到让前端处理不过来所以这个问题不突出。但如果你的前端渲染逻辑很重建议在服务端做一个简单的节流把短时间内收到的多个 delta 合并成一个事件再发送。这样既减少事件数量又降低前端渲染频率。合并的粒度可以按时间比如每 50ms 合并一次或按字符数比如攒够 20 个字符发一次。3.4 前端 fetch 流式读取的正确姿势用EventSource有个硬伤它只支持 GET 请求没法带复杂的请求体。而对话场景通常需要 POST 一段 JSON包含历史消息、参数等。所以生产环境我更推荐用fetchReadableStream手动读取。const controller new AbortController(); const response await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message: userInput }), signal: controller.signal }); const reader response.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); // 按 \n\n 切分事件块处理完整事件保留半截 const parts buffer.split(\n\n); buffer parts.pop(); for (const part of parts) { // 解析 data 字段并渲染 } }这段代码里有两个关键点decoder.decode(value, { stream: true })的stream: true参数保证多字节字符比如中文不会被截断成乱码buffer的保留机制保证跨 chunk 的事件块能正确拼接。这两点如果漏了中文场景下大概率出现乱码或丢字。4. 实操过程与核心环节实现4.1 后端接口骨架搭建先搭一个最小的流式对话接口。核心是返回SseEmitter并在独立线程里读取模型流、推送事件。RestController RequestMapping(/api/chat) public class ChatController { private final ExecutorService streamExecutor Executors.newCachedThreadPool(); private final LlmClient llmClient; public ChatController(LlmClient llmClient) { this.llmClient llmClient; } PostMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter streamChat(RequestBody ChatRequest request) { SseEmitter emitter new SseEmitter(300_000L); emitter.onCompletion(() - log.info(SSE completed, session{}, request.getSessionId())); emitter.onTimeout(() - { log.warn(SSE timeout, session{}, request.getSessionId()); emitter.complete(); }); emitter.onError(e - log.error(SSE error, session{}, request.getSessionId(), e)); streamExecutor.submit(() - { try { llmClient.streamGenerate(request.getMessage(), new StreamCallback() { Override public void onDelta(String delta) { try { emitter.send(SseEmitter.event().name(delta).data(delta)); } catch (IOException e) { // 前端已断开抛出以终止模型流 throw new StreamAbortedException(e); } } Override public void onComplete() { try { emitter.send(SseEmitter.event().name(done).data([DONE])); emitter.complete(); } catch (IOException ignored) { } } Override public void onError(Throwable t) { try { emitter.send(SseEmitter.event().name(error).data(t.getMessage())); } catch (IOException ignored) { } emitter.completeWithError(t); } }); } catch (StreamAbortedException e) { log.info(Stream aborted by client, session{}, request.getSessionId()); } }); return emitter; } }这段代码里有几个设计决策值得说明。第一produces明确声明text/event-stream让 Spring 知道这是 SSE 响应。第二onCompletion、onTimeout、onError三个回调都注册了方便观测连接生命周期。第三也是最关键的——在onDelta里如果emitter.send()抛IOException说明前端已经断开此时我抛出一个自定义异常StreamAbortedException让模型流的读取循环感知到并停止。这就是前面说的“双向中断”的落地方式。4.2 模型流读取与中断传播LlmClient的streamGenerate方法负责调用模型并逐段回调。伪代码大致如下public void streamGenerate(String prompt, StreamCallback callback) { try (StreamChunk stream modelClient.stream(prompt)) { IteratorChunk it stream.iterator(); while (it.hasNext()) { Chunk chunk it.next(); String delta chunk.getContent(); if (delta ! null !delta.isEmpty()) { callback.onDelta(delta); } } callback.onComplete(); } catch (StreamAbortedException e) { throw e; // 向上传播让外层知道是主动中断 } catch (Exception e) { callback.onError(e); } }关键点在于callback.onDelta抛出的StreamAbortedException会沿着调用栈向上冒泡跳出while循环从而停止继续读取模型流。同时try-with-resources保证流被关闭释放底层连接。这样前端一断后端在下一个 delta 到达时就会感知并停止不会白白消耗资源。注意中断的感知有延迟——必须等到下一个 delta 到达、尝试发送失败时才会触发。如果模型生成很慢中间有较长间隔中断响应也会慢。这是 SSE 方案的固有特性无法做到“立即”中断但通常延迟在几百毫秒内用户无感。4.3 前端完整渲染流程前端部分我把接收、解析、渲染、中断串成一个完整流程。核心是维护fullText和buffer两个变量。async function sendMessage(userInput) { const controller new AbortController(); currentController controller; let fullText ; let buffer ; const response await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ message: userInput, sessionId: currentSessionId }), signal: controller.signal }); const reader response.body.getReader(); const decoder new TextDecoder(utf-8); try { while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const events buffer.split(\n\n); buffer events.pop(); for (const raw of events) { const lines raw.split(\n); let eventName message; let data ; for (const line of lines) { if (line.startsWith(event:)) eventName line.slice(6).trim(); else if (line.startsWith(data:)) data line.slice(5).trim(); } if (eventName delta) { fullText data; renderMarkdown(fullText); } else if (eventName done) { finishRender(fullText); } else if (eventName error) { showError(data); } } } } catch (e) { if (e.name AbortError) { // 用户主动中断保留已生成内容 finishRender(fullText); } else { showError(e.message); } } } function abortGeneration() { if (currentController) { currentController.abort(); currentController null; } }这段代码里renderMarkdown(fullText)每次都用完整文本重新渲染而不是追加。虽然看起来“浪费”但保证了 Markdown 语法的正确解析——比如一个代码块跨了多个 delta只有拿到完整文本才能正确渲染。实际测试下来只要渲染函数本身性能过关这个开销完全可以接受。4.4 参数选择与性能权衡几个关键参数我列个表方便对照调整参数建议值说明SseEmitter 超时300000ms覆盖长回答避免中途断开模型流读取线程池缓存线程池或固定 20-50 线程视并发量调整避免阻塞 Web 线程前端事件合并粒度50ms 或 20 字符渲染重时启用减少重排abort 感知延迟取决于 delta 间隔通常 500ms缓冲区大小无硬限制按事件切分注意半截事件处理线程池大小这块我的经验是如果模型调用是 IO 密集型大部分时间在等网络线程数可以设得比 CPU 核数大不少20 到 50 是常见区间。但要注意模型服务端通常有并发限制线程开太多反而会触发限流。所以实际值要结合你的模型服务配额来定。5. 常见问题与排查技巧实录5.1 中文乱码十有八九是解码姿势不对前端收到data后显示成乱码最常见的原因是TextDecoder没有用stream: true。UTF-8 的中文占 3 个字节如果一次read()恰好把某个汉字切成两半不用流式解码就会得到替换字符。加上{ stream: true }后解码器会缓存不完整的字节序列等下一个 chunk 到达再拼。这个坑我踩过不止一次排查时先看这里。5.2 回答中途截断检查超时和容器配置用户反馈“回答说到一半就没了”排查顺序是先看SseEmitter超时是不是设太短再看容器Tomcat/Nginx有没有自己的超时或缓冲配置。特别是 Nginx 反代场景默认proxy_buffering on会把 SSE 响应缓冲起来导致前端迟迟收不到数据看起来像卡住。需要显式关闭location /api/chat/stream { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; proxy_set_header Connection ; proxy_http_version 1.1; }proxy_buffering off是 SSE 场景的必配项漏了它前端可能一直等到响应结束才收到全部内容流式效果直接失效。5.3 中断后资源没释放onCompletion 里要做清理前端 abort 后后端如果只在onDelta抛异常时停止模型流还有一种情况没覆盖模型流刚好在两次 delta 之间前端断了但后端还没感知。这时候onCompletion回调会被触发连接关闭时可以在这里设置一个标志位让模型流读取循环在下次检查时退出。我一般用一个AtomicBoolean aborted配合双保险。5.4 常见问题速查表现象可能原因排查方向中文乱码解码未用 stream 模式检查 TextDecoder 参数回答截断超时过短或 Nginx 缓冲检查超时配置和 proxy_buffering流式变一次性后端调模型用了同步接口确认模型调用是流式中断无效后端未感知连接断开检查 onCompletion 和异常传播连接数暴涨线程池或连接未释放检查 complete 调用和线程池配置首字延迟高模型首 token 慢属模型侧问题可加 loading 提示5.5 几个我踩过的坑第一个坑是忘记调用emitter.complete()。模型流结束后如果不显式 complete连接会一直挂着直到超时才释放。并发一高连接数就爆了。所以onComplete里必须调 complete。第二个坑是在 Web 线程里直接读模型流。早期我图省事直接在 Controller 方法里循环读流并 send结果一个请求占一个 Tomcat 线程几十个并发就把线程池占满了。后来改成独立线程池才解决。第三个坑是前端没处理done事件。有些实现里流结束后前端还在等界面上的“生成中”状态不消失。所以后端一定要发一个明确的结束事件前端收到后清理状态。第四个坑是Markdown 渲染时机。如果每个 delta 都触发一次完整的 Markdown 解析delta 密集时 CPU 会飙高。我的优化是加一个简单的节流用requestAnimationFrame或setTimeout把渲染频率限制在每秒 30 次左右视觉上依然流畅CPU 压力小很多。6. 关于这套方案的一些延伸想法这套 SSE 流式方案跑通之后其实还能往几个方向扩展。比如多轮对话的上下文管理可以在ChatRequest里带上历史消息列表后端拼接后传给模型比如生成过程中的“停止”按钮本质上就是前端 abort 加后端中断传播已经包含在这套实现里了再比如把流式能力复用到知识库问答场景模型先检索再生成检索阶段可以先推一个“正在检索”的事件让用户知道系统在干活。我个人在实际项目里体会最深的一点是流式对话的难点从来不在“怎么把字推出去”而在“怎么在推的过程中保持状态一致”。连接断了要能恢复、用户中断要能停干净、异常了要能给出明确反馈——这些边界情况的处理才是决定这套东西能不能上生产的关键。上面那些排查技巧基本都是被线上问题逼出来的希望能帮你少熬几个夜。
返回列表