Skip to content

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
// ch01-basics/src/main/java/com/aitech/basics/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
// ch01-basics/src/main/java/com/aitech/basics/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

验证一:正常对话

  1. 页面输入「讲讲 JVM 垃圾回收」,点发送;
  2. 观察文字逐字出现;
  3. 结束后收到 {"type":"done","signal":"onComplete"}。

验证二:中断(本节核心)

  1. 输入「写一篇五千字的 Spring 源码分析」,点发送;
  2. 输出约 2 秒后点「停止」;
  3. 立刻收到 {"type":"stopped"} 和 {"type":"done","signal":"cancel"};
  4. 输出停止,不再有新内容追加。

验证三:连接关闭时的兜底

  1. 生成过程中直接关闭浏览器标签;
  2. 服务端日志出现 stop: <sessionId> 且 running 中已无该会话。
检查项通过标准
中断及时性点停止后 1 秒内不再有新内容
done 帧 signal中断时为 cancel,正常结束为 onComplete
会话清理中断或断连后 running 中不留残留
重复中断连续点两次停止不报错

7. 生产避坑 ​

  1. WebSocket 让服务变成有状态,直接和弹性伸缩冲突。多副本下,用户在 A 副本建立的连接,重连时可能落到 B 副本,会话状态丢失。做法是:要么配置网关粘性会话,要么把会话状态外置到 Redis——不要只放内存,本地跑得通不代表生产跑得通。
  2. 代理和网关默认会掐断空闲 WebSocket。Nginx 默认 proxy_read_timeout 60s,模型思考时间长一点就被断开,表现为「回答到一半突然没了」。要么调大超时,要么像上面那样做 25 秒心跳。推荐两者都做,因为超时配置往往不在你手里。
  3. session.send() 在连接已关闭时调用会抛异常。高并发下中断和关闭是竞态的,必须在发送前判 session.isOpen(),并且用 onErrorResume 兜住。忽略这点会在压测时出现大量无意义错误日志,淹没真正的问题。

8. 延伸与锚点 ​

  • 思考题:现在中断了,但重新开始就要重发整个对话。怎么做到「中断后能从断点继续」?(提示:需要把历史消息存下来——下一课时)
  • 代码锚点:git checkout ch01-05-websocket
  • 下一课时:01-06 多轮对话与 Memory 持久化
  • 对应课件:L01-05 WebSocket 双向流