Appearance
流式输出 SSE:让回答一个字一个字地出来
1. 本节产出
把同步接口改造成流式:GET /api/chat/stream 返回 SSE,浏览器端用 EventSource 收到打字机效果,并且客户端断开连接后服务端会立刻取消订阅,不再继续烧 Token。
2. 前置依赖
- 01-02 第一个 ChatClient:已有可调用的同步接口
- 01-03 国内多模型统一适配层:厂商兼容端点已验证
- 了解
Flux的基本操作符(map/takeUntil/doFinally)即可,不需要精通 Reactor
3. 为什么必须做流式
先看数字。旗舰模型生成一个 500 字的回答,首 Token 延迟约 0.8 秒,完整生成约 12 秒。
| 模式 | 用户感知 | 实际后果 |
|---|---|---|
| 同步 | 转圈 12 秒,然后整段出现 | 用户以为卡死,刷新页面重试 → 又烧一次钱 |
| 流式 | 0.8 秒后出现第一个字,持续输出 | 用户在读内容,等待感消失 |
真正的损失不是体验,是钱:同步模式下用户中途刷新或关闭页面,服务端已经把整段生成完了,Token 照扣,用户什么也没看到。流式模式下客户端断开可以立刻取消,用的是多少算多少。
一个容易被忽略的事实:流式不会让模型变快,总时长几乎一样。它改变的是「首字节时间(TTFB)」和「可中断性」。讲课时务必说清这点,否则学员会以为流式是性能优化。
4. 核心原理
4.1 数据是怎么流出来的
模型侧 Spring AI 浏览器
───────── ────────── ─────────
chunk: "你" ──▶ Flux<String> ──▶ SSE: data: 你
chunk: "好" ──▶ SSE: data: 好
chunk: "," ──▶ SSE: data: ,
[DONE] ──▶ onComplete ──▶ EventSource 关闭模型厂商在流式模式下返回的是 Server-Sent Events 格式的分片,每个分片携带一小段增量文本。Spring AI 把它转成 Flux<String>(或 Flux<ChatResponse>),你只需要把这个 Flux 直接作为响应体返回。
4.2 SSE vs WebSocket 怎么选
| 维度 | SSE | WebSocket |
|---|---|---|
| 方向 | 服务端 → 客户端(单向) | 双向 |
| 协议 | HTTP,天然穿透代理和网关 | 需要升级握手,部分网关不支持 |
| 断线重连 | 浏览器 EventSource 自动重连 | 需自己实现心跳与重连 |
| 复杂度 | 极低,Controller 返回 Flux 即可 | 需要配置 Handler/握手/会话管理 |
| 适用场景 | LLM 逐字输出(95% 的场景) | 语音实时对话、多人协同、需要客户端随时打断 |
结论:先做 SSE。 只有当业务真的需要「用户在模型输出中途插入新指令」时,才上 WebSocket(见 01-05)。大多数项目声称要 WebSocket,其实只是要打字机效果。
4.3 背压与取消
浏览器关闭页面
│
▼
TCP 连接断开 → Spring WebFlux 检测到
│
▼
Flux 收到 cancel 信号
│
▼
.doFinally(...) 执行 → 释放资源、记录用量.doFinally() 是必须加的:它是你唯一能观察到「流被中断」的钩子。中断处理三件事——记日志、记用量、释放与该请求绑定的资源。
5. 代码走查
5.1 Controller:关键在 produces
java
// ch01-basics/src/main/java/com/aitech/basics/controller/StreamController.java
@RestController
@RequestMapping("/api")
public class StreamController {
private final ChatClient chatClient;
public StreamController(ChatClient chatClient) {
this.chatClient = chatClient;
}
@GetMapping(value = "/chat/stream",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> stream(@RequestParam String message,
@RequestParam(defaultValue = "default") String convId) {
return chatClient.prompt()
.user(message)
.stream() // 注意:不是 .call()
.content() // Flux<String>
.doFinally(signal -> {
// signal 可能是 onComplete / onError / cancel
log.info("stream finished: convId={}, signal={}", convId, signal);
});
}
}三个必须写对的地方:
| 写法 | 为什么 |
|---|---|
produces = TEXT_EVENT_STREAM_VALUE | 不加这个,浏览器不会按 SSE 解析,会一次性等到结束 |
.stream() 而非 .call() | call() 返回聚合后的单个结果,没有流 |
.doFinally(...) | 唯一能观测中断的位置,用于记账与清理 |
5.2 需要更丰富信息时用 ChatResponse
java
// 需要拿到 usage(用量)时,用 chatResponse() 而不是 content()
@GetMapping(value = "/chat/stream/full",
produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<ServerSentEvent<String>> streamFull(@RequestParam String message) {
return chatClient.prompt()
.user(message)
.stream()
.chatResponse() // Flux<ChatResponse>
.map(resp -> {
String text = resp.getResult().getOutput().getText();
return ServerSentEvent.<String>builder()
.data(text == null ? "" : text)
.build();
})
.concatWith(Flux.defer(() -> // 流结束后补一条用量事件
Flux.just(ServerSentEvent.<String>builder()
.event("usage")
.data("{\"done\":true}")
.build())));
}提示:流式模式下多数厂商不返回 usage 字段(因为服务端还没统计完)。想精确计费,要么在流结束后再查一次用量接口,要么用 01-01 的估算器兜底。这一点在做成本治理时非常关键,别等上线才发现。
5.3 前端消费
html
<!-- src/main/resources/static/stream.html -->
<div id="out"></div>
<script>
const out = document.getElementById('out');
const es = new EventSource('/api/chat/stream?message=' + encodeURIComponent(q));
es.onmessage = (e) => { out.textContent += e.data; };
es.onerror = (e) => { console.warn('stream error', e); es.close(); };
// 页面关闭时必须主动关,否则浏览器会重连
window.addEventListener('beforeunload', () => es.close());
</script>beforeunload 里 es.close() 是必须的。EventSource 默认会在连接断开后自动重连——如果用户在回答生成完后关闭页面又打开,或者网络抖动,会触发一次全新的请求,又烧一遍钱。这是新手最常踩的坑。
5.4 阻断与超时
java
// 限制最多接收 N 个分片,防止异常长的输出
.chatResponse()
.takeUntil(resp -> resp.getResult().getOutput().getText().contains("。END"))
.timeout(Duration.ofSeconds(60)) // 整体超时
.onErrorResume(TimeoutException.class, ex -> {
log.warn("stream timeout");
return Flux.just(fallbackResponse());
});6. 跑起来
bash
git checkout ch01-04-streaming-sse
mvn spring-boot:run验证一:命令行看原始 SSE
bash
curl -N -H "Accept: text/event-stream" \
"http://localhost:8080/api/chat/stream?message=介绍Spring事务传播"期望输出(逐行、逐字出现,不是一次性):
data: Spring
data: 的事务
data: 传播机制
data: 定义了
...验证二:浏览器看打字机
http://localhost:8080/stream.html验证三:中断测试(本节最重要的一步)
bash
# 启动请求后 2 秒按 Ctrl+C
timeout 2 curl -N "http://localhost:8080/api/chat/stream?message=写一篇两千字的技术文章"然后看应用日志,必须出现类似:
stream finished: convId=default, signal=cancel| 检查项 | 通过标准 |
|---|---|
| 首字延迟 | 明显早于完整回答(1~2 秒内出现) |
| 中断日志 | signal=cancel,不是 onComplete |
| 中断后开销 | 服务端日志显示未继续生成剩余内容 |
7. 生产避坑
- 反向代理会缓冲 SSE。Nginx 默认开启
proxy_buffering,会把你的流式响应攒着一起发,打字机效果直接消失。必须配proxy_buffering off;和proxy_read_timeout 300s;。这个坑在内网演示时不会出现,一上生产就爆。 - 客户端断开 ≠ 服务端立刻停止。取消信号要穿过整个调用链才能生效,中间任何一环做了阻塞或缓存都会延迟取消。做法是在
doFinally里记录signal,并监控cancel占比——如果这个比例异常高,说明大量请求被中途放弃,要查是不是响应太慢。 - 别忘了 SSE 也占用连接数。浏览器对同一域名的并发连接数有上限(HTTP/1.1 下通常 6 个)。用户在多个标签页同时对话会把连接占满,其他请求开始排队。生产上要么上 HTTP/2,要么给流式接口单独域名。
8. 延伸与锚点
- 思考题:流式已经够了,什么业务真的需要 WebSocket?(提示:需要用户中途打断、或服务端主动推送——答案在下一课时)
- 代码锚点:
git checkout ch01-04-streaming-sse - 下一课时:01-05 WebSocket 双向流
- 对应课件:L01-04 流式输出 SSE