云计算百科
云计算领域专业知识百科平台

Spring AI 流式对话踩坑:SSE 已关闭,为什么大模型还在继续生成?

大模型在执行过程中停止的企业实战:从踩坑到解决方案


摘要

最近在企业项目里踩了一个大坑:用户在AI对话过程中点"停止",但模型还在后台疯狂输出。这坑不只是浪费Token,更严重的是——资源泄漏。排查了一圈源码,最后发现真正的问题不在前端,而在架构设计。这篇文章分享完整的解决方案,包括Agent架构设计、流式输出停止机制、以及为什么要这样设计。看完你就能明白:为什么逻辑和架构比代码更重要。


一、为什么会写这篇

最近在做企业知识库项目,用的Spring AI Alibaba + ReactAgent。

有一天客户反馈:

“我点停止了,AI怎么还在回答?”

一开始我以为是前端问题,让前端同事排查。

前端说:

“按钮都disabled了,事件都解绑了,怎么还在输出?”

后来查日志才发现——

SSE连接确实关闭了,但底层的大模型HTTP请求还在跑。

更坑的是:

用户刷新页面、关浏览器、切后台……

服务端完全没感知。

一晚上下来,服务器内存飙到90%,全是僵尸连接。

Bug没解决,CPU温度先上来了。

这个坑我是真的踩过,而且踩得很痛。

所以写篇文章,把方案完整梳理出来。


二、问题复现

2.1 环境信息

组件版本
Spring Boot 3.2.0
Spring AI Alibaba 1.1.2
ReactAgent 官方版本
大模型API 通义千问 / OpenAI
前端 React + SSE

2.2 业务场景

用户通过SSE流式对话,调用大模型。

正常流程:

用户提问 → 意图分析 → Agent执行 → SSE流式返回

问题场景:

用户点停止 → 前端断开连接 → 后端还在跑 → Token继续扣 → 资源泄漏

2.3 问题根因

我们当时的代码长这样:

@PostMapping("/stream")
public SseEmitter chatStream(@RequestBody ChatRequest request) {
SseEmitter emitter = new SseEmitter(0L);

// 异步执行
executor.execute(() -> {
Flux<NodeOutput> stream = agent.stream(request.getMessage());
stream.subscribe(
output -> emitter.send(output),
error -> emitter.completeWithError(error),
() -> emitter.complete()
);
});

return emitter;
}

看似没问题。

能跑。

但埋了雷。

