优化 SSE 流式输出,核心是四件事:让第一个 token 尽快到前端、连接全程不中断、碎片 token 合并后推送、结束时明确告知客户端。工程上我们用 SseEmitter 承载 Spring AI 的 Flux 流,15 秒心跳保活,200ms 微缓冲合并碎片,多阶段生成用 event 区分事件类型。首 token 耗时是首要指标,页面拿到首字就渲染,用户体感从”等结果”变成”看生成”。
一、为什么用 SSE 而不是轮询和 WebSocket
| 方案 | 方向 | 实时性 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 轮询 | 请求-响应 | 差 | 低 | 低频状态查询 |
| SSE | 单向服务端推送 | 高 | 低 | 流式文本、进度条 |
| WebSocket | 双向 | 高 | 高 | 聊天、协同编辑 |
文章生成是纯服务端到客户端的单向流,SSE 走普通 HTTP,浏览器原生 EventSource 就能接,没有升级握手和额外协议开销。
二、基础链路:ChatClient.stream() → Flux → SseEmitter
项目用的 Spring AI 的 ChatClient,流式返回一个 Flux。Controller 返回 SseEmitter,在订阅回调里把每个 token 推给前端:
@PostMapping(value = "/api/article/generate/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter stream(@RequestBody GenerateRequest request) {
SseEmitter emitter = new SseEmitter(300_000L);
Flux<String> tokenStream = chatClient.prompt()
.user(buildPrompt(request))
.stream()
.content();
tokenStream.subscribe(
token -> safeSend(emitter, SseEmitter.event().name("token").data(token)),
emitter::completeWithError,
() -> {
safeSend(emitter, SseEmitter.event().name("done").data("[DONE]"));
emitter.complete();
}
);
return emitter;
}
订阅的 onError 和 onComplete 对应 emitter 的 completeWithError 与 complete,两端的生命周期必须对齐,否则连接悬着不释放。
三、四个关键优化点
3.1 心跳保活
AI 生成一次几十秒,中间可能长时间没有输出,代理层和浏览器容易判定超时断开。启动一个定时任务,每 15 秒发一个注释事件:
ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor();
scheduler.scheduleAtFixedRate(() -> {
try {
emitter.send(SseEmitter.event().comment("heartbeat"));
} catch (IOException e) {
scheduler.shutdown(); // 发送失败说明连接断了,停掉心跳
}
}, 0, 15, TimeUnit.SECONDS);
onComplete / onTimeout / onError 三个回调里都要 scheduler.shutdown(),防止心跳线程泄漏。
3.2 碎片 token 微缓冲
流式返回里一个 token 经常才几个字符,逐个推送,网络包多、渲染抖动,前端每帧都在重排 DOM。用 bufferTimeout 把 200ms 内的 token 攒成一小段再发:
tokenStream
.bufferTimeout(10, Duration.ofMillis(200))
.map(chunk -> String.join("", chunk))
200ms 的缓冲人眼无感,前端渲染次数却减少一个数量级。多阶段生成时,阶段标题这类关键节点不做缓冲,立即推送,进度切换才跟手。
3.3 多阶段生成的事件分发
创作器一次生成经历选题、大纲、初稿、润色四个阶段。每个阶段开始推一个 stage 事件,前端收到后切换 loading 文案:
safeSend(emitter, SseEmitter.event().name("stage")
.data("{\"stage\":\"draft\",\"text\":\"正在生成正文\"}"));
前端 addEventListener(“stage”) 更新进度条,addEventListener(“token”) 追加正文,两套渲染逻辑互不干扰。
3.4 代理层与浏览器端的配合
Nginx 默认开 proxy_buffering,SSE 数据会在代理层攒批,后端推一个字前端可能要等好久。SSE 路径必须显式关掉缓冲:
location /api/article/generate/stream {
proxy_pass http://backend;
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 3600s;
add_header X-Accel-Buffering no;
}
gzip 也会把流式内容压缩进缓冲,SSE 接口排除压缩。服务端对 SSE 响应显式返回 X-Accel-Buffering: no,双重保险。
四、前端怎么收
浏览器 EventSource 不支持自定义请求头,带 Token 的接口改用 fetch + ReadableStream 手写解析,控制力更强:
const res = await fetch('/api/article/generate/stream', {
method: 'POST',
headers: { 'Content-Type': 'application/json', 'Authorization': `Bearer ${token}` },
body: JSON.stringify(params),
});
const reader = res.body.getReader();
const decoder = new TextDecoder();
let buf = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buf += decoder.decode(value, { stream: true });
for (const line of buf.split('\n')) {
if (line.startsWith('data:')) onToken(line.slice(5));
}
buf = '';
}
断开重连用 Last-Event-ID 语义:前端记录已收到的最后一段,重连时带回去,服务端把它作为上下文前缀重新发起生成,弥补 SSE 不支持断点续传的缺陷。
五、SseEmitter 与 WebFlux 怎么选
SseEmitter 每个连接占一个 servlet 线程,Tomcat 线程池小的话并发一高就拒连。企业级 AI 网关项目用的是 WebFlux 响应式栈,事件循环模型,几千并发连接不吃线程。单机文章创作器并发量不大,SseEmitter 够用且改造成本低;面向大规模对话网关,优先 WebFlux。
常见问题(FAQ)
Q1:SSE 连接断了能断点续传吗?
SSE 协议不支持,客户端带已收内容重连,服务端重新生成。
Q2:心跳注释事件会被前端误当数据吗?
不会,注释行不以 data 开头,EventSource 直接忽略。
Q3:流式接口一定要关 gzip 吗?
要。压缩缓冲会攒批,破坏逐字输出,SSE 路径排除压缩。