Appearance
WebSocket 双向流:让用户能随时打断模型
1. 本节产出
一个 WebSocket 端点 /ws/chat:客户端发一句话,服务端逐字回推;客户端可以在任意时刻发送 {"type":"stop"} 立刻中断生成,服务端立即停止调用模型并回报已消耗的内容。
2. 前置依赖
- 01-04 流式输出 SSE:已理解
Flux与流式中止 - 依赖
spring-boot-starter-websocket - 前端能写基本的 WebSocket 调用
3. 为什么还要 WebSocket
SSE 已经能打字机了,为什么还要多这一节?因为有一个 SSE 解决不了的场景:
用户在模型输出过程中想插话。
| 场景 | SSE 能做吗 | WebSocket |
|---|---|---|
| 逐字输出 | 可以 | 可以 |
| 用户中途点「停止」 | 勉强:要另发一个 HTTP 请求,还要把 stop 信号和流关联起来 | 天然:同一条连接上发一个 stop 帧 |
| 用户中途追加条件 | 不行:只能断开重连,丢失上下文 | 可以:在原连接上追加指令 |
| 服务端主动推送(比如后台任务进度) | 不行 | 可以 |
| 语音实时对话(边说边出字) | 不行 | 可以 |
判断标准:如果你的交互是「一问一答」,SSE 足够,不要上 WebSocket。只有需要「客户端在生成过程中反向控制服务端」时,WebSocket 才值得。多一层协议就多一层运维复杂度——代理配置、心跳、重连、会话状态,全都要自己管。
4. 核心原理
4.1 会话生命周期
WS 握手 → 建立会话(分配 sessionId)
│
├─ 收到 chat 消息 → 订阅 Flux → 逐帧回推
│ └─ 收到 stop 消息 → disposable.dispose() → Flux 取消
│
└─ 连接关闭 → 清理会话、记录用量关键对象是 Disposable:订阅 Flux 时返回的句柄,调用它的 dispose() 就能取消整个流。这是整个「可打断」能力的支点。
4.2 中断是怎么传导的
客户端 send {"type":"stop"}
│
▼
服务端找到该 sessionId 对应的 Disposable
│
▼
dispose() → Flux 发出 cancel 信号
│
▼
Spring AI 关闭与厂商的 SSE 连接
│
▼
厂商停止生成 → 已生成的 Token 计费,未生成的不计费注意最后一步:中断只能省下「还没生成的部分」。已经生成的内容照样计费。所以「打断」是止损手段,不是免费手段——这点要跟学员讲清楚,否则他们会以为随便打断不花钱。
4.3 并发与状态管理
| 问题 | 解法 |
|---|---|
| 一个用户开多个标签页 | 以 sessionId 而非 userId 为单位管理 Disposable |
| 用户断线后重连 | 会话状态不能只放内存,需持久化(见 01-06) |
| 服务端多副本 | WebSocket 是有状态连接,需要粘性会话或集中式会话存储 |
第三点是生产上最大的坑:WebSocket 把「有状态」带回了你的服务,这会和 K8s 的弹性伸缩直接冲突。做法是把会话状态外置到 Redis,让任意副本都能接管——这条在 04-03 弹性伸缩里会详细展开。
5. 代码走查
5.1 依赖与配置
xml
<!-- pom.xml -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>java
// src/main/java/com/example/aibasics/config/WebSocketConfig.java
@Configuration
public class WebSocketConfig {
@Bean
public HandlerMapping webSocketMapping(ChatWebSocketHandler handler) {
Map<String, WebSocketHandler> map = Map.of("/ws/chat", handler);
SimpleUrlHandlerMapping mapping = new SimpleUrlHandlerMapping();
mapping.setUrlMap(map);
mapping.setOrder(1); // 必须优先于静态资源映射
return mapping;
}
@Bean
public WebSocketHandlerAdapter handlerAdapter() {
return new WebSocketHandlerAdapter();
}
}提示:
setOrder(1)容易被漏掉。不设的话可能被其他 HandlerMapping 抢先,表现是「连接能建立但立刻关闭」,排查起来很费时间。
5.2 会话与中断的核心实现
java
// src/main/java/com/example/aibasics/ws/ChatWebSocketHandler.java
@Component
public class ChatWebSocketHandler implements WebSocketHandler {
private final ChatClient chatClient;
/** sessionId → 当前正在进行的订阅,用于中断 */
private final ConcurrentHashMap<String, Disposable> running = new ConcurrentHashMap<>();
@Override
public Mono<Void> handle(WebSocketSession session) {
String id = session.getId();
Flux<String> inbound = session.receive()
.map(WebSocketMessage::getPayloadAsText)
.doOnNext(text -> handleInbound(session, id, text))
.doFinally(sig -> stop(id)); // 连接断开时兜底清理
// outbound 保持空,输出通过 session.send 主动推送
return session.send(Mono.never()).and(inbound).then();
}
private void handleInbound(WebSocketSession session, String id, String text) {
Message msg = Json.decode(text, Message.class);
if ("stop".equals(msg.type())) {
stop(id);
send(session, "{\"type\":\"stopped\"}");
return;
}
if ("chat".equals(msg.type())) {
Disposable d = chatClient.prompt()
.user(msg.content())
.stream()
.content()
.doOnNext(chunk -> send(session, chunk))
.doFinally(sig -> {
running.remove(id);
send(session, "{\"type\":\"done\",\"signal\":\"" + sig + "\"}");
})
.subscribe(); // 必须显式 subscribe
running.put(id, d);
}
}
private void stop(String id) {
Disposable d = running.remove(id);
if (d != null && !d.isDisposed()) d.dispose();
}
private void send(WebSocketSession session, String payload) {
if (session.isOpen()) {
session.send(Mono.just(session.textMessage(payload))).subscribe();
}
}
public record Message(String type, String content) {}
}三个必须注意的点:
| 点 | 说明 |
|---|---|
ConcurrentHashMap 存 Disposable | 中断必须能找到正在跑的那个订阅 |
doFinally 里 running.remove(id) | 不清理会导致内存泄漏,且下次中断无效 |
显式 subscribe() | 在命令式写法里 Flux 不会自动执行,忘了就是「什么都没发生」 |
5.3 客户端
javascript
const ws = new WebSocket('ws://localhost:8080/ws/chat');
ws.onopen = () => ws.send(JSON.stringify({type:'chat', content:'讲讲 JVM 垃圾回收'}));
ws.onmessage = (e) => {
const d = e.data;
if (d.startsWith('{')) { console.log('控制帧', JSON.parse(d)); return; }
document.getElementById('out').textContent += d; // 正文帧直接拼接
};
// 用户点「停止」
function stopGeneration() {
ws.send(JSON.stringify({type:'stop'}));
}
// 心跳:防止代理掐断空闲连接
setInterval(() => ws.readyState === 1 && ws.send('{"type":"ping"}'), 25000);6. 跑起来
bash
git checkout ch01-05-websocket
mvn spring-boot:run
# 打开 http://localhost:8080/ws.html验证一:正常对话
- 页面输入「讲讲 JVM 垃圾回收」,点发送;
- 观察文字逐字出现;
- 结束后收到
{"type":"done","signal":"onComplete"}。
验证二:中断(本节核心)
- 输入「写一篇五千字的 Spring 源码分析」,点发送;
- 输出约 2 秒后点「停止」;
- 立刻收到
{"type":"stopped"}和{"type":"done","signal":"cancel"}; - 输出停止,不再有新内容追加。
验证三:连接关闭时的兜底
- 生成过程中直接关闭浏览器标签;
- 服务端日志出现
stop: <sessionId>且running中已无该会话。
| 检查项 | 通过标准 |
|---|---|
| 中断及时性 | 点停止后 1 秒内不再有新内容 |
| done 帧 signal | 中断时为 cancel,正常结束为 onComplete |
| 会话清理 | 中断或断连后 running 中不留残留 |
| 重复中断 | 连续点两次停止不报错 |
7. 生产避坑
- WebSocket 让服务变成有状态,直接和弹性伸缩冲突。多副本下,用户在 A 副本建立的连接,重连时可能落到 B 副本,会话状态丢失。做法是:要么配置网关粘性会话,要么把会话状态外置到 Redis——不要只放内存,本地跑得通不代表生产跑得通。
- 代理和网关默认会掐断空闲 WebSocket。Nginx 默认
proxy_read_timeout 60s,模型思考时间长一点就被断开,表现为「回答到一半突然没了」。要么调大超时,要么像上面那样做 25 秒心跳。推荐两者都做,因为超时配置往往不在你手里。 session.send()在连接已关闭时调用会抛异常。高并发下中断和关闭是竞态的,必须在发送前判session.isOpen(),并且用onErrorResume兜住。忽略这点会在压测时出现大量无意义错误日志,淹没真正的问题。
8. 延伸与锚点
- 思考题:现在中断了,但重新开始就要重发整个对话。怎么做到「中断后能从断点继续」?(提示:需要把历史消息存下来——下一课时)
- 代码锚点:
git checkout ch01-05-websocket - 下一课时:01-06 多轮对话与 Memory 持久化
- 对应课件:L01-05 WebSocket 双向流