☰
LLM流式对话实战:Spring Boot+SSE实现打字机效果与abort中断
2026/9/28 23:37:08 网站建设 项目流程

1. 为什么“流式”是 LLM 对话产品的分水岭

做过对话类产品的朋友大概率都有过这种体验:用户点下发送按钮,界面转圈转了七八秒,然后“啪”一下整段答案全冒出来。功能上没毛病,但用起来就是别扭——像对着一个沉默寡言的人说话,对方憋半天才一次性回你一大段。而 ChatGPT 那种一个字一个字往外蹦的效果,哪怕总耗时一样,主观感受也快得多、自然得多。这中间的差别,就是流式对话。

我这次要聊的,就是围绕LLM 流式对话做的一套前后端联接的完整实现。核心目标很明确:让大模型的回答像打字机一样实时渲染到页面上,而不是等全部生成完再一次性返回。技术栈上,后端用Spring Boot承载业务与模型调用,传输层用SSE(Server-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循环 flush,SseEmitter帮我们处理了事件格式、连接关闭、异常回调这些琐事,代码干净很多。

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(包含历史消息、参数等)。所以生产环境我更推荐用fetch+ReadableStream手动读取。

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 (Stream<Chunk> stream = modelClient.stream(prompt)) { Iterator<Chunk> 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 加后端中断传播,已经包含在这套实现里了;再比如把流式能力复用到知识库问答场景,模型先检索再生成,检索阶段可以先推一个“正在检索”的事件,让用户知道系统在干活。

我个人在实际项目里体会最深的一点是:流式对话的难点从来不在“怎么把字推出去”,而在“怎么在推的过程中保持状态一致”。连接断了要能恢复、用户中断要能停干净、异常了要能给出明确反馈——这些边界情况的处理,才是决定这套东西能不能上生产的关键。上面那些排查技巧,基本都是被线上问题逼出来的,希望能帮你少熬几个夜。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询