Spring Boot WebSocket实战:构建实时视频更新通知系统

📅 2026/8/20 10:09:33
Spring Boot WebSocket实战:构建实时视频更新通知系统
最近在开发一个短视频推荐系统时遇到了一个典型问题如何将后端服务生成的新内容通知高效、实时地推送给前端用户并营造一种“围观”和“互动”的氛围传统的轮询Polling方式不仅浪费资源延迟也高用户体验不佳。而 WebSocket 作为一种全双工通信协议正是解决这类实时通知、聊天、弹幕等场景的利器。本文将以一个模拟的“Funk汤姆猫视频更新”通知场景为例手把手带你从零搭建一个完整的 Spring Boot WebSocket 实时消息推送服务。无论你是想为你的应用添加实时功能还是单纯想学习 WebSocket 的集成与使用这篇涵盖原理、搭建、实战、安全及生产级优化的教程都能让你获得一套可直接复用的解决方案。1. WebSocket 核心概念与为何选择它在开始敲代码之前我们有必要搞清楚 WebSocket 是什么以及它为何比传统 HTTP 更适合我们的“视频上线通知”场景。1.1 什么是 WebSocketWebSocket 是一种在单个 TCP 连接上进行全双工通信的网络协议。它于 2011 年被 IETF 标准化为 RFC 6455。与我们熟悉的 HTTP 协议最大的不同在于HTTP无状态、单向通信。客户端发起请求服务器响应后连接立即关闭。每次交互都是独立的。WebSocket有状态、双向通信。在初次握手基于 HTTP 升级建立连接后连接会保持打开状态。此后服务器和客户端可以在任何时候主动向对方发送数据无需重复建立连接。你可以把它想象成一条一直开通的“电话线”双方可以随时说话而不需要每次通话前都先拨号。1.2 为何是实时通知的最佳选择对于“最新视频已上线快来围观”这样的场景需求非常明确实时性视频发布后通知需要几乎无延迟地触达在线用户。低开销系统可能有成千上万的在线用户连接需要高效维持。服务端主动通知的触发权在服务端视频发布而非客户端不断询问。对比几种方案短轮询Short Polling客户端每隔几秒问一次“有新视频吗”。简单但延迟高无效请求多服务器压力大。长轮询Long Polling客户端发起请求服务器持有直到有数据或超时。改善了实时性但每次请求仍要重建连接复杂度高。Server-Sent Events (SSE)服务器可以向客户端单向推送数据。适合通知流但只能是单向且浏览器兼容性略逊于 WebSocket。WebSocket建立一次连接双向自由通信。完美契合低延迟、高并发、服务端主动推送的需求是此类实时交互功能的首选。2. 环境准备与项目初始化我们将使用 Spring Boot 来快速集成 WebSocket这是 Java 生态中最主流的方案。2.1 技术栈与版本说明JDK: 1.8 或更高版本本文使用 JDK 11Spring Boot: 2.7.x 或 3.x.x本文使用 2.7.18与 Spring Boot 3.x 配置略有不同文中会注明构建工具: Maven 或 Gradle本文使用 MavenIDE: IntelliJ IDEA, Eclipse, VS Code 等任选前端测试用: 一个简单的 HTML 页面使用原生 WebSocket API 进行测试。重要版本提示Spring Boot 2.x 和 3.x 在 WebSocket 依赖上完全兼容但如果你使用 Spring Boot 3.x请确保你的 JDK 版本 17。本文代码在两者上均可运行配置通用。2.2 创建 Spring Boot 项目使用 Spring Initializr 或 IDE 的创建向导生成一个基础项目。所需依赖Spring Web: 提供 Web MVC 能力包含 WebSocket 所需的底层支持。Spring Boot DevTools (可选): 开发热重启。你的pom.xml依赖部分应类似如下?xml version1.0 encodingUTF-8? project xmlnshttp://maven.apache.org/POM/4.0.0 xmlns:xsihttp://www.w3.org/2001/XMLSchema-instance xsi:schemaLocationhttp://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd modelVersion4.0.0/modelVersion parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version !-- 或 3.1.5 -- relativePath/ /parent groupIdcom.example/groupId artifactIdwebsocket-demo/artifactId version0.0.1-SNAPSHOT/version namewebsocket-demo/name descriptionDemo project for Spring Boot WebSocket/description properties java.version11/java.version !-- Spring Boot 3.x 请使用 17 -- /properties dependencies !-- Web 支持包含 WebSocket -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- 可选用于JSON消息转换 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-json/artifactId /dependency !-- 开发工具 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-devtools/artifactId scoperuntime/scope optionaltrue/optional /dependency !-- 测试 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-test/artifactId scopetest/scope /dependency /dependencies build plugins plugin groupIdorg.springframework.boot/groupId artifactIdspring-boot-maven-plugin/artifactId /plugin /plugins /build /project项目创建完成后目录结构大致如下websocket-demo ├── src/main/java/com/example/websocketdemo │ ├── config/ # 配置类 │ ├── controller/ # HTTP控制器 │ ├── handler/ # WebSocket处理器 │ ├── interceptor/ # 拦截器可选 │ └── WebsocketDemoApplication.java # 启动类 ├── src/main/resources │ ├── static/ # 静态资源如测试HTML │ └── application.properties └── pom.xml3. WebSocket 核心配置与处理器编写Spring 提供了强大的spring-websocket模块已包含在spring-boot-starter-web中我们通过几个核心组件来启用它。3.1 启用 WebSocket 支持配置类首先创建一个配置类来启用 WebSocket 并注册端点。// 文件路径src/main/java/com/example/websocketdemo/config/WebSocketConfig.java package com.example.websocketdemo.config; import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.config.annotation.EnableWebSocket; import org.springframework.web.socket.config.annotation.WebSocketConfigurer; import org.springframework.web.socket.config.annotation.WebSocketHandlerRegistry; Configuration EnableWebSocket // 关键注解启用WebSocket支持 public class WebSocketConfig implements WebSocketConfigurer { Override public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) { // 注册一个WebSocket处理器指定连接路径和允许的源CORS registry.addHandler(myWebSocketHandler(), /ws/video-notification) .setAllowedOrigins(*); // 生产环境应指定具体域名如 http://yourdomain.com } // 将处理器声明为Spring Bean Bean public MyWebSocketHandler myWebSocketHandler() { return new MyWebSocketHandler(); } }关键点解释EnableWebSocket: 必须的注解告诉 Spring Boot 启用 WebSocket 自动配置。WebSocketConfigurer: 接口用于自定义配置。addHandler(): 将我们待会编写的处理器与一个 URL 路径绑定。客户端将通过ws://localhost:8080/ws/video-notification连接。setAllowedOrigins(*): 设置 CORS跨域资源共享。*表示允许所有来源这在开发测试时方便但生产环境必须替换为具体的前端域名这是重要的安全措施。3.2 处理连接与消息核心处理器处理器是 WebSocket 的核心它负责处理连接建立、接收消息、发送消息和连接关闭等生命周期事件。// 文件路径src/main/java/com/example/websocketdemo/handler/MyWebSocketHandler.java package com.example.websocketdemo.handler; import com.fasterxml.jackson.databind.ObjectMapper; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Component; import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.TextMessage; import org.springframework.web.socket.WebSocketSession; import org.springframework.web.socket.handler.TextWebSocketHandler; import java.io.IOException; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; Component Slf4j public class MyWebSocketHandler extends TextWebSocketHandler { // 使用线程安全的Map来在线保存所有会话。Key可以为用户ID这里先用Session ID演示。 private static final MapString, WebSocketSession SESSIONS new ConcurrentHashMap(); private final ObjectMapper objectMapper new ObjectMapper(); /** * 连接建立成功后被调用 */ Override public void afterConnectionEstablished(WebSocketSession session) throws Exception { String sessionId session.getId(); SESSIONS.put(sessionId, session); log.info(新的WebSocket连接建立Session ID: {}, 当前在线人数: {}, sessionId, SESSIONS.size()); // 可选连接建立后立即发送一条欢迎消息 MapString, Object welcomeMsg Map.of( type, system, content, 连接成功当‘Funk汤姆猫’有新视频时你会第一时间收到通知。 ); sendMessageToSession(session, welcomeMsg); } /** * 处理客户端发送来的文本消息 */ Override protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception { String payload message.getPayload(); log.info(收到来自 Session [{}] 的消息: {}, session.getId(), payload); // 这里可以解析客户端消息根据不同的指令进行不同操作 // 例如客户端可以发送 {type: heartbeat} 作为心跳 try { Map?, ? msgMap objectMapper.readValue(payload, Map.class); String msgType (String) msgMap.get(type); if (heartbeat.equals(msgType)) { // 处理心跳维持连接 MapString, Object pong Map.of(type, pong, timestamp, System.currentTimeMillis()); sendMessageToSession(session, pong); } // 可以扩展其他消息类型... } catch (Exception e) { log.warn(消息解析失败: {}, payload, e); // 可以给客户端返回错误信息 MapString, Object errorMsg Map.of(type, error, content, 消息格式错误); sendMessageToSession(session, errorMsg); } } /** * 连接关闭后被调用 */ Override public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception { String sessionId session.getId(); SESSIONS.remove(sessionId); log.info(WebSocket连接关闭Session ID: {}, 关闭原因: {}, 当前在线人数: {}, sessionId, status, SESSIONS.size()); } /** * 向单个会话发送消息的通用方法 */ public void sendMessageToSession(WebSocketSession session, Object messageObj) throws IOException { if (session ! null session.isOpen()) { String jsonMessage objectMapper.writeValueAsString(messageObj); session.sendMessage(new TextMessage(jsonMessage)); } } /** * 向所有在线会话广播消息核心推送方法 * 模拟“视频发布”通知 */ public void broadcastNewVideoNotification(String videoTitle, String videoUrl) { MapString, Object notification Map.of( type, new_video, title, videoTitle, url, videoUrl, content, 【Funk汤姆猫】最新视频《 videoTitle 》已上线快来围观, timestamp, System.currentTimeMillis() ); String jsonNotification; try { jsonNotification objectMapper.writeValueAsString(notification); } catch (Exception e) { log.error(构建通知JSON失败, e); return; } TextMessage message new TextMessage(jsonNotification); // 遍历所有会话并发送 for (WebSocketSession session : SESSIONS.values()) { try { if (session.isOpen()) { session.sendMessage(message); log.debug(已向 Session [{}] 发送视频通知, session.getId()); } } catch (IOException e) { log.error(向 Session [{}] 发送消息失败, session.getId(), e); // 可以考虑移除已关闭的会话 } } log.info(视频通知广播完成总计尝试发送给 {} 个会话, SESSIONS.size()); } /** * 获取当前在线连接数可用于监控 */ public int getOnlineCount() { return SESSIONS.size(); } }代码深度解析继承TextWebSocketHandler我们处理文本消息所以继承此类。如果处理二进制消息则继承BinaryWebSocketHandler。会话管理 (SESSIONS)使用ConcurrentHashMap存储所有活跃的WebSocketSession。这是实现广播功能的关键。生产环境中如果服务是多实例部署此 Map 无法跨节点共享需要引入 Redis 或消息中间件下文会讨论。生命周期方法afterConnectionEstablished: 连接建立保存会话可发送欢迎语。handleTextMessage: 处理客户端发来的消息。我们演示了如何解析 JSON 并处理“心跳”消息这是维持连接健康、防止被代理服务器断开的常见做法。afterConnectionClosed: 连接关闭清理会话。核心广播方法broadcastNewVideoNotification这是模拟业务逻辑的方法。当“视频发布”事件触发时调用此方法它会构建一个结构化的通知消息JSON格式并遍历所有在线会话进行发送。异常处理在发送消息时进行try-catch防止因某个连接异常导致整个广播循环中断。4. 模拟业务触发视频发布控制器我们需要一个 HTTP 接口来模拟“视频发布”这个业务事件从而触发 WebSocket 广播。// 文件路径src/main/java/com/example/websocketdemo/controller/VideoPublishController.java package com.example.websocketdemo.controller; import com.example.websocketdemo.handler.MyWebSocketHandler; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.web.bind.annotation.PostMapping; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; import java.util.Map; RestController RequestMapping(/api/video) Slf4j public class VideoPublishController { Autowired private MyWebSocketHandler myWebSocketHandler; /** * 模拟视频发布接口 * 请求体示例{title: 汤姆猫的搞笑日常#99, url: /videos/99} */ PostMapping(/publish) public MapString, Object publishVideo(RequestBody MapString, String videoInfo) { String title videoInfo.get(title); String url videoInfo.get(url); if (title null || title.isEmpty() || url null || url.isEmpty()) { return Map.of(success, false, message, 视频标题和URL不能为空); } log.info(收到视频发布请求: title{}, url{}, title, url); // 这里是你的业务逻辑例如保存视频信息到数据库... // 模拟业务处理 try { Thread.sleep(100); // 模拟处理耗时 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } // 关键步骤触发 WebSocket 广播通知所有在线用户 myWebSocketHandler.broadcastNewVideoNotification(title, url); int onlineUsers myWebSocketHandler.getOnlineCount(); log.info(视频发布成功已通知 {} 位在线用户, onlineUsers); return Map.of( success, true, message, 视频发布成功已推送通知, data, Map.of(title, title, url, url), notifiedUsers, onlineUsers ); } }这个控制器很简单它接收一个包含视频标题和URL的 POST 请求在处理完业务逻辑这里用sleep模拟后调用MyWebSocketHandler的广播方法向所有在线 WebSocket 连接发送通知。5. 前端测试页面体验实时通知为了完整演示我们创建一个简单的前端页面来连接 WebSocket 并接收通知。将以下 HTML 文件放在src/main/resources/static/目录下命名为index.html。启动应用后访问http://localhost:8080即可。!DOCTYPE html html langzh-CN head meta charsetUTF-8 meta nameviewport contentwidthdevice-width, initial-scale1.0 titleFunk汤姆猫 - 视频更新实时通知/title style body { font-family: sans-serif; margin: 20px; background-color: #f5f5f5; } .container { max-width: 800px; margin: auto; background: white; padding: 20px; border-radius: 8px; box-shadow: 0 2px 10px rgba(0,0,0,0.1); } h1 { color: #333; } .status { padding: 10px; margin: 10px 0; border-radius: 4px; } .connected { background-color: #d4edda; color: #155724; border: 1px solid #c3e6cb; } .disconnected { background-color: #f8d7da; color: #721c24; border: 1px solid #f5c6cb; } #messageLog { height: 300px; overflow-y: auto; border: 1px solid #ccc; padding: 10px; background-color: #fafafa; margin-top: 20px; font-family: monospace; font-size: 0.9em; } .log-entry { margin-bottom: 5px; padding: 5px; border-left: 3px solid #007bff; } .log-system { border-left-color: #6c757d; } .log-notification { border-left-color: #28a745; font-weight: bold; } .log-error { border-left-color: #dc3545; } button { padding: 10px 15px; margin: 5px; background-color: #007bff; color: white; border: none; border-radius: 4px; cursor: pointer; } button:hover { background-color: #0056b3; } button:disabled { background-color: #6c757d; cursor: not-allowed; } .simulate-area { background-color: #e9ecef; padding: 15px; border-radius: 4px; margin-top: 20px; } input { padding: 8px; margin: 5px; width: 300px; } /style /head body div classcontainer h1 Funk汤姆猫 - 视频更新实时通知测试/h1 p本页面模拟用户端实时接收服务器推送的新视频上线通知。/p div idstatus classstatus disconnected状态未连接/div div button idconnectBtn onclickconnectWebSocket()连接 WebSocket/button button iddisconnectBtn onclickdisconnectWebSocket() disabled断开连接/button button idclearLogBtn onclickclearLog()清空日志/button /div div classsimulate-area h3模拟后台发布新视频 (HTTP API调用)/h3 input typetext idvideoTitle placeholder输入视频标题如汤姆猫的搞笑日常#100 value【Funk汤姆猫】最新舞蹈挑战 input typetext idvideoUrl placeholder输入视频URL value/videos/latest-dance button onclicksimulateVideoPublish()发布视频并推送通知/button psmall点击此按钮会调用后端 /api/video/publish 接口触发全局广播。/small/p /div h3实时消息日志/h3 div idmessageLog/div /div script let socket null; const statusDiv document.getElementById(status); const messageLogDiv document.getElementById(messageLog); const connectBtn document.getElementById(connectBtn); const disconnectBtn document.getElementById(disconnectBtn); // WebSocket 服务器地址根据你的后端地址调整 const wsUrl ws://${window.location.host}/ws/video-notification; function logMessage(type, content) { const entry document.createElement(div); entry.className log-entry log-${type}; const time new Date().toLocaleTimeString(); entry.innerHTML strong[${time}]/strong ${content}; messageLogDiv.appendChild(entry); // 自动滚动到底部 messageLogDiv.scrollTop messageLogDiv.scrollHeight; } function updateStatus(isConnected) { if (isConnected) { statusDiv.textContent 状态已连接 (${wsUrl}); statusDiv.className status connected; connectBtn.disabled true; disconnectBtn.disabled false; } else { statusDiv.textContent 状态未连接; statusDiv.className status disconnected; connectBtn.disabled false; disconnectBtn.disabled true; } } function connectWebSocket() { if (socket socket.readyState WebSocket.OPEN) { logMessage(system, 已经连接了。); return; } logMessage(system, 正在连接 ${wsUrl} ...); socket new WebSocket(wsUrl); socket.onopen function(event) { logMessage(system, ✅ WebSocket 连接成功); updateStatus(true); // 连接成功后可以发送一个初始消息或心跳 sendHeartbeat(); }; socket.onmessage function(event) { try { const data JSON.parse(event.data); console.log(收到消息:, data); switch(data.type) { case system: logMessage(system, 系统消息: ${data.content}); break; case new_video: logMessage(notification, 新视频通知标题: em${data.title}/em); logMessage(notification, ${data.content}); logMessage(notification, a href${data.url} target_blank点击观看/a (模拟链接)); // 在实际应用中这里可以触发页面弹窗、播放提示音等 break; case pong: logMessage(system, ❤️ 心跳响应: ${new Date(data.timestamp).toLocaleTimeString()}); break; case error: logMessage(error, 错误: ${data.content}); break; default: logMessage(system, 未知消息类型: ${JSON.stringify(data)}); } } catch (e) { logMessage(error, 消息解析失败: ${event.data}); } }; socket.onerror function(error) { logMessage(error, WebSocket 错误: ${error}); updateStatus(false); }; socket.onclose function(event) { logMessage(system, ❌ WebSocket 连接关闭。代码: ${event.code}, 原因: ${event.reason || 无}); updateStatus(false); socket null; // 可选尝试重连 // setTimeout(connectWebSocket, 3000); }; } function disconnectWebSocket() { if (socket) { socket.close(1000, 用户主动断开); logMessage(system, 正在断开连接...); } } function sendHeartbeat() { if (socket socket.readyState WebSocket.OPEN) { socket.send(JSON.stringify({ type: heartbeat })); // 每隔30秒发送一次心跳 setTimeout(sendHeartbeat, 30000); } } function simulateVideoPublish() { const title document.getElementById(videoTitle).value || 默认标题; const url document.getElementById(videoUrl).value || /videos/default; logMessage(system, 模拟调用发布接口: title${title}, url${url}); fetch(/api/video/publish, { method: POST, headers: { Content-Type: application/json, }, body: JSON.stringify({ title, url }) }) .then(response response.json()) .then(data { logMessage(system, 后端响应: ${JSON.stringify(data)}); }) .catch(error { logMessage(error, 调用发布接口失败: ${error}); }); } function clearLog() { messageLogDiv.innerHTML ; } // 页面加载后自动连接可选 // window.onload connectWebSocket; /script /body /html这个页面功能完整连接/断开手动控制 WebSocket 连接。消息展示以不同样式展示系统消息、新视频通知、错误信息。心跳机制连接后每30秒发送一次心跳维持连接活性。模拟发布通过一个表单和按钮调用我们写好的/api/video/publish接口触发全局广播。你可以打开多个浏览器标签页模拟多个在线用户体验广播效果。6. 运行与测试启动后端运行WebsocketDemoApplication的main方法。控制台应无报错Spring Boot 启动在默认的 8080 端口。访问测试页打开浏览器访问http://localhost:8080。建立连接点击页面的“连接 WebSocket”按钮状态应变为“已连接”并收到一条欢迎消息。测试广播在页面的“模拟后台发布新视频”区域填写标题和URL点击“发布视频并推送通知”。稍等片刻当前页面以及所有其他已连接的标签页的日志区域会立即收到一条格式化的新视频通知。测试心跳观察日志每隔30秒会有一条“心跳响应”消息。测试多客户端再打开一个浏览器窗口或标签页重复步骤3。然后在第一个页面发布视频观察两个页面是否同时收到通知。至此一个完整的、可运行的 WebSocket 实时通知 demo 就完成了。7. 常见问题与排查思路在实际集成中你可能会遇到以下问题问题现象可能原因排查思路与解决方案前端无法连接ws://...报错WebSocket connection failed1. 后端服务未启动或端口不对。2. WebSocket 端点路径配置错误。3. 防火墙或网络策略阻止。4. 使用了https页面连接ws应使用wss。1. 检查后端日志确认服务已启动在正确端口。2. 核对WebSocketConfig中addHandler的路径前端连接的URL必须完全一致。3. 本地开发可暂时关闭防火墙测试。4. 如果前端是 HTTPSWebSocket 必须使用wss://并且后端需要配置 SSL。连接建立后立即断开状态码 10061. 后端未正确处理Origin头CORS 问题。2. 代理服务器如 Nginx未正确配置 WebSocket 代理。3. 心跳机制缺失连接被中间节点如负载均衡器超时断开。1. 检查setAllowedOrigins确保包含了前端的源。生产环境不要用*。2. 如果经过 Nginx需添加配置proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade;。3. 实现前端心跳机制定期发送消息保活。能连接但收不到广播消息1. 广播逻辑有误未遍历到所有 session。2. Session 在发送时已关闭或无效。3. 消息格式错误前端onmessage未正确处理。4. 多实例部署session 未共享。1. 在broadcastNewVideoNotification方法中加日志检查SESSIONS大小和遍历过程。2. 发送前务必检查session.isOpen()。3. 使用浏览器开发者工具 Network - WS 标签查看原始收到的消息数据。4. 单机测试无此问题多实例需参考下文“生产级考量”。后端报错java.lang.IllegalStateException: Async support must be enabledSpring Boot 版本或配置问题异步处理未启用。确保使用的是 Spring Boot 2.x 或 3.x。如果使用spring-boot-starter-web默认已启用。如果自定义了WebMvcConfigurer检查是否关闭了异步支持。高并发下连接数上不去或内存溢出1. 未限制单个 Session 的缓冲区大小。2. 未妥善处理异常断开导致SESSIONSMap 内存泄漏。1. 在配置中可设置setSendBufferSizeLimit和setSendTimeLimit。2. 确保afterConnectionClosed方法被正确调用并移除 session。考虑使用SessionCleanupTask定期清理无效 session。8. 生产级考量与最佳实践上面的 Demo 适用于学习和单机环境。要上线生产必须考虑以下方面8.1 会话共享多实例部署问题当你的服务通过负载均衡部署了多个实例时用户A连接到实例1用户B连接到实例2。如果视频发布请求打到实例1那么broadcastNewVideoNotification只能通知到实例1上的用户实例2上的用户B收不到通知。解决方案引入一个中心化的消息广播机制。方案一Redis Pub/Sub每个 WebSocket 处理器实例订阅一个共同的 Redis 频道Channel例如channel:video_notification。当视频发布时调用接口的实例向该频道发布Publish一条消息。所有订阅了该频道的实例都会收到消息然后各自向连接在自己身上的用户进行广播。方案二消息队列如 RabbitMQ、Kafka原理类似使用消息队列的广播或工作队列模式确保所有实例都能消费到“视频发布”这个事件。示例Redis Pub/Sub 思路// 1. 在 Handler 中注入 RedisTemplate Autowired private RedisTemplateString, Object redisTemplate; // 2. 在 afterConnectionEstablished 和 afterConnectionClosed 中可以将用户信息非session对象存入/移除一个全局的 Redis Set用于统计。 // 3. 新增一个方法用于监听 Redis 频道 PostConstruct public void init() { // 订阅频道 redisTemplate.getConnectionFactory().getConnection().subscribe((message, pattern) - { String channel new String(message.getChannel()); String body new String(message.getBody()); // 解析 body触发本实例的广播 broadcastNewVideoNotificationLocal(body); }, channel:video_notification.getBytes()); } // 4. 修改发布视频的控制器不再直接调用 handler.broadcast而是向 Redis 频道发布消息 // videoPublishController 中 // myWebSocketHandler.broadcastNewVideoNotification(title, url); // 注释掉 redisTemplate.convertAndSend(channel:video_notification, notificationJson); // 改为发布消息8.2 连接认证与安全问题Demo 中任何知道地址的人都可以连接这很不安全。解决方案在 WebSocket 握手阶段进行拦截和认证。实现HandshakeInterceptorpublic class AuthHandshakeInterceptor implements HandshakeInterceptor { Override public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, WebSocketHandler wsHandler, MapString, Object attributes) throws Exception { // 1. 可以将HTTP Session中的用户信息传递过来 // 2. 更常见的是通过URL参数或Header传递Token例如 ws://...?tokenxxx String token request.getURI().getQuery(); // 简单演示需解析 if (!isValidToken(token)) { // 返回false拒绝握手 return false; } // 认证通过可以将用户ID存入attributes后续在Handler中通过session.getAttributes()获取 attributes.put(userId, extractUserIdFromToken(token)); return true; } // ... afterHandshake 方法 }在配置中注册拦截器registry.addHandler(myWebSocketHandler(), /ws/video-notification) .addInterceptors(new AuthHandshakeInterceptor()) .setAllowedOrigins(https://your-frontend.com);Handler 中使用用户信息在afterConnectionEstablished中可以通过session.getAttributes().get(userId)获取用户ID并用它作为SESSIONSMap 的 key实现精准推送如只推给特定用户组。8.3 心跳、断线重连与容错心跳前端已实现后端也需处理。防止长时间无通信导致连接被运营商或防火墙断开。断线重连前端 WebSocket 的onclose事件中应实现指数退避的重连逻辑。容错广播消息时对每个session.sendMessage()进行try-catch并将发送失败的 session 从 Map 中移除防止后续继续尝试。8.4 性能与可扩展性连接数限制一个实例能承载的连接数有限受内存、线程资源限制。需要监控连接数必要时进行水平扩展。使用SendTo和SubscribeMapping注解对于更简单的广播场景Spring 提供了基于 STOMP 子协议的高级抽象可以简化代码。但 STOMP 会引入额外开销对于自定义协议要求高的场景原生 WebSocket 更灵活。考虑使用 Netty对于超大规模数十万级以上的并发连接Spring Boot 内嵌的 Tomcat/Jetty 可能成为瓶颈可以考虑使用 Netty 直接编写 WebSocket 服务器性能更高。8.5 监控与日志记录连接数、消息收发速率、错误类型。使用 Micrometer 等工具将指标对接 Prometheus Grafana。关键业务日志如视频发布、大规模广播要记录便于问题回溯。从简单的单机 Demo 到支撑生产环境的实时系统WebSocket 的集成只是第一步。围绕它的会话管理、消息广播、安全认证和集群扩展构成了一个健壮的实时通信架构的核心。希望这篇从场景出发贯穿原理、实战到进阶优化的长文能帮助你扎实地掌握这项技术并成功应用到你的“Funk汤姆猫”或任何需要实时交互的项目中去。