流式输出已是 AI 应用标配,而「客户端断开后的资源治理」直接决定你的服务会不会在高峰期裸奔烧钱。本文从底层(TCP / 响应式取消传播)讲清原理,并给出 Go / PHP / Python / Java 四种语言的可运行实现。
你不需要「主动去感知」前端断开。只要链路不阻断,TCP 断开会以 cancel 信号自动从下游传到上游。
把整条链路摊开,取消信号是这样传播的:
后端怎么知道前端走了?靠 TCP 连接断开(RST/FIN)被服务端网络层立即捕获,而不是靠前端发「我要走了」的通知。强杀、断网、锁屏都不会触发 beforeunload,但一定会断 TCP。
把断开变成「取消信号」,一路传到最上游——真正发起对 LLM 的 HTTP 请求的那一层,把那条连接关掉,让模型停止生成,停止烧钱。
前端关页面往往意味着会话「半途而废」。会话上下文、并发计数、临时状态、连接池资源都要在 统一收尾点清掉,否则下轮对话错乱、并发数虚高触发误限流、连接泄漏最终 OOM。
感知方式取决于技术栈,主流有三条路。关键词:别信 beforeunload。
Reactor Netty 监听 TCP 状态,客户端断开瞬间内核发来 RST/FIN → 最下游 Subscriber 调用 subscription.cancel() → cancel 沿 Flux 链一路向上传播。你只在链上挂监听即可。
// Spring WebFlux + Spring AI return chatClient.prompt(request.message()) .stream().content() .doOnCancel(() -> log.warn("客户端断开,触发取消 | session={}", sessionId)) // 感知点 .doFinally(signal -> { if (signal == SignalType.CANCEL) log.warn("流被取消——前端跑了"); else if (signal == SignalType.ON_COMPLETE) log.info("流正常结束"); else if (signal == SignalType.ON_ERROR) log.error("流异常终止"); cleanupSession(sessionId); // 统一收尾 });
MVC 下客户端断开不会自动传播 cancel——你是在「下一次 send 失败」时才后知后觉。dispose() 是手动取消模型订阅的核心动作,每个回调都必须带上。
SseEmitter emitter = new SseEmitter(120_000L); // 必须设超时! Disposable sub = chatClient.prompt(req.message()).stream().content() .subscribe( chunk -> { try { emitter.send(SseEmitter.event().data(chunk)); } catch (IOException e) { log.warn("send 失败,客户端已断开"); sub.dispose(); } }, emitter::completeWithError, emitter::complete); emitter.onCompletion(() -> { log.info("SSE 关闭"); sub.dispose(); }); emitter.onTimeout(() -> { log.warn("SSE 超时,强制取消"); sub.dispose(); emitter.complete(); }); emitter.onError(e -> { log.warn("SSE 异常: {}", e.getMessage()); sub.dispose(); });
dispose() 漏掉任何一个回调 = 一条「烧钱流」。
若中间隔着代理,TCP 断连感知会被延迟甚至吞掉。服务端每 15~30s 发一个 SSE 注释行 :heartbeat\n\n,一旦 send 心跳失败即判定连接死亡,主动取消模型调用。
Flux<ServerSentEvent<String>> heartbeat =
Flux.interval(Duration.ofSeconds(15))
.map(i -> ServerSentEvent.<String>builder().comment("keep-alive").build());
// 与数据流 merge,任一被取消都能感知
return Flux.merge(dataStream, heartbeat)
.doOnCancel(() -> log.warn("心跳或数据流被取消"));
感知到了,关键动作是让 LLM 那边的 HTTP 请求真正断掉。
下游 cancel → ChatClient 的 Flux 被取消 → WebClient 响应 Flux 被取消 → Reactor Netty 客户端关闭与 LLM 的 TCP 连接 → 模型端收到断开停止生成。你不用写任何「取消模型」的代码,只要保证信号能传上去。
下游 cancel ↓ ChatClient 的 Flux 被取消 ↓ WebClient 的响应 Flux 被取消 ↓ Reactor Netty 关闭与 LLM 的 TCP 连接 ↓ 模型端收到连接关闭 → 停止生成
block() 会卡死传播链,cancel 传不上来。要用响应式 API(返回 Mono,由 flatMap 串联)。
// ❌ 错误 .flatMap(chunk -> { String r = someBlockingDbQuery(chunk); return Flux.just(r); }) // ✅ 正确 .flatMap(chunk -> reactiveDbQuery(chunk))
mono.subscribe() 在 flatMap 里创建了与外部链脱节的订阅,cancel 传不到,模型调用变「孤儿请求」。应返回链内对象由框架串联。
// ❌ 错误:外部 cancel 传不到 .flatMap(chunk -> { someAsyncCall(chunk).subscribe(); return Flux.just(chunk); }) // ✅ 正确 .flatMap(chunk -> someAsyncCall(chunk))
HttpClient.send()(同步阻塞)调模型,cancel 根本传不进去。要么换成 sendAsync() / WebClient,要么手动取消。
ExecutorService 提交,中断正在执行 HTTP 调用的线程。requestId,前端断开时后端写「取消表」;模型响应返回后发现有取消标记就直接丢弃结果。模型可能仍在烧钱,但不污染业务状态。取消只是手段,清理才是目的。
| 要清理的东西 | 为什么必须清 | 不清的后果 |
|---|---|---|
| 会话上下文(ChatMemory) | 半截对话塞进历史会错乱 | 用户下次对话前言不搭后语 |
| 计数 / 并发状态 | 流断了计数没减 | 并发数虚高 → 触发限流误伤 |
| DB / 缓存临时状态 | 半成品记录残留 | 脏数据、下次读到过期状态 |
| 线程池 / 连接池 | 订阅未释放 | 连接泄漏 → 高峰期 OOM / 无连接 |
doFinally,而不是 doOnCancel 或 onError。因为 doFinally 无论流是正常完成、被取消、还是异常都会执行,保证必达。
.doFinally(signal -> { chatMemory.clear(sessionId); // 清会话 concurrency.decrementAndGet(); // 减并发计数 releaseResources(sessionId); // 释放连接 });
点击下方语言切换。每个实现都覆盖:感知(连接断开)→ 取消(传给上游 HTTP)→ 清理(收尾)。
Go 的 context.Context 就是 cancel 信号的载体。HTTP/2 流(SSE / 流式)在客户端断开时,r.Context().Done() 会被关闭。只要把这个 ctx 一路透传到上游 LLM 的 HTTP 请求,模型调用会被自动取消——原理和响应式 cancel 传播一模一样,只是用 context 表达。
func streamChat(w http.ResponseWriter, r *http.Request) { ctx := r.Context() // 客户端断开 → ctx 自动 Done flusher := w.(http.Flusher) w.Header().Set("Content-Type", "text/event-stream") // 1) 向上游 LLM 发起流式请求,携带同一个 ctx req, _ := http.NewRequestWithContext(ctx, "POST", llmURL, body) resp, err := httpClient.Do(req) if err != nil { return } defer resp.Body.Close() // 2) 监听 ctx.Done(),及时收尾 go func() { <-ctx.Done() log.Println("客户端断开,context 取消,上游 LLM 流将关闭") }() // 3) 边读上游边写给客户端;客户端断开时 resp.Body.Read 立即报错 buf := make([]byte, 4096) for { select { case <-ctx.Done(): // 客户端断开 cleanupSession(r) // 清理 return default: n, err := resp.Body.Read(buf) if n > 0 { w.Write(buf[:n]); flusher.Flush() } if err != nil { // 上游结束或 ctx 取消导致读失败 cleanupSession(r); return } } } }
httpClient.Do(req) 用的是带 ctx 的请求,ctx 被取消时,底层的 HTTP/2 连接会被关闭,LLM 端停止生成。你不需要手动去「取消模型」——和 WebFlux 的 cancel 传播同构。defer resp.Body.Close() 在 ctx 取消时也会释放连接,正好对应「清理层」。
PHP(常驻进程,如 Swoole / FrankenPHP / Workerman)才能真正感知断开;传统 FPM 模式下脚本随请求结束,无法在客户端断开后继续「取消上游」。现代方案用 Swoole 的 Server->close 事件或 connection_aborted() + 心跳检测。
// 伪代码:基于 Swoole 常驻 worker 的 SSE 流式 use Swoole\Http\Response; function streamChat(Request $req, Response $resp) { $resp->header('Content-Type', 'text/event-stream'); $resp->header('Cache-Control', 'no-cache'); $fd = $req->fd; // 1) 用 cURL(CURLOPT_TIMEOUT_MS)或 Swoole Coroutine\Http\Client 发起 LLM 流式请求 $client = new Co\Http\Client($host, 443, true); $client->set(['timeout' => 120]); $client->post('/v1/chat/stream', $body); // 2) 边读上游边写给客户端;每次写入检测客户端是否还在 while ($chunk = $client->recv()) { // 关键检测:客户端连接是否已断开 if (!Server::exist($fd) || Server::getClientInfo($fd) === false) { log("客户端断开,关闭上游 LLM 连接"); $client->close(); // 取消模型调用 cleanupSession($sessionId); // 清理 return; } $resp->write("data: $chunk\n\n"); $resp->flush(); } $client->close(); cleanupSession($sessionId); }
connection_aborted() 只在有输出 flush 后才灵敏,所以必须频繁 flush + 心跳,否则要等下次写才发现断开。传统栈下这就是「感知延迟」问题,和 MVC 的 SseEmitter 一样。
FastAPI / Starlette 的 Request.is_disconnected() 会在客户端断开后置为 True;配合 asyncio 的 task.cancel() 或 httpx.AsyncClient 流式读取时客户端断开会抛异常,从而取消上游调用。
from fastapi import Request import httpx, asyncio async def stream_chat(request: Request): async with httpx.AsyncClient(timeout=120) as client: # 1) 向上游 LLM 发起流式请求(异步,可被取消) async with client.stream("POST", llm_url, json=body) as resp: async def gen(): async for chunk in resp.aiter_text(): # 2) 每次产出前检测客户端是否断开 if await request.is_disconnected(): break # 退出生成器 → resp 被关闭 → 上游 LLM 停止 yield f"data: {chunk}\n\n" # 3) 生成器退出即清理 cleanup_session(session_id) return StreamingResponse(gen(), media_type="text/event-stream")
asyncio.create_task(),客户端断开时 task.cancel(),asyncio 会在最近的 await 点抛出 CancelledError,httpx 随即关闭与 LLM 的 TCP 连接——模型停止生成。这和响应式 cancel 传播等价,只是用协程取消表达。注意:is_disconnected() 在没网络写入时不主动触发,建议加心跳或依赖底层读异常。
Java 是参考文章的主场。两个栈做法完全不同:响应式靠 cancel 自动传播,MVC 靠手动 dispose。下面给出两版完整接口。
// 取消信号自动从下游传到上游 WebClient,无需手动取消 LLM @PostMapping(value = "/stream", produces = TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<String>> stream(@RequestBody ChatRequest req) { concurrency.incrementAndGet(); return chatClient.prompt(req.message()).stream().content() .map(c -> ServerSentEvent.<String>builder().data(c).build()) .doOnCancel(() -> log.warn("前端断开,取消模型 | {}", req.sessionId())) .doFinally(sig -> { concurrency.decrementAndGet(); if (sig == SignalType.CANCEL) { chatMemory.clear(req.sessionId()); // 清理半截会话 log.warn("清理会话 | {}", req.sessionId()); } }); }
SseEmitter emitter = new SseEmitter(120_000L); // 必须设超时 Disposable sub = chatClient.prompt(req.message()).stream().content() .subscribe( c -> { try { emitter.send(SseEmitter.event().data(c)); } catch (IOException e) { sub.dispose(); } }, // send 失败=断开 emitter::completeWithError, emitter::complete); emitter.onCompletion(() -> sub.dispose()); emitter.onTimeout(() -> { sub.dispose(); emitter.complete(); }); emitter.onError(e -> sub.dispose());
dispose(),漏一个就是一条烧钱流;③ 操作链里禁止 block()、禁止 mono.subscribe() 独立订阅,否则 cancel 传播被阻断。| 栈 | 感知断开的方式 | 取消上游的机制 | 清理收尾点 |
|---|---|---|---|
| Go (net/http) | r.Context().Done() 自动关闭 | ctx 透传到 http.NewRequestWithContext,底层连接关闭 | defer + ctx.Done 监听 |
| PHP (Swoole) | 连接存在性检测 / 心跳 | 手动 $client->close() | 检测分支里清理 |
| Python (asyncio) | request.is_disconnected() / 读异常 | 退出生成器 / task.cancel() | 生成器退出 / finally |
| Java (WebFlux) | TCP 断开 → cancel 自动传播 | 自动(WebClient 取消) | doFinally |
| Java (MVC) | 下一次 send 失败 / 回调 | 手动 subscription.dispose() | 各回调 + dispose |
| 维度 | SSE | WebSocket |
|---|---|---|
| 方向 | 服务端→客户端 单向 | 双向 |
| 协议 | 基于 HTTP,兼容好 | 独立协议,需握手升级 |
| 自动重连 | 浏览器原生支持 | 需自己实现 |
| 断连感知 | TCP 层即可 | 需心跳 ping/pong |
| 大模型流式 | ✅ 首选(单向推送够用) | 需双向时才用 |
mono.subscribe() 扔在 flatMap 里,cancel 传不到 = 孤儿请求。
结论:核心是响应式链路下的取消传播——TCP 断开会以 cancel 信号自动从下游传到上游,我不用主动感知,但要保证链路不被阻断,并在关键节点做感知和清理。
机制:WebFlux 下 Reactor Netty 检测到连接断开触发下游 cancel,信号沿 Flux 链传播到 WebClient,自动关闭与 LLM 的 HTTP 连接,模型停止生成。我在 doOnCancel 感知,在 doFinally 统一清理会话与并发计数。
兜底:若是 MVC + SseEmitter,则通过三兄弟回调手动 dispose();若是同步调用,用取消表做软取消。加分项:清理放 doFinally 保证必达、心跳应对 CDN、断开只是止损。