问题在哪?

  • 没有停止接口:前端点了停止,后端根本不知道
  • 没有会话管理:sessionId到Flux的映射没维护
  • 浏览器断连无感知:用户关浏览器,服务端无动于衷
  • 没有兜底清理:异常场景没人管,全成僵尸连接
  • 更关键的是:

    SseEmitter.complete() 只关闭了SSE连接。

    底层的Flux订阅、HTTP连接、大模型API调用——

    全还在跑。


    三、Agent架构设计(让流程可视化)

    在说解决方案之前,先说说我们的Agent架构。

    为什么?

    因为架构决定了停止方案的可实现性。

    3.1 整体架构图

    #mermaid-svg-yvpEFpC4CTuIhXZx{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-yvpEFpC4CTuIhXZx .error-icon{fill:#552222;}#mermaid-svg-yvpEFpC4CTuIhXZx .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-yvpEFpC4CTuIhXZx .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-yvpEFpC4CTuIhXZx .marker{fill:#333333;stroke:#333333;}#mermaid-svg-yvpEFpC4CTuIhXZx .marker.cross{stroke:#333333;}#mermaid-svg-yvpEFpC4CTuIhXZx svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-yvpEFpC4CTuIhXZx p{margin:0;}#mermaid-svg-yvpEFpC4CTuIhXZx .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster-label text{fill:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster-label span{color:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster-label span p{background-color:transparent;}#mermaid-svg-yvpEFpC4CTuIhXZx .label text,#mermaid-svg-yvpEFpC4CTuIhXZx span{fill:#333;color:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx .node rect,#mermaid-svg-yvpEFpC4CTuIhXZx .node circle,#mermaid-svg-yvpEFpC4CTuIhXZx .node ellipse,#mermaid-svg-yvpEFpC4CTuIhXZx .node polygon,#mermaid-svg-yvpEFpC4CTuIhXZx .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-yvpEFpC4CTuIhXZx .rough-node .label text,#mermaid-svg-yvpEFpC4CTuIhXZx .node .label text,#mermaid-svg-yvpEFpC4CTuIhXZx .image-shape .label,#mermaid-svg-yvpEFpC4CTuIhXZx .icon-shape .label{text-anchor:middle;}#mermaid-svg-yvpEFpC4CTuIhXZx .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-yvpEFpC4CTuIhXZx .rough-node .label,#mermaid-svg-yvpEFpC4CTuIhXZx .node .label,#mermaid-svg-yvpEFpC4CTuIhXZx .image-shape .label,#mermaid-svg-yvpEFpC4CTuIhXZx .icon-shape .label{text-align:center;}#mermaid-svg-yvpEFpC4CTuIhXZx .node.clickable{cursor:pointer;}#mermaid-svg-yvpEFpC4CTuIhXZx .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-yvpEFpC4CTuIhXZx .arrowheadPath{fill:#333333;}#mermaid-svg-yvpEFpC4CTuIhXZx .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-yvpEFpC4CTuIhXZx .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-yvpEFpC4CTuIhXZx .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-yvpEFpC4CTuIhXZx .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-yvpEFpC4CTuIhXZx .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-yvpEFpC4CTuIhXZx .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster text{fill:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx .cluster span{color:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-yvpEFpC4CTuIhXZx .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-yvpEFpC4CTuIhXZx rect.text{fill:none;stroke-width:0;}#mermaid-svg-yvpEFpC4CTuIhXZx .icon-shape,#mermaid-svg-yvpEFpC4CTuIhXZx .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-yvpEFpC4CTuIhXZx .icon-shape p,#mermaid-svg-yvpEFpC4CTuIhXZx .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-yvpEFpC4CTuIhXZx .icon-shape .label rect,#mermaid-svg-yvpEFpC4CTuIhXZx .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-yvpEFpC4CTuIhXZx .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-yvpEFpC4CTuIhXZx .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-yvpEFpC4CTuIhXZx :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

    基础设施层

    流会话管理层

    Service层

    前端层

    POST /stream

    POST /stop

    ConcurrentHashMap

    创建 SseEmitter

    1.意图分析

    查询智能体配置

    构建 ChatModel

    返回 IntentType

    2.路由分发

    CHAT

    WRITE_ARTICLE

    注册 Disposable

    注册 Future

    stop sessionId

    doCleanup

    Controller层

    浏览器/客户端

    AiChatController.chatStream

    AiChatController.stopStream

    AiChatServiceImpl.streamChat

    StreamRegistry

    StreamSession

    MySQL

    ChatModelFactory

    StreamHelper

    IntentAnalyzer.analyze

    AgentHandlerFactory

    ChatAgentHandler

    WritingAgentHandler

    取消Flux订阅中断执行线程关闭SSE连接

    3.2 流式对话流程

    SseEmitter

    大模型

    AgentHandler

    IntentAnalyzer

    AiChatServiceImpl

    AiChatController

    前端

    SseEmitter

    大模型

    AgentHandler

    IntentAnalyzer

    AiChatServiceImpl

    AiChatController

    前端

    #mermaid-svg-yGwjDiNjDLP7vodo{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-yGwjDiNjDLP7vodo .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-yGwjDiNjDLP7vodo .error-icon{fill:#552222;}#mermaid-svg-yGwjDiNjDLP7vodo .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-yGwjDiNjDLP7vodo .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-yGwjDiNjDLP7vodo .marker{fill:#333333;stroke:#333333;}#mermaid-svg-yGwjDiNjDLP7vodo .marker.cross{stroke:#333333;}#mermaid-svg-yGwjDiNjDLP7vodo svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-yGwjDiNjDLP7vodo p{margin:0;}#mermaid-svg-yGwjDiNjDLP7vodo .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-yGwjDiNjDLP7vodo text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-yGwjDiNjDLP7vodo .actor-line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-yGwjDiNjDLP7vodo .innerArc{stroke-width:1.5;stroke-dasharray:none;}#mermaid-svg-yGwjDiNjDLP7vodo .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-yGwjDiNjDLP7vodo .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-yGwjDiNjDLP7vodo #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-yGwjDiNjDLP7vodo .sequenceNumber{fill:white;}#mermaid-svg-yGwjDiNjDLP7vodo #sequencenumber{fill:#333;}#mermaid-svg-yGwjDiNjDLP7vodo #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-yGwjDiNjDLP7vodo .messageText{fill:#333;stroke:none;}#mermaid-svg-yGwjDiNjDLP7vodo .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-yGwjDiNjDLP7vodo .labelText,#mermaid-svg-yGwjDiNjDLP7vodo .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-yGwjDiNjDLP7vodo .loopText,#mermaid-svg-yGwjDiNjDLP7vodo .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-yGwjDiNjDLP7vodo .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-yGwjDiNjDLP7vodo .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-yGwjDiNjDLP7vodo .noteText,#mermaid-svg-yGwjDiNjDLP7vodo .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-yGwjDiNjDLP7vodo .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-yGwjDiNjDLP7vodo .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-yGwjDiNjDLP7vodo .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-yGwjDiNjDLP7vodo .actorPopupMenu{position:absolute;}#mermaid-svg-yGwjDiNjDLP7vodo .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-yGwjDiNjDLP7vodo .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-yGwjDiNjDLP7vodo .actor-man circle,#mermaid-svg-yGwjDiNjDLP7vodo line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-yGwjDiNjDLP7vodo :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

    创建 SseEmitter (timeout=0)

    loop

    [流式事件]

    POST /api/ai/chat/stream

    streamChat(request)

    analyze(request)

    创建 IntentAnalyzerAgent.call(message)

    返回 CHAT/WRITE_ARTICLE

    IntentType

    handle(request, emitter)

    查询配置 + 构建模型

    创建 Agent.stream()

    Flux<NodeOutput>

    NodeOutput

    sendOpenAiChunk

    流式响应

    sendOpenAiDone ([DONE])

    3.3 架构设计的核心思路

    这个架构不是拍脑袋想的。

    而是踩了坑之后,一点点优化出来的。

    第一版架构:没有意图分析

    直接调用ChatModel。

    问题:

    用户问"帮我写一篇关于AI的文章"和"今天天气怎么样",走的是同一个处理逻辑。

    结果:

    写文章的需求被当成普通聊天,回答太短。

    第二版架构:有了意图分析,但没有路由

    在Service里写if-else。

    问题:

    代码耦合严重,新增意图要改Service。

    违反开闭原则。

    第三版架构(当前):策略模式 + 工厂模式

    // 新增意图,只需两步:
    // 1. 添加 IntentType 枚举
    // 2. 实现 AgentHandler 接口
    // 无需修改现有代码

    public interface AgentHandler {
    void handle(ChatRequest request, SseEmitter emitter);
    }

    设计模式不是装逼用的。

    而是真真切切能减少Bug的。


    四、大模型停止输出的方案

    4.1 为什么需要停止方案?

    企业场景下,这是刚需:

    场景问题后果
    用户点停止 大模型还在输出 浪费Token,成本增加
    用户关闭浏览器 服务端无感知 资源泄漏,内存泄漏
    用户刷新页面 旧连接没清理 建立新连接,累积消耗
    长时间对话 超时无兜底 僵尸连接,服务崩溃

    Token是真金白银。

    资源泄漏会搞垮服务。

    所以停止方案不是可选,而是必选。

    4.2 核心技术原理

    这个方案的灵魂就一句话:

    利用Reactor的取消传播机制。

    什么意思?

    看调用链:

    ReactAgent.stream()
    └→ OpenAiChatModel.stream(prompt)
    └→ OpenAiApi.chatCompletionStream(request)
    └→ webClient.post()
    └→ retrieve()
    └→ bodyToFlux()

    关键点:

    • OpenAiApi用WebClient发起流式HTTP请求
    • bodyToFlux() 返回 Flux<ChatCompletionChunk>
    • 这个Flux底层绑定了一个TCP连接

    取消传播流程:

    Flux.subscribe() 返回 Disposable

    Disposable.dispose()

    Reactor 向 Flux 链下发 cancel 信号

    WebClient 捕获取消信号

    Reactor Netty 关闭 TCP 连接

    效果:

  • 客户端侧:Flux订阅取消,TCP连接关闭
  • 服务端侧:大模型API检测到客户端断开,停止生成Token
  • 验证方法:

    用Wireshark抓包:

    # 过滤器
    tcp.port == 443 && ip.addr == <LLM_API_IP>

    # 观察到dispose()后出现:
    # – TCP RST (客户端主动重置)
    # – 或 TCP FIN (客户端主动关闭)

    这就是为什么 Disposable.dispose() 能真正停止模型输出。

    不是玄学,是Reactor的设计。

    4.3 方案设计:三层兜底保障

    既然取消能传播,那问题就变成:

    如何在正确的时机调用 Disposable.dispose()?

    我们设计了三层保障:

    #mermaid-svg-mkBDETYLxTDvfslz{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-mkBDETYLxTDvfslz .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-mkBDETYLxTDvfslz .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-mkBDETYLxTDvfslz .error-icon{fill:#552222;}#mermaid-svg-mkBDETYLxTDvfslz .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-mkBDETYLxTDvfslz .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-mkBDETYLxTDvfslz .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-mkBDETYLxTDvfslz .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-mkBDETYLxTDvfslz .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-mkBDETYLxTDvfslz .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-mkBDETYLxTDvfslz .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-mkBDETYLxTDvfslz .marker{fill:#333333;stroke:#333333;}#mermaid-svg-mkBDETYLxTDvfslz .marker.cross{stroke:#333333;}#mermaid-svg-mkBDETYLxTDvfslz svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-mkBDETYLxTDvfslz p{margin:0;}#mermaid-svg-mkBDETYLxTDvfslz .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-mkBDETYLxTDvfslz .cluster-label text{fill:#333;}#mermaid-svg-mkBDETYLxTDvfslz .cluster-label span{color:#333;}#mermaid-svg-mkBDETYLxTDvfslz .cluster-label span p{background-color:transparent;}#mermaid-svg-mkBDETYLxTDvfslz .label text,#mermaid-svg-mkBDETYLxTDvfslz span{fill:#333;color:#333;}#mermaid-svg-mkBDETYLxTDvfslz .node rect,#mermaid-svg-mkBDETYLxTDvfslz .node circle,#mermaid-svg-mkBDETYLxTDvfslz .node ellipse,#mermaid-svg-mkBDETYLxTDvfslz .node polygon,#mermaid-svg-mkBDETYLxTDvfslz .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-mkBDETYLxTDvfslz .rough-node .label text,#mermaid-svg-mkBDETYLxTDvfslz .node .label text,#mermaid-svg-mkBDETYLxTDvfslz .image-shape .label,#mermaid-svg-mkBDETYLxTDvfslz .icon-shape .label{text-anchor:middle;}#mermaid-svg-mkBDETYLxTDvfslz .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-mkBDETYLxTDvfslz .rough-node .label,#mermaid-svg-mkBDETYLxTDvfslz .node .label,#mermaid-svg-mkBDETYLxTDvfslz .image-shape .label,#mermaid-svg-mkBDETYLxTDvfslz .icon-shape .label{text-align:center;}#mermaid-svg-mkBDETYLxTDvfslz .node.clickable{cursor:pointer;}#mermaid-svg-mkBDETYLxTDvfslz .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-mkBDETYLxTDvfslz .arrowheadPath{fill:#333333;}#mermaid-svg-mkBDETYLxTDvfslz .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-mkBDETYLxTDvfslz .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-mkBDETYLxTDvfslz .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-mkBDETYLxTDvfslz .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-mkBDETYLxTDvfslz .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-mkBDETYLxTDvfslz .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-mkBDETYLxTDvfslz .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-mkBDETYLxTDvfslz .cluster text{fill:#333;}#mermaid-svg-mkBDETYLxTDvfslz .cluster span{color:#333;}#mermaid-svg-mkBDETYLxTDvfslz div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-mkBDETYLxTDvfslz .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-mkBDETYLxTDvfslz rect.text{fill:none;stroke-width:0;}#mermaid-svg-mkBDETYLxTDvfslz .icon-shape,#mermaid-svg-mkBDETYLxTDvfslz .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-mkBDETYLxTDvfslz .icon-shape p,#mermaid-svg-mkBDETYLxTDvfslz .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-mkBDETYLxTDvfslz .icon-shape .label rect,#mermaid-svg-mkBDETYLxTDvfslz .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-mkBDETYLxTDvfslz .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-mkBDETYLxTDvfslz .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-mkBDETYLxTDvfslz :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

    浏览器断连

    用户点停止

    连接超时

    会话超时

    触发停止

    触发方式

    onError回调自动触发

    POST /stop API主动触发

    onTimeout回调自动触发

    定时扫描兜底触发

    StreamRegistry.stop

    doCleanup三重释放

    第一层:SseEmitter生命周期回调(最常用)

    emitter.onCompletion(() -> streamRegistry.remove(sessionId));
    emitter.onError(ex -> streamRegistry.stop(sessionId));
    emitter.onTimeout(() -> streamRegistry.stop(sessionId));

    • onCompletion:正常结束,只移除映射
    • onError:浏览器断连/I/O异常,需要释放资源
    • onTimeout:Tomcat连接超时,需要释放资源

    第二层:主动停止API

    @PostMapping("/stop/{sessionId}")
    public Result<Void> stopStream(@PathVariable String sessionId) {
    boolean stopped = streamRegistry.stop(sessionId);
    if (stopped) {
    return Result.success();
    }
    return Result.error(404, "未找到活跃的流式输出");
    }

    前端点击停止按钮,调用这个接口。

    第三层:定时兜底清理

    @Scheduled(fixedRateString = "#{@streamConfigProperties.cleanup.scanInterval.toMillis()}")
    public void cleanupStaleSessions() {
    registry.forEachExpired(Instant.now(), session -> {
    StreamRegistry.doCleanup(session, "兜底清理-超时");
    });
    }

    每30秒扫描一次,清理超时残留。

    为什么需要三层?

    因为现实世界很复杂。

    • 用户可能关浏览器(第一层兜底)
    • 用户可能点停止(第二层兜底)
    • 代码可能抛异常(第三层兜底)

    三层缺一不可。

    4.4 关键设计:SessionId前后端传递

    停止的前提是:后端能找到那个正在运行的Flux。

    怎么找?

    用sessionId作为唯一标识。

    SseEmitter

    StreamRegistry

    AiChatController

    前端

    SseEmitter

    StreamRegistry

    AiChatController

    前端

    #mermaid-svg-7mnFHojNT8sgNrWK{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-7mnFHojNT8sgNrWK .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-7mnFHojNT8sgNrWK .error-icon{fill:#552222;}#mermaid-svg-7mnFHojNT8sgNrWK .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-7mnFHojNT8sgNrWK .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-7mnFHojNT8sgNrWK .marker{fill:#333333;stroke:#333333;}#mermaid-svg-7mnFHojNT8sgNrWK .marker.cross{stroke:#333333;}#mermaid-svg-7mnFHojNT8sgNrWK svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-7mnFHojNT8sgNrWK p{margin:0;}#mermaid-svg-7mnFHojNT8sgNrWK .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7mnFHojNT8sgNrWK text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-7mnFHojNT8sgNrWK .actor-line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-7mnFHojNT8sgNrWK .innerArc{stroke-width:1.5;stroke-dasharray:none;}#mermaid-svg-7mnFHojNT8sgNrWK .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-7mnFHojNT8sgNrWK .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-7mnFHojNT8sgNrWK #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-7mnFHojNT8sgNrWK .sequenceNumber{fill:white;}#mermaid-svg-7mnFHojNT8sgNrWK #sequencenumber{fill:#333;}#mermaid-svg-7mnFHojNT8sgNrWK #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-7mnFHojNT8sgNrWK .messageText{fill:#333;stroke:none;}#mermaid-svg-7mnFHojNT8sgNrWK .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7mnFHojNT8sgNrWK .labelText,#mermaid-svg-7mnFHojNT8sgNrWK .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-7mnFHojNT8sgNrWK .loopText,#mermaid-svg-7mnFHojNT8sgNrWK .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-7mnFHojNT8sgNrWK .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-7mnFHojNT8sgNrWK .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-7mnFHojNT8sgNrWK .noteText,#mermaid-svg-7mnFHojNT8sgNrWK .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-7mnFHojNT8sgNrWK .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7mnFHojNT8sgNrWK .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7mnFHojNT8sgNrWK .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-7mnFHojNT8sgNrWK .actorPopupMenu{position:absolute;}#mermaid-svg-7mnFHojNT8sgNrWK .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-7mnFHojNT8sgNrWK .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-7mnFHojNT8sgNrWK .actor-man circle,#mermaid-svg-7mnFHojNT8sgNrWK line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-7mnFHojNT8sgNrWK :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

    SSE 起始事件携带 sessionId

    前端提取并存储 sessionId

    POST /stream

    创建 sessionId

    sendOpenAiStartChunk(sessionId)

    {"delta": {"session_id": "xxx"}}

    POST /stop/{sessionId}

    stop(sessionId)

    从Map中找到会话

    doCleanup释放资源

    true

    Result.success()

    关键代码:

    后端发送sessionId:

    public void sendOpenAiStartChunk(SseEmitter emitter, String model, long created, String sessionId) {
    Map<String, Object> delta = new ConcurrentHashMap<>();
    delta.put("role", "assistant");
    if (sessionId != null) {
    delta.put("session_id", sessionId);
    }
    // … 发送SSE事件
    }

    前端提取sessionId:

    eventSource.onmessage = (event) => {
    const data = JSON.parse(event.data);
    if (data.choices && data.choices[0].delta.session_id) {
    this.sessionId = data.choices[0].delta.session_id;
    }
    };

    前端停止请求:

    const stopStream = async () => {
    await fetch(`/api/ai/chat/stop/${this.sessionId}`, { method: 'POST' });
    };

    SessionId是连接前后端的桥梁。

    没有它,停止就成了空谈。

    4.5 三重释放策略:doCleanup

    这是方案的核心。

    static void doCleanup(StreamSession session, String reason) {
    String sessionId = session.sessionId();
    log.info("释放资源 sessionId={}, reason={}", sessionId, reason);

    // ① 取消Flux订阅
    if (session.disposable() != null && !session.disposable().isDisposed()) {
    session.disposable().dispose();
    }

    // ② 中断执行线程
    if (session.future() != null && !session.future().isDone()) {
    session.future().cancel(true);
    }

    // ③ 关闭SSE连接
    if (session.emitter() != null) {
    try {
    session.emitter().complete();
    } catch (Exception e) {
    log.warn("关闭SSE连接失败", e);
    }
    }
    }

    为什么要三重?

    步骤资源操作效果
    Disposable (Flux订阅) dispose() 取消信号传播到WebClient,关闭大模型HTTP连接
    Future (执行线程) cancel(true) 中断Agent工作流线程,终止工具执行
    SseEmitter (SSE连接) complete() 通知前端SSE流结束

    少一步都有问题:

    • 只做①:线程还在跑,Agent工作流没停止
    • 只做②:HTTP连接没关,大模型还在输出
    • 只做③:客户端以为结束了,服务端还在消耗资源

    三重释放,确保资源100%释放。


    五、完整代码实现

    5.1 StreamSession(会话记录)

    public record StreamSession(
    String sessionId,
    SseEmitter emitter,
    Disposable disposable,
    Future<?> future,
    Instant createdAt,
    Duration timeout
    ) {
    public boolean isExpired(Instant now) {
    return Duration.between(createdAt, now).compareTo(timeout) > 0;
    }
    }

    用Record,简洁。

    5.2 StreamRegistry(注册表)

    @Slf4j
    @Component
    public class StreamRegistry {

    private final ConcurrentHashMap<String, StreamSession> activeStreams = new ConcurrentHashMap<>();

    /**
    * 注册会话(幂等)
    */

    public boolean register(String sessionId, SseEmitter emitter,
    Disposable disposable, String interfaceKey) {
    Duration timeout = config.getTimeout().getTimeout(interfaceKey);
    StreamSession session = new StreamSession(
    sessionId, emitter, disposable, Instant.now(), timeout);

    // 幂等校验:防止重复注册
    StreamSession existing = activeStreams.putIfAbsent(sessionId, session);
    if (existing != null && existing.disposable() != null && !existing.disposable().isDisposed()) {
    log.warn("重复注册被拒绝 sessionId={}", sessionId);
    return false;
    }

    log.info("会话已注册 sessionId={}", sessionId);
    return true;
    }

    /**
    * 停止会话(释放资源)
    */

    public boolean stop(String sessionId) {
    StreamSession session = activeStreams.remove(sessionId);
    if (session == null) {
    return false;
    }
    doCleanup(session, "主动停止");
    return true;
    }

    /**
    * 移除会话(仅清理Map)
    */

    public void remove(String sessionId) {
    activeStreams.remove(sessionId);
    }

    /**
    * 三重释放
    */

    static void doCleanup(StreamSession session, String reason) {
    log.info("释放资源 sessionId={}, reason={}", session.sessionId(), reason);

    // ① 取消Flux订阅
    if (session.disposable() != null && !session.disposable().isDisposed()) {
    session.disposable().dispose();
    }

    // ② 中断执行线程
    if (session.future() != null && !session.future().isDone()) {
    session.future().cancel(true);
    }

    // ③ 关闭SSE连接
    if (session.emitter() != null) {
    try {
    session.emitter().complete();
    } catch (Exception e) {
    log.warn("关闭SSE连接失败", e);
    }
    }
    }
    }

    5.3 StreamCleaner(定时清理)

    @Slf4j
    @Component
    public class StreamCleaner {

    private final StreamRegistry registry;

    @Scheduled(fixedRateString = "#{@streamConfigProperties.cleanup.scanInterval.toMillis()}")
    public void cleanupStaleSessions() {
    if (registry.activeCount() == 0) {
    return;
    }

    Instant now = Instant.now();
    registry.forEachExpired(now, session -> {
    log.warn("兜底清理:会话超时 sessionId={}", session.sessionId());
    StreamRegistry.doCleanup(session, "兜底清理-超时");
    });
    }
    }

    5.4 Controller层改造

    @PostMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
    public SseEmitter chatStream(@Validated @RequestBody ChatRequest request) {
    String sessionId = ensureSessionId(request);
    SseEmitter emitter = new SseEmitter(0L);

    // 先绑回调
    emitter.onCompletion(() -> streamRegistry.remove(sessionId));
    emitter.onError(ex -> streamRegistry.stop(sessionId));
    emitter.onTimeout(() -> streamRegistry.stop(sessionId));

    // 再注册
    streamRegistry.register(sessionId, emitter, null, "chat");

    // 异步执行
    executor.execute(() -> {
    Flux<NodeOutput> stream = agent.stream(request.getMessage());
    Disposable disposable = subscriber.subscribe(stream);

    // 回填Disposable
    streamRegistry.updateDisposable(sessionId, disposable);

    // 发送sessionId到前端
    streamHelper.sendOpenAiStartChunk(emitter, modelName, created, sessionId);
    });

    return emitter;
    }

    @PostMapping("/stop/{sessionId}")
    public Result<Void> stopStream(@PathVariable String sessionId) {
    boolean stopped = streamRegistry.stop(sessionId);
    if (stopped) {
    return Result.success();
    }
    return Result.error(404, "未找到活跃的流式输出");
    }

    注册顺序非常重要:

    先绑回调 → 再注册 → 最后订阅

    为什么?

    如果先订阅,订阅过程中抛异常,但回调还没绑,就没人清理了。


    六、方案的优点和为什么

    6.1 核心优点

    优点说明为什么
    资源零泄漏 三重释放确保100%清理 每个资源都有释放路径
    Token成本降低 用户点停止,立刻中断大模型 Reactor取消传播到HTTP连接
    用户体验好 前端点停止,后端立刻响应 sessionId映射实现精准定位
    代码解耦 StreamRegistry、StreamCleaner职责单一 单一职责原则
    扩展性强 新增意图只需实现AgentHandler 策略模式 + 工厂模式
    可配置 超时时间、扫描间隔全可配 @ConfigurationProperties
    可验证 日志、Wireshark、集成测试多维度验证 不是玄学,是科学

    6.2 为什么这样设计?

    问题1:为什么不用线程中断?

    线程中断是Java层面的。

    但大模型调用是HTTP请求。

    HTTP请求在Reactor Netty线程池里。

    中断Java线程,不等同于关闭HTTP连接。

    所以要用Reactor的取消机制,而不是Java的线程中断。

    问题2:为什么需要三层兜底?

    因为用户行为不可控。

    • 用户可能点停止
    • 用户可能关浏览器
    • 用户可能刷新页面
    • 代码可能抛异常

    一层兜不住所有场景。

    问题3:为什么注册顺序这么重要?

    看反例:

    // 错误示范
    Disposable disposable = stream.subscribe(...); // 先订阅
    emitter.onError(ex -> registry.stop(sessionId)); // 后绑回调
    registry.register(sessionId, emitter, disposable); // 再注册

    // 问题:订阅过程中如果抛异常,回调还没绑,资源泄漏

    正确顺序:

    // 正确示范
    emitter.onError(ex -> registry.stop(sessionId)); // 先绑回调
    registry.register(sessionId, emitter, null); // 再注册
    Disposable disposable = stream.subscribe(...); // 最后订阅
    registry.updateDisposable(sessionId, disposable); // 回填

    先绑回调,再注册,最后订阅。

    消除竞态窗口。


    👇 三连支持,动力源泉

    如果这篇文章帮你省下了踩坑的时间,欢迎:

    🔹 点赞 —— 让更多人看到这篇干货 🔹 在看 —— 你的认可是我持续输出的动力 🔹 转发 —— 分享给身边正在做AI Agent的朋友

    你的每一个小动作,对我都很重要 ❤️


    🙏 关于作者

    你好,我是 空门技术栈,一个常年和Bug战斗、持续填坑的Java开发者。

    专注分享:

    • ✅ Java / Spring Boot / Spring AI Alibaba 企业级实战
    • ✅ RAG知识库、AI Agent、多智能体协作落地经验
    • ✅ Docker部署、微服务架构、线上问题排查
    • ✅ 偶尔聊聊「如何保住头发」这类程序员终极话题 😂

    不搞水文,不贩卖焦虑,只写能跑通、能落地、能帮你少加班的实战内容。

    关注我,咱们一起少踩坑,多写优雅代码。


    📖 更多干货推荐

    • 告别手动复制接口文档!Apifox MCP + AI 自动测试让开发效率起飞
    • MySQL MCP Server 从零安装到使用实战,AI 直接查询数据库
    • Spring Event 用了三年,同事一句话把我问懵了
    • Java 抽象类(Abstract Class)彻底讲透:从基础到多态实战
    • Spring AI Alibaba 多智能体(Multi-agent)实战:6 大协作模式 + 完整代码
    • Spring AI Alibaba 智能体作为工具实战:别再让主 Agent 当"人肉路由器"了
    • 一文搞懂 Spring AI Alibaba Workflow:10 个实战案例带你彻底掌握 AI 工作流编排
    • RAG 知识库为什么越更新越乱?一文讲透生产级文档更新方案
    • Transformers VS vLLM:大模型部署到底该选谁?从本地运行到生产上线完整解析
    • LangChain Agent终于讲透了:短期记忆、Redis持久化、Middleware企业级实战,一篇带你从入门到生产
    • LangChain 流式输出终于讲透了:6 种 stream_mode 一篇全搞懂

    🤝 项目合作 / 技术咨询

    平时工作之余,也会接一些技术项目和咨询,主要方向:

    ⚔️ 企业级开发

    • Java / Spring Boot 项目开发与重构
    • 微服务架构设计与落地
    • 系统性能调优、线上问题排查

    🤖 AI 应用落地(这是我最近的主力方向)

    • Spring AI Alibaba / RAG / Agent 应用开发
    • 企业私有知识库搭建
    • AI能力接入现有业务系统
    • 大模型本地化部署与调优

    🛠️ 技术顾问 / 疑难Bug排查

    • 项目架构评审与方案设计
    • 线上疑难问题定位解决
    • 技术选型与团队培训

    如果你正遇到以下情况,欢迎找我聊聊:

    • ✅ 想做AI项目,但技术方案拿不准
    • ✅ 项目卡在某个Bug上很久,团队搞不定
    • ✅ 想把AI接入现有业务,不知道从哪下手
    • ✅ 需要靠谱的开发外包或长期技术顾问

    📮 联系渠道(按回复速度排序):

  • 最快:私信空门技术栈
  • 邮件:2929119150@qq.com(请注明来意和具体需求)
  • 一个人踩坑,是事故;一群人踩坑,就是《避坑宝典》。

    —— IT 空门,与诸君共修技术大道 😎

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » Spring AI 流式对话踩坑:SSE 已关闭,为什么大模型还在继续生成?
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!