在 AI 对话框里敲下一句话,回车,字开始一个一个地往外蹦——不是一次性刷出一整段,而是像有人在打字一样。这个"打字机"效果几乎成了 AI 产品的标配,但它背后藏着的,是一整条横跨浏览器、Java 后端、Python 服务的流式链路。
这篇文章就把这条链路从"用户敲下回车"到"最后一个字打完"的每一跳都拆开讲。你会发现,"打字机"只是最表面的那一层——真正有意思的,是这串字符是怎么越过浏览器、穿过 Nginx、被 Java 中转、再被 Python 里的 Agent 和一串工具加工,最后原路返回还带上账单的。
▲ 图:用户眼里的 AI 对话——打字机、思考步骤、令牌余额,都在这一个屏里
先记住这串接力棒,后面每一节都是其中一跳:
EventSource 连上后端;MessageQaController 收到请求,先做额度预检、建会话、开 SseEmitter;WebClient 转发到 Python 服务,Python 里 ReactAgent 调一堆工具(搜站、查资料、读图片、看天气……)拿到答案;text/event-stream 一截一截吐回来,Java 中转把"状态"和"内容"分拣成两类事件,再推给前端;注意这六跳里有两套传输:SSE(服务器单向推给浏览器,走内容)和 WebSocket(双向,走令牌余额)。一个 AI 对话功能,为什么同时用两种?这正是这篇文章想讲透的第一个问题。
这两套传输的边界很清楚:SSE 是"一次性"的,一条问答一条流,答完就关;WebSocket 是"常驻"的,整个页面生命周期都在。所以内容走 SSE(跟问答同生共死),余额走 WebSocket(跟页面同生共死)。后面每一跳都会围绕这个边界展开。
还有一件事得先说清楚:这条链不是"一次成型"的,而是被线上问题逼着一层一层补出来的——最初可能只是前端 EventSource 直连一个后端接口,然后发现 Nginx 缓冲要关、中间要中转、要心跳、要多标签接力、要记账。后面每一节的"设计",回头看都是"踩坑的产物"。读的时候不妨想:这一层,如果不做,会出什么问题?想通了,你就知道它为什么存在。
前端的第一跳选的是 EventSource,而不是 fetch 配 ReadableStream。为什么?因为 AI 对话的流式是纯单向的——浏览器只负责收,不需要中途往服务器发东西(发请求走普通 HTTP 就行)。单向推送,SSE 是原生最合适的:一个 URL、一个 new EventSource(url),服务器不断吐事件就行,还能自动重连。
const eventSource = new EventSource(url, { withCredentials: true })
currentEventSource = eventSource // 保存实例,用于打断
eventSource.addEventListener('aiStatusUpdate', (event) => {
const dataObj = JSON.parse(event.data)
if (dataObj?.status) {
aiMessagePlaceholder.statusSteps.push(dataObj.status) // "正在思考"的步骤
aiMessagePlaceholder.isThinking = true
}
})
eventSource.addEventListener('aiAnswerChunk', (event) => {
const dataObj = JSON.parse(event.data)
const token = dataObj.token
if (token == null) return
if (dataObj.isSync) {
// 多标签接力时,直接同步缓存全文,不走打字机
fullContent = token
displayedContent = token
return
}
fullContent += token
})
两个自定义事件名值得注意:aiStatusUpdate 是"思考状态",aiAnswerChunk 是"回答内容"。这是后端在转发 Python 流时特意分成两类事件——状态和内容分开,前端才能一边显示"正在搜索资料、正在整理答案"的步骤提示,一边让打字机专心打字,互不干扰。
withCredentials: true 也很关键:SSE 请求默认不带 Cookie,而鉴权靠 Sa-Token(Cookie 里),不带凭证后端直接拒绝。这一行决定了"登录用户"能不能连上。
再往深一层看,EventSource 还有一个常被忽略的优点:自动重连。网络抖一下断了,EventSource 会按协议自动重连,还会带上 Last-Event-ID 头让服务器知道该从哪续。虽然这里的打断逻辑是主动管理连接的(下一节讲),但 SSE 协议自带的这个能力,仍然比手写 fetch 流省心不少。
还有个对比值得说:为什么不用 WebSocket 推内容?因为内容是"一次问答"的产物,用 WebSocket 意味着要自己维护一条常驻双向连接、自己管心跳、自己管重连——而 SSE 把这些都包了。单向推流用 SSE,双向交互才上 WebSocket,这是"用什么传输"的第一性判断。这个站里,聊天室用 WebSocket(因为要双向发消息),AI 对话用 SSE(因为只是单向收),各得其所。
收到 aiAnswerChunk 只是把内容攒进了 fullContent。真正显示到屏幕上,用的是 requestAnimationFrame 驱动的"打字机",每帧决定打几个字:
const typeNext = () => {
if (displayedContent.length < fullContent.length) {
const remaining = fullContent.length - displayedContent.length
// 基于剩余字符动态决定每帧打几个字:保持丝滑,积压多时加速
let charsToAdd = 1
if (remaining > 50) charsToAdd = Math.ceil(remaining / 5)
else if (remaining > 20) charsToAdd = 3
else if (remaining > 5) charsToAdd = 2
displayedContent += fullContent.substring(displayedContent.length,
displayedContent.length + charsToAdd)
aiMessagePlaceholder.content = displayedContent
scrollToBottomRealtime()
typeFrame = requestAnimationFrame(typeNext)
}
}
这段的核心是"动态每帧字数"。如果固定每帧打一个字符,模型吐得快的时候,前端会越来越落后,用户看着像卡住;模型吐得慢的时候,又显得断断续续。所以这里按积压量调速:积压超过 50 个字符,每帧最多打剩余量的五分之一,快速追平;积压少了,回到一帧一两个字的正常节奏。打字机"丝滑"的秘诀,就在这个小小的 charsToAdd 计算里。
scrollToBottomRealtime 每帧都被调用,保证长回答时聊天框始终跟着最新内容往下滚。这类高频操作要注意性能——它内部用的是直接操作 DOM 的滚动,不经过 Vue 响应式,否则每帧一次响应式更新会把页面拖垮。
这里还藏着一个"占位符"的设计:AI 消息在收到第一个字之前,前端就已经在列表里放了一条"正在思考"的占位消息,用 isStreaming 标记。打字机是把 fullContent 一帧帧渲染进这条占位消息里。为什么先占位?因为列表渲染是响应式的,如果你等收到第一个字才把新消息 push 进列表,整个列表会闪一下;先占位再填内容,列表纹丝不动,只有内容在长。这个细节用户说不出哪里好,但对比过的人都知道,"回答没有突然冒出来"的感觉有多顺。
scrollToBottomRealtime 每帧被调用也有讲究:它跳过 Vue 响应式,直接操作 DOM 的 scrollTop。因为每帧一次的响应式更新(触发整个列表重新渲染)足以让低端机掉帧。这里"哪里该用响应式、哪里该直插 DOM"的取舍,是打字机不掉帧的关键——Vue 的响应式是便利,不是万能的,高频操作就得绕开它。
AI 回答到一半,用户想打断——这是聊天产品的刚需。实现靠的是一个保存全局的 currentEventSource:
const stopStreaming = () => {
if (currentEventSource) {
currentEventSource.close() // 关闭 SSE 连接
currentEventSource = null
}
// 把占位的 AI 消息从"流式中"状态里摘出来
const sIndex = streamingMessageIds.value.indexOf(Number(id))
if (sIndex > -1) streamingMessageIds.value.splice(sIndex, 1)
finalizeMessage()
}
streamingMessageIds 是一份"正在流式输出"的消息 ID 列表,页面里可能同时存在多条流式中的消息(比如并发问两个问题)。打断时不仅要把连接关掉,还要把这条消息从流式列表里摘出来、把占位符收尾,否则 UI 上会一直显示"输入中"的假光标。
有个细节:currentEventSource 是模块级变量,而不是绑定在单个消息上。为什么?因为同一时刻只允许一条流在走——如果用户在上一条还在打字机时又发了一条,前端会先打断上一条。这是产品层面的约束,也是实现上"全局只存一个当前流"的原因。
还有个细节:用户主动打断和网络意外断开,是两种完全不同的收尾。主动打断是 stopStreaming——关连接、摘消息、正常收尾,服务器那边可能还在生成,但前端不管了;意外断开是 EventSource 的 onerror——连接没了,currentEventSource 还是旧实例,这时要清空引用、把消息标记为"未完成",等用户下次操作再恢复。这两种情况如果不区分,会出现"打断后消息还挂着输入光标"或者"断线后状态卡死在思考中"的 bug。
服务器侧的打断也值得一提:前端关掉 EventSource,后端 SseEmitter 的 onCompletion 会触发,把 emitter 从会话列表里移除、停掉心跳任务。也就是说,"打断"在前后端是联动的一套动作——前端关连接,后端收到完成事件做清理。中间任何一环漏了,都会留下一个"幽灵连接"占着资源,攒多了迟早把服务器拖垮。
EventSource 连的是后端 /api/rag/qa/message/send/stream。后端用 Spring 的 SseEmitter 来承载这个长连接。这个方法里藏着一个杀不死 SSE 的坑,一上来就先处理它:
response.setHeader("X-Accel-Buffering", "no"); // 让 Nginx 别缓冲 SSE
response.setHeader("Cache-Control", "no-cache");
response.setHeader("Connection", "keep-alive");
response.setContentType("text/event-stream");
response.setCharacterEncoding("UTF-8");
SseEmitter emitter = new SseEmitter(180000L); // 180秒超时
X-Accel-Buffering: no 这行是血泪换来的。Nginx 默认会对上游响应做缓冲——它把 Java 吐出来的流攒起来,攒够一批才转发给浏览器。对普通页面这没问题,但对 SSE 是致命的:流会被 Nginx 卡住,直到攒满缓冲或者连接超时,前端一个字都收不到。X-Accel-Buffering: no 告诉 Nginx"这个响应别缓冲,来一块发一块"。没有这行,前面所有的流式功夫都白搭。
SseEmitter 的生命周期还有三个回调要处理:onCompletion、onError、onTimeout。它们做的事高度一致——停掉心跳任务、把自己从会话的 emitter 列表里移除。为什么不共用一段代码?因为这三个回调触发时,emitter 的状态不同,但清理逻辑确实重复了。这里用重复换来了清晰,是可以接受的——三个回调各自负责一种退出路径,出问题能一眼定位是哪种。
180 秒的超时也有讲究。它比模型最慢的回答时间长,又不会让连接无限挂着。真正的风险在于:如果超时发生在回答中途,onTimeout 会把连接 complete 掉,前端收到 EOF——这时前端已经把已收到的内容显示出来了,剩下的半截就没了。所以心跳和超时要配合:心跳让连接"一直有事干",超时只是最后的兜底。设置超时,不是期待它触发,而是怕它该触发的时候没触发。
SseEmitter 的超时是 180 秒。AI 回答一般几秒到几十秒,180 秒够用,但如果真超时了,后端的 onTimeout 回调会负责清理。这个超时要跟后面的"5 秒心跳"配合着看——心跳的存在,正是为了防止长回答在中间被某种网关超时掐断。
SSE 是长连接,但中间任何一层代理(Nginx、负载均衡)都可能设空闲超时——连接太久没数据,被当僵尸掐掉。AI 回答长的时候,模型思考那几十秒里,浏览器和后端之间是安静的,容易被误杀。解法是后端定时往连接里塞一个"注释":
ScheduledFuture<?> heartbeatTask = heartbeatExecutor.scheduleAtFixedRate(() -> {
try {
emitter.send(SseEmitter.event().comment("ping"));
} catch (Exception ignored) {
}
}, 5, 5, TimeUnit.SECONDS);
每 5 秒发一个 comment("ping")。注意是 comment(注释)而不是 data(数据)——SSE 协议里注释以 : 开头,浏览器解析时直接忽略,不会触发任何 onmessage 或事件监听。它纯粹是"保持连接有流量"的哨兵,前端完全无感。
心跳的 5 秒间隔也不是随便定的。太短,每 5 秒一个 TCP 包,虽然负载小但无谓;太长,比如 30 秒,中间层的空闲超时可能已经掐了连接。5 秒是"既不会让中间层误判空闲,又不至于太频繁"的折中。这个值在不同部署环境下可能要调——如果前面有云负载均衡,它的空闲超时是多少,你的心跳间隔就要明显小于它。X-Accel-Buffering 关的是 Nginx 缓冲,心跳管的是"空闲超时",两者是一对兄弟:一个管"攒着不发",一个管"太静被砍",缺一不可。很多 SSE 应用只处理了缓冲没处理空闲,或者反过来,最后都栽在半路。
有个细节值得记住:心跳的负载必须是浏览器"看不见"的。如果发的是普通 data,前端每个心跳都要多一次事件处理,还得特判忽略;发 comment 则浏览器根本不当作事件,零成本。这是 SSE 心跳和 WebSocket 心跳(那个要发真实消息)的一个区别。
用户在 A 标签问了一个问题,又在 B 标签打开同一个会话——B 标签的 EventSource 也连上了同一个 sessionId。如果 B 从头开始等,它会错过前面已经输出的内容。后端的处理是"接力":
if (sessionIdToEmitters.containsKey(actualSessionId)) {
CopyOnWriteArrayList<SseEmitter> emitters = sessionIdToEmitters.get(actualSessionId);
emitters.add(emitter); // 后到的标签,加到同一个分发列表
// 把已经攒下的内容一次性同步给后到的连接
StringBuilder cache = sessionIdToCache.get(actualSessionId);
if (cache != null && cache.length() > 0) {
Map<String, Object> tokenMap = new HashMap<>();
tokenMap.put("token", cache.toString());
tokenMap.put("isSync", true); // 前端收到 isSync 直接整段显示,不走打字机
emitter.send(SseEmitter.event().name("aiAnswerChunk").data(tokenMap));
}
return emitter;
}
关键在 sessionIdToEmitters(当前会话的所有连接)和 sessionIdToCache(这个会话已经输出的全文缓存)这两张表。新连接加入时,把缓存全文用 isSync: true 一次性发过去,前端收到 isSync 就整段显示、跳过打字机。之后 distributeEvent 会把后续的新内容同时分发给列表里的所有连接——A、B 两个标签同步往下走。
这个设计回答了一个容易被忽略的问题:一个后端会话,可以同时挂多个 SseEmitter。CopyOnWriteArrayList 保证并发安全的遍历,distributeEvent 逐个发送,发送失败的(标签已关)顺手从列表里移除。
这张 sessionIdToEmitters 和 sessionIdToCache 是存在哪的?它们是 MessageQaController 里的静态 Map——也就是挂在 JVM 上的。这意味着一个会话的"进行中的流"信息,所有请求线程都能访问到,这是多标签同步的前提。但也要注意它的生命周期:会话结束后,这两条 Map 里的条目如果不清理,会一直占着内存。所以 onCompletion 里除了移除 emitter,还要在最后一个 emitter 离开时把整个会话的条目删掉。静态 Map + 手动清理,是这个设计里"记得收尾"的地方——只想着建,忘了清,内存迟早给你颜色看。
后端拿到问题后,真正干活的是 Python 服务。Java 用 WebClient(WebFlux 的响应式客户端)去消费 Python 的 SSE 流,把结果转成一个一个的 type 事件:
ragService.chatStream(userId, actualSessionId, reqContent, saToken, model, result -> {
String type = (String) result.get("type");
String chunkContent = (String) result.get("content");
if ("token".equals(type)) {
// 内容块 → 推给前端打字机
Map<String, String> data = new HashMap<>();
data.put("token", chunkContent);
distributeEvent(actualSessionId, "aiAnswerChunk", data);
} else if ("status".equals(type)) {
// 思考状态 → 推给前端"正在思考"提示
Map<String, String> data = new HashMap<>();
data.put("status", chunkContent);
distributeEvent(actualSessionId, "aiStatusUpdate", data);
}
}, ...);
这一段是整条链路的"分拣台":Java 把 Python 吐上来的原始流,按 token 和 status 两种类型,分别转成前端能识别的两个事件。distributeEvent 再把它广播给这个会话的所有连接(多标签)。
为什么要 Java 中转而不是前端直连 Python?两个原因。第一,鉴权在 Java——Python 服务在独立网络里,不能直接暴露给公网;前端只认识 Java,Java 再把凭据 saToken 传给 Python。第二,Java 要记账——流结束时要算 token、写数据库、更新会话上下文,这些必须发生在后端。
WebClient 的初始化还有一个细节:增大默认缓冲区到 10MB。因为 AI 的回答里可能夹带音频(TTS 生成的语音),长文本合成音频流的体积可能超过 WebClient 默认的 256KB 缓冲上限,超了会直接抛错。这行注释是踩坑留下的:"防止长文本合成音频流超出 256KB 限制报错"。流式传输的每一层,都要考虑"单块数据到底能多大",这是流式架构和普通接口一个隐蔽的区别——普通接口的响应是一整个对象,流式接口的一块数据可能只是整个响应的一小段,但那一小段就可能超出默认缓冲。
在调 Python 之前,Java 还会做两件上下文工程:一是把会话历史(短期记忆)转成消息列表;二是取"长时记忆"LTM 和用户画像。LTM 是把过去多轮会话里值得记住的要点单独存的一份摘要,跨会话生效;用户画像是根据历史行为生成的用户特征描述。这两样东西拼进提示词,让模型的回答"记得你是谁"。一个"记得你"的 AI 和"每次从零开始"的 AI,体验差距是巨大的,而这份差距就藏在这段 Java 代码里——它不是魔法,是把上下文工程老老实实做在前端调用之前。
Python 那边才是真正的"大脑"。ReactAgent 用 LangChain 的 create_agent 搭了一个 ReAct 风格 Agent,挂着一串工具:
from agent.tools import available_tools
agent = create_agent(
model=chat_model,
tools=available_tools,
system_prompt=system_prompt,
)
available_tools 是这一串"手":搜索站点(search_site)、读取我的个人数据(get_my_personal_data)、RAG 总结(rag_summarize)、分析图片(analyze_image)、联网搜索(tavily_search_tool)、查时间、播报天气、搜 Pexels 素材、查询系统状态……几十个工具。Agent 的推理过程就是"判断该用什么工具 → 调用 → 拿到结果 → 再判断",这就是 ReAct 循环。
还有个细节:系统提示词(system_prompt)是从文件读的,模型可以动态切换(get_agent(model_name) 按模型名缓存不同的 Agent 实例)。也就是说,同一个 Agent 逻辑,换个模型就是换一个"脑子",工具不变。这个"多模型动态切换"是这个 AI 服务的一个卖点,也是后面令牌倍率的伏笔——不同模型,计费倍率差十倍以上。
Python 的流式生成器把 Agent 的思考过程暴露成事件:Agent 输出里带 __STATUS__: 标记的行被识别成状态,__TOKEN_USAGE__: 标记的行被识别成 token 消耗。这些标记是 Agent 和流式层之间约定的"暗号",靠它们在文本流里夹带结构化信息。
工具是怎么工作的?以 search_site 为例:Agent 判断"用户可能想问站内的事",就把搜索词作为参数调用这个工具,工具返回站内搜索结果,Agent 再基于结果组织回答。这就是 ReAct 循环的"思考-行动-观察"。工具和 Agent 的耦合点是"工具描述"——每个工具注册时带的 docstring,就是 Agent 决定"什么时候用哪个工具"的依据。工具描述写得好不好,直接决定 Agent 用不用它,这是 prompt 工程在工具层的体现。一个工具,参数说明写得含糊,Agent 就可能不会调用或者传错参数,等于这个工具白注册。
还有一个细节:系统提示词从文件读,而不是写死在代码里。好处是改提示词不用改代码、不用重新部署服务,改个文件就行。对一个要频繁调 prompt 的 AI 服务来说,这是运维上的巨大便利——你可以在不打断线上服务的前提下,反复调整"这个 AI 的人格设定"。AI 产品的迭代,很大一部分不是改代码,是改这段文本。
AI 回答最大的体验问题是"黑盒"——用户不知道它到底在干嘛。这条链路用 aiStatusUpdate 事件把它撕开了一条缝:Agent 每走一步(搜索资料、整理答案),就把进度作为一个状态推给前端,前端展示"正在搜索资料""正在整理答案"这样的步骤提示,配合一个"思考中"的动画。
if (dataObj?.status) {
aiMessagePlaceholder.statusSteps.push(dataObj.status)
aiMessagePlaceholder.isThinking = true
}
这一步的价值不能低估。没有它,用户面对的是一个静止的"输入中",等久了就以为死了,狂点重发;有了它,用户能看到"它正在搜索、正在总结"的推进过程,等待的耐心指数级提升。把思考过程可视化,是 AI 产品体验的隐形杠杆——实现成本不高,收益巨大。
为什么 AI 产品的"思考中"提示这么重要?这背后是"感知延迟"心理学:人对"等待"的耐心,取决于等待时能否看到进展。同样是 10 秒,面对一个静止的加载圈和面对"正在搜索资料→正在整理答案→正在生成回复"的步骤推进,后者的体感时间短得多。这个功能不增加任何"实际价值",但它大幅改善了"感知价值",是所有 AI 产品性价比最高的 UX 投资之一。
还有个实现细节:statusSteps 是个数组,一步步往里推,UI 上渲染成一条时间线。为什么用数组而不是覆盖式地只显示最后一步?因为"步骤的推进感"本身就传递"它在认真工作"的信号——看到三步走完,比一直停在第一步更有说服力。UI 上的这些小心思,和前面的工程细节一样,都是"让用户觉得它真的在做事"的一部分。
这条链路最"现实"的一环是记账。AI 每次调用都有成本,这个站把它做成了"配额制":用户有额度,用完要升级。额度检查在流式开始前就做:
try {
aiTokenRecordService.checkTokenQuota(userId); // 预检额度,超出直接结束流
} catch (BusinessException e) {
// 发一条系统提示,正常结束流,而不是报错
...
emitter.complete();
return emitter;
}
额度是"动态窗口"的,用 Redis 的 ZSet 实现:5 小时内和一周内两个滚动窗口,各算各的用量,超了就拒。这种"滚动窗口 + ZSet"比简单计数器先进——它天然支持"滑出窗口的旧消耗自动作废",不用定时任务清数据。
真正值钱的是流结束时的真实计量:
int consumeToken = (totalTokens != null && totalTokens > 0) ? totalTokens
: (reqContent.length() + finalContent.length()); // 拿不到真实值就按长度估
int finalConsumeToken = (int) Math.ceil(consumeToken * getModelTokenMultiplier(model));
aiTokenRecordService.recordTokenUsage(userId, finalConsumeToken);
getModelTokenMultiplier 是个一票否决般的存在:
private double getModelTokenMultiplier(String model) {
if (model == null) return 1.0;
String lowerModel = model.toLowerCase();
if (lowerModel.contains("deepseek-v4-pro")) return 16.36;
if (lowerModel.contains("qwen3.5-plus")) return 2.55;
if (lowerModel.contains("deepseek-v4-flash")) return 1.36;
return 1.0;
}
同一个问题,用 Pro 模型消耗的 token 是普通模型的 16 倍多。这个倍率不是拍脑袋——它大致对应不同模型的实际单价。一个"输入几个字"的动作,因为选的模型不同,可能消耗几十到上千个"计量 token"。不乘倍率的记账,等于让用户用最便宜的模型价格,享受最贵的模型服务,这是所有 AI 产品都要算清楚的账。
再看额度检查的实现。checkTokenQuota 用的是两个 Redis 键:ai:limit:5h: 和 ai:limit:week:,都是 ZSet。5 小时窗口的 ZSet 里,每个元素是"某次消耗的记录",score 是时间戳。检查时先 removeRangeByScore 把窗口外的旧记录清掉,再 range 求和,和配额比。这套逻辑的精髓是"滚动窗口"——不是固定从月初或今天零点算,而是从"当前时刻往前推 5 小时",谁先消耗谁先滑出窗口,配额自然回收。
配额还分会员档位:普通用户、会员、管理员,各有各的 5 小时上限和每周上限。这套"按身份给额度"的设计,把 AI 成本从"一视同仁的亏损"变成了"分级运营的资源"——普通用户限得紧一点,付费会员给得多一点,成本跟着身份走。
还有一层:recordTokenUsage 记录消耗时,Redis ZSet 的 score 用时间戳,value 用消耗数。查询余额时,getTokenUsage 分别算 5 小时和一周的消耗总和,减去配额就是剩余。整个计量体系没有用一张 MySQL 表做实时加减,而是纯 Redis 滚动窗口——因为"消耗"是高频追加,Redis 的写性能比 MySQL 强一个量级,而且滚动窗口天然带过期语义,不用定时任务清数据。选型不是追求花哨,是"高频写 + 自动过期"这个场景,Redis 就是比 MySQL 合适。
还有一层兜底:如果 Python 没返回真实的 total_tokens,就用"问题长度 + 答案长度"估算。真实计量拿不到的时候,长度估算虽然粗糙,但至少不会零成本。
记账是后端的事,但用户要实时看到余额变化。这里开了第二条连接——一个独立的 WebSocket(/ws/ai-token),专门推令牌用量。后端在用户连上时立刻推一次当前用量,之后每次记账后推一次:
public void sendUsageToUser(String userId) {
WebSocketSession session = SESSIONS.get(userId);
if (session != null && session.isOpen()) {
AiTokenUsageVO tokenUsage = aiTokenRecordService.getTokenUsage(Long.valueOf(userId));
Map<String, Object> response = new HashMap<>();
response.put("type", "ai_token_usage_response");
response.put("data", tokenUsage);
sendMessage(session, response);
}
}
所以一个 AI 对话功能里,SSE 和 WebSocket 各司其职:SSE 单向推内容,WebSocket 双向(虽然这里主要也是推)管余额。为什么不把余额也塞进 SSE?因为 SSE 是"一次问答"的连接,问答结束就关了;余额是"整个页面生命周期"都要显示的东西,需要一条常驻连接。两条连接,两种生命周期,各管各的——这是"为什么同时用两种实时技术"的答案。
前端这边 aiTokenWebSocketService 是单例(和 a19 的聊天列表 WebSocket 一个套路),带心跳、带重连、带 ai_token_usage_response 事件监听,收到就更新页面顶部的余额显示。
为什么余额要单独开一条 WebSocket,而不是在 SSE 流的 done 事件里带一下?因为 SSE 流是一次问答的,答完就关了;而余额是常驻的——用户在页面上看余额、切走再回来、问下一轮,余额要一直对。如果你只在流结束时推一次,用户切到别的页面再回来,余额就是旧的。WebSocket 常驻连接能在"任何时间点"把最新余额推给"还开着的页面",这是它的生命周期决定的。
这条 WebSocket 的心跳和重连,跟聊天那篇(a19)里是同一个套路:单例、指数退避重连、心跳保活。值得注意的一点是,它对 EOF 异常做了特殊处理——客户端正常切页或 App 后台化导致的 EOFException 被降级为 DEBUG 日志,不污染错误告警。这个细节说明:不是所有"断开"都是事故,把正常断开和异常断开分开,日志才有意义。不然每次用户切个页,告警就响一次,早晚把你对日志的信任磨光。
流走完了,不等于结束了。后端在流结束后还要做三件事,顺序还不能错:
// 1. 把占位符消息(内容是"...")替换成真实回答
RagSessionMessage finalMsg = new RagSessionMessage();
finalMsg.setId(aiMessageId);
finalMsg.setContent(finalContent);
ragSessionMessageService.updateById(finalMsg);
// 2. 从回答里提取并保存资源(图片、音频等)
aiResourceService.extractAndSaveResources(finalContent, aiMessageId, userId);
// 3. 把这次的问答存进上下文,供下次提问引用
saveSessionContext(sessionId, updatedContext);
第一件,占位符替换。注意前端显示的是流式渲染的字,但数据库里的占位消息最初只有"...",流结束后才把真实全文写进去——这是为了让中途断掉时,数据库里至少有一条"半成品"而不是空记录。第二件,资源提取——AI 的回答里可能夹着图片 URL、音频 URL,要捞出来单独建资源记录。第三件,上下文保存——下一次提问时,这些历史要拼进消息列表给模型当上下文(还有一层"长时记忆"LTM 专门存跨会话的要点)。
这三件都是"收尾但不迟到"的活:晚了,用户刷新页面看到的是半截占位符;漏了,下一次提问就少了上下文。所以它们都放在 onComplete 回调里,和流式传输串行执行。
这里还有个取舍:资源提取是"同步做"还是"异步做"?如果同步,流结束的 onComplete 回调会等提取完才返回,用户看到的"回答完成"会延迟;如果异步,提取失败就静默丢了。项目选了同步——因为回答里夹的资源(图片、音频)是回答的一部分,宁可慢一点也要保证落库,否则用户刷新后看到的是"回答了但图片丢了"的诡异状态。这个取舍没有标准答案,取决于"资源是核心产出还是点缀"——对这个图片社区来说,AI 生成的图片就是产出本身,所以同步值得。
这条链路的坑,比任何功能都密集。我回头分成三组——连接层的、前端流式的、成本与监控的。每组的问题根源,往往是同一个。
连接层的三跤。 上线初期前端一个字都收不到,浏览器却显示连接是通的——排查半天,是 Nginx 默认缓冲把流攒住了,X-Accel-Buffering: no 一行解决,但那些小时的绝望只有经历过的人才懂。withCredentials 忘写是第二跤:EventSource 不带凭证,后端鉴权直接拒绝,表现为"连上了马上被断开"。第三跤是 HTTPS 下的混合内容:页面 https、EventSource 却是 http,浏览器直接拒。这三跤的修法都是"统一"——头固化成模板、凭证统一加、协议跟着页面走。SSE 上生产就死,先查这三样。
前端流式的纠缠。 打字机初版每帧固定打一个字符,模型吐得快时积压几十上百字符,用户看着像卡死——改成按积压量动态调速才丝滑,这个 charsToAdd 是调出来的不是算出来的。多标签最初各开各的连接,第二个标签打开同一会话历史一片空白——sessionIdToCache 接力是后补的,补完才发现"一个后端会话挂多个 SseEmitter"早该有。还有流结束的竞态:finalizeMessage 和 typeNext 抢,流结束信号先到、打字机还没打完最后一帧,最后一个字就丢了——修法是 finalizeMessage 只在打字机停了才执行。前端流式的坑,全是"状态没对齐"。
成本与监控的教训。 记账最初没乘模型倍率,用户全用 Pro 模型、消耗和 Flash 一个价,成本直接漏——getModelTokenMultiplier 一加,Pro 16.36 倍,账才立住。AI 产品的成本控制,必须落在"真实计量 + 模型定价"两件事上。心跳那跤更隐蔽:看日志以为连接很健康,其实心跳在自动续命,掩盖了上游模型变慢的事实——心跳只证明"TCP 通",不证明"服务没卡"。监控里除了心跳,还得有"业务心跳"(多少秒内必须有内容事件),才抓得住"连接活着但服务僵死"。
把这条 AI 对话链路收起来,其实就是六个动作:EventSource 接流 → SseEmitter 转发 → WebClient 中转 → Agent 加工 → 分拣成状态和内容 → 记账并推送余额。每一个动作单独看都不算难,难的是让它们首尾相接、出错不崩、断点能续。
这里想特别强调两件事。第一,SSE 和 WebSocket 不是二选一,而是按生命周期分工:一次性的内容流用 SSE,常驻的余额推送用 WebSocket。第二,"打字机"只是面子,里子是整条链路每一环的可靠性——Nginx 不缓冲、心跳保活、多标签接力、断点收尾,这些看不见的部分,才是一个 AI 对话"丝滑"的真正原因。
如果哪天你要做一个 AI 对话产品,不妨照着这条链逐环检查:你的流,是从浏览器一路通到模型,中间每一层代理都不缓存、不掐断吗?你的钱,是按模型真实价格算的吗?你的"思考中",是真的让用户看见了吗?这三个问题想清楚,你的 AI 对话就比大多数 demo 强了。
最后说句掏心窝的:这篇拆的东西,单独看每一环都是"别人家也这么做"的常识——EventSource、SseEmitter、WebClient、Agent、Redis 记账。但把它们串成一条链,并且让每一环在极端情况下(断线、超时、多标签、额度耗尽)都不崩,这中间踩过的坑、补过的洞,才是真正值钱的积累。技术不值钱,把技术用对的判断值钱。希望这篇能把一些"判断"传给你。如果这篇讲的是"AI 怎么回答你的问题",那下一篇想聊的就是"你站里的东西怎么被找到"——搜索与热搜。一个是把用户的意图变成答案,一个把站里的内容变成入口,都是同一件事的两面:让该被看到的东西,真的被看到。
最后补一句对"打字机"的感想:它几乎是 AI 产品最表面的那一层,可恰恰是这一层决定了用户愿不愿意等。很多人觉得打字机是给用户"看效果"的装饰,其实它是在给模型一个"边想边说"的缓冲——模型不用憋一整段再吐出来,前端不用一次性渲染一大块。技术上的每一环让步,最后都落在用户那一秒的等待上。做 AI 产品,多想想"用户看到的是什么",比多堆一层架构有用。