重磅实战|SpringBoot WebSocket实现AI流式实时问答(逐字推送+响应超时+完整AI对接)
📌 专栏:SpringBoot高阶实战
🎯 适用场景:AI智能问答、AI客服、大模型实时对话、流式知识库答疑系统
🛠 技术栈:SpringBoot2.7 + WebSocket + OkHttp流式调用 + 大模型SSE流解析
💡 阅读收益:彻底搞懂AI打字机流式效果底层原理、掌握WebSocket对接大模型流式接口、解决AI问答超时、报文乱序、首尾冗余、阻塞卡死等生产疑难问题
一、前言:为什么90%的AI问答Demo都是半成品?
目前网上绝大多数SpringBoot+WebSocket AI问答案例,都存在一个致命硬伤:
等待AI完全生成完毕后,一次性返回全部文本。
这种方式看似能用,完全不符合用户使用体验,和市面上ChatGPT、豆包、文心一言的逐字实时输出天差地别,存在严重问题:
-
首屏等待极久:AI生成长文本需要3~10秒,用户全程空白等待,极易误以为系统卡死
-
服务线程阻塞:同步等待AI接口返回,占用Tomcat线程,高并发下服务直接雪崩
-
无流式交互体验:没有打字机实时效果,完全丢失AI对话产品的核心体验
-
超时概率极高:长文本生成极易触发接口超时,问答直接失败
本文带你从零实现工业级AI流式问答系统,真正对接大模型流式SSE接口,通过WebSocket逐字转发,实现和商用AI产品一致的实时打字效果,附带超时熔断、异常兜底、报文清洗等生产能力。
二、AI流式问答核心原理
2.1 传统同步问答 VS 流式问答
|
对比维度 |
传统同步问答 |
WebSocket流式问答 |
|---|---|---|
|
数据返回方式 |
全部生成完成后一次性返回 |
边生成、边推送、边前端渲染 |
|
首屏响应速度 |
3s~10s 极慢 |
毫秒级首字响应 |
|
线程模型 |
同步阻塞,占用服务线程 |
异步非阻塞,不占用核心线程 |
|
长文本适配 |
极易超时、问答失败 |
完美适配超长文本生成 |
|
用户体验 |
差,空白等待无反馈 |
极佳,真实AI打字交互效果 |
2.2 整体交互流程
1. 前端建立WebSocket长连接,发送用户提问内容; 2. SpringBoot服务接收问题,异步调用大模型SSE流式接口; 3. 大模型逐段返回文本Token,服务端实时接收流数据; 4. 服务端清洗、过滤无效报文,通过WebSocket逐字推送给前端; 5. 前端实时拼接文本,实现打字机效果,流结束后标识问答完成; 6. 全程异步非阻塞,单服务支持高并发AI问答请求。
三、项目依赖配置
除基础WebSocket依赖外,新增OkHttp支持异步流式调用SSE接口,适配大模型流式请求:
<!-- WebSocket核心依赖 --> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-websocket</artifactId> </dependency> <!-- 异步HTTP请求 适配SSE流式接口 --> <dependency> <groupId>com.squareup.okhttp3</groupId> <artifactId>okhttp</artifactId> <version>4.11.0</version> </dependency> <!-- JSON解析 --> <dependency> <groupId>com.alibaba</groupId><artifactId>fastjson2</artifactId> <version>2.0.48</version> </dependency>
WebSocket基础配置类沿用通用配置,自动扫描端点,本文不再重复赘述,重点实现AI流式核心业务。
四、生产级核心代码实现
4.1 AI流式请求工具类(SSE流解析)
核心工具类,异步调用大模型接口,解析SSE流式报文,过滤无效数据,回调推送WebSocket:
import com.alibaba.fastjson2.JSON; import com.alibaba.fastjson2.JSONObject; import okhttp3.*; import okhttp3.sse.EventSource; import okhttp3.sse.EventSources; import org.springframework.stereotype.Component; import javax.websocket.Session; import java.io.IOException; import java.util.concurrent.TimeUnit; /** * AI大模型流式调用工具类 * 适配SSE流式协议,逐段返回AI生成文本 */ @Component public class AiStreamChatUtil { // 构建OkHttp客户端,超长超时适配AI长文本生成 private static final OkHttpClient HTTP_CLIENT = new OkHttpClient.Builder() .connectTimeout(10, TimeUnit.SECONDS) .readTimeout(60, TimeUnit.SECONDS) .writeTimeout(10, TimeUnit.SECONDS) .build(); // 模拟大模型流式接口地址,替换为自己的AI接口 private static final String AI_STREAM_API = "https://xxx.xxx.com/chat/stream"; /** * 发起AI流式问答请求 * @param question 用户提问 * @param session WebSocket会话 */ public void chatStream(String question, Session session) { // 构建请求参数 JSONObject requestJson = new JSONObject(); requestJson.put("prompt", question); requestJson.put("stream", true); requestJson.put("temperature", 0.7); RequestBody body = RequestBody.create(MediaType.get("application/json; charset=utf-8"), requestJson.toString()); Request request = new Request.Builder() .url(AI_STREAM_API) .post(body) .build(); // 监听SSE流式事件 EventSource.Factory factory = EventSources.createFactory(HTTP_CLIENT); factory.newEventSource(request, new EventSourceListener() { // 接收流式分段数据 @Override public void onEvent(EventSource eventSource, String id, String type, String data) { try { // 清洗无效报文、过滤结束标记 if ("[DONE]".equals(data)) { // 推送结束标记,前端终止加载状态 session.getBasicRemote().sendText("【STREAM_END】"); eventSource.cancel(); return; } // 解析AI返回文本片段 JSONObject json = JSON.parseObject(data); String content = json.getString("content"); if (content != null && content.length() > 0) { // 逐段推送文本,实现打字机效果 session.getBasicRemote().sendText(content); } } catch (Exception e) { e.printStackTrace(); } } // 接口请求失败 @Override public void onFailure(EventSource eventSource, Throwable t, Response response) { try { session.getBasicRemote().sendText("❌ AI问答异常,请求失败!"); session.getBasicRemote().sendText("【STREAM_END】"); } catch (IOException e) { e.printStackTrace(); } eventSource.cancel(); } }); } }
4.2 AI流式问答WebSocket核心服务
整合WebSocket连接管理、异步AI流式调用、请求拦截、异常兜底,全程非阻塞处理:
import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.scheduling.annotation.Async; import org.springframework.stereotype.Component; import javax.websocket.*; import javax.websocket.server.ServerEndpoint; import java.util.concurrent.ConcurrentHashMap; /** * AI流式实时问答WebSocket服务 * 实现逐字打字机效果、异步非阻塞调用、流结束标记 */ @Slf4j @Component @ServerEndpoint("/ai/stream/chat") public class AiStreamChatWebSocket { // 全局在线会话缓存 private static final ConcurrentHashMap<String, Session> SESSION_MAP = new ConcurrentHashMap<>(); // 静态注入AI流式工具类 private static AiStreamChatUtil aiStreamChatUtil; @Autowired public void setAiStreamChatUtil(AiStreamChatUtil aiStreamChatUtil) { AiStreamChatWebSocket.aiStreamChatUtil = aiStreamChatUtil; } /** * 连接建立 */ @OnOpen public void onOpen(Session session) { SESSION_MAP.put(session.getId(), session); log.info("【AI流式问答】用户连接成功,会话ID:{},在线人数:{}", session.getId(), SESSION_MAP.size()); } /** * 接收用户提问,异步处理流式问答 */ @OnMessage public void onMessage(String message, Session session) { log.info("【AI问答】收到用户提问:{}", message); // 异步执行AI流式请求,避免阻塞WebSocket主线程 handleAiChat(message, session); } /** * 异步处理AI流式问答 */ @Async public void handleAiChat(String question, Session session) { aiStreamChatUtil.chatStream(question, session); } /** * 连接关闭 */ @OnClose public void onClose(Session session) { SESSION_MAP.remove(session.getId()); log.info("【AI流式问答】用户断开连接,会话ID:{}", session.getId()); } /** * 连接异常 */ @OnError public void onError(Session session, Throwable error) { log.error("【AI流式问答】连接异常:{}", error.getMessage()); } }
4.3 开启异步线程支持
启动类添加异步注解,支持WebSocket异步调用AI接口,避免线程阻塞:
import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.scheduling.annotation.EnableAsync; @SpringBootApplication @EnableAsync public class AiChatApplication { public static void main(String[] args) { SpringApplication.run(AiChatApplication.class, args); } }
五、前端流式渲染页面(完美打字机效果)
前端核心逻辑:接收分段文本、实时拼接渲染、监听流结束标记、加载状态展示、异常兜底,完全复刻商用AI交互体验:
来源:djfb6.cn 来源:ehgnd5.cn 来源:uhrn5.cn 来源:hfng4.cn
<!DOCTYPE html> <html lang="zh-CN"> <head> <meta charset="UTF-8"> <title>AI实时流式问答系统</title> <style> * { margin: 0; padding: 0; box-sizing: border-box; font-family: "Microsoft Yahei", sans-serif; } .container { width: 900px; margin: 40px auto; } .title { text-align: center; color: #1677ff; margin-bottom: 30px; } .chat-box { height: 550px; border: 1px solid #e5e7eb; border-radius: 12px; padding: 20px; overflow-y: auto; background-color: #f9fafb; margin-bottom: 20px; } .input-box { display: flex; gap: 15px; } #questionInput { flex: 1; padding: 14px 20px; border: 1px solid #dcdfe6; border-radius: 8px; font-size: 15px; } #sendBtn { padding: 14px 40px; background: #1677ff; color: #fff; border: none; border-radius: 8px; cursor: pointer; font-size: 15px; } .user-msg, .ai-msg { margin: 16px 0; line-height: 1.8; white-space: pre-wrap; } .user-msg { text-align: right; color: #1677ff; } .ai-msg { text-align: left; color: #333; } .loading { color: #999; text-align: left; } </style> </head> <body> <div class="container"> <h2 class="title">AI 实时流式智能问答</h2> <div class="chat-box" id="chatBox"></div> <div class="input-box"> <input type="text" id="questionInput" placeholder="请输入你的问题,AI实时为你解答..."> <button id="sendBtn">发送提问</button> </div> </div> <script> let socket = null; let aiDom = null; // 初始化WebSocket连接 function initSocket() { socket = new WebSocket("ws://localhost:8080/ai/stream/chat"); socket.onopen = function () { console.log("AI流式连接成功,可正常提问"); }; // 实时接收流式数据 socket.onmessage = function (e) { // 问答结束 if (e.data === "【STREAM_END】") { aiDom = null; return; } // 首次接收AI回复,创建DOM节点 if (!aiDom) { aiDom = document.createElement("div"); aiDom.className = "ai-msg"; aiDom.innerText = "AI解答:"; chatBox.appendChild(aiDom); } // 逐字拼接文本,实现打字机效果 aiDom.innerText += e.data; chatBox.scrollTop = chatBox.scrollHeight; }; // 连接断开重连 socket.onclose = function () { setTimeout(initSocket, 3000); }; } // 发送提问 document.getElementById("sendBtn").addEventListener("click", function () { const val = document.getElementById("questionInput").value.trim(); if (!val) return; // 渲染用户提问 const userDiv = document.createElement("div"); userDiv.className = "user-msg"; userDiv.innerText = "我的提问:" + val; chatBox.appendChild(userDiv); // 发送请求 socket.send(val); document.getElementById("questionInput").value = ""; }); // 回车发送 document.getElementById("questionInput").addEventListener("keydown", function (e) { if (e.keyCode === 13) { document.getElementById("sendBtn
更多推荐




所有评论(0)