
1. 項目概述從零構建一個健壯的WebSocket推送服務最近在重構一個后臺管理系統的實時通知模塊需求很簡單當管理員在后臺操作后前端頁面能立刻收到一條消息提示。聽起來像是WebSocket的典型應用場景對吧但真做起來你會發現遠不止在SpringBoot里加個ServerEndpoint注解那么簡單。消息怎么確保送到連接斷了怎么辦用戶成千上萬時如何高效地給特定人群發消息這些問題都是把一個“玩具級”的WebSocket服務升級為“生產級”推送服務必須跨過的坎。這個項目就是一次完整的實戰記錄。我們不只關注如何用SpringBoot和WebSocket建立連接更要深入解決校驗、心跳、分組這些核心的工程問題。我會帶你從零開始搭建一個具備用戶身份校驗、自動心跳保活、支持用戶分組廣播的WebSocket服務。過程中我會分享那些官方文檔里不會寫的坑比如為什么你的心跳PING-PONG機制總是不生效如何優雅地處理用戶上下線以及在大規模連接下如何做簡單的性能優化。無論你是想為你的應用增加實時能力還是正在被WebSocket的各種不穩定所困擾這篇內容應該都能給你提供一套可直接復用的解決方案。2. WebSocket核心機制與SpringBoot集成選型在動手寫代碼之前我們得先搞清楚WebSocket到底是什么以及為什么在眾多實時通信方案中我們選擇了它。這決定了我們后續所有技術決策的底層邏輯。2.1 WebSocket與HTTP、SSE的對比很多人會把WebSocket和HTTP長輪詢、Server-Sent Events (SSE)搞混。簡單來說HTTP請求就像你每次都要打電話問客服“有新消息嗎”問完就掛斷。長輪詢是電話不掛等客服有消息了再告訴你然后掛斷你再打下一個。這種方式開銷大延遲高。SSE是服務器向瀏覽器單向推送數據的技術它基于HTTP瀏覽器通過一個持久的連接監聽服務器發來的事件流。它的優點是協議簡單天然支持斷線重連。但缺點是單向只能服務器推給瀏覽器并且在一些老式瀏覽器上支持不佳。而WebSocket則是在HTTP握手成功后建立了一個全雙工的TCP長連接。就像你和客服之間拉了一條專線雙方隨時可以主動說話沒有請求-響應的概念。這對于需要頻繁雙向通信的場景如聊天、實時協作、游戲是最高效的。我們的消息推送雖然主要是服務端推但客戶端的心跳確認PONG和可能的業務ACK也需要這個雙向通道。2.2 為什么是SpringBoot 原生WebSocket APISpringBoot集成WebSocket主要有兩種方式一是使用Spring提供的WebSocketHandler抽象和STOMP子協議二是直接使用JSR-356定義的javax.websocket標準API即ServerEndpoint注解。STOMP在Spring生態中很強大它相當于在WebSocket之上定義了一套消息格式和路由規則非常適合復雜的消息代理場景比如結合RabbitMQ或Kafka。但它也帶來了額外的復雜性和學習成本。對于我們這個相對純粹的消息推送服務——核心是連接管理、心跳和分組廣播——STOMP顯得有些“重”了。直接使用ServerEndpoint讓我們能更精細地控制每一個連接的生命周期實現自定義的心跳、校驗邏輯代碼也更直觀。因此我選擇了原生API方案并通過Spring的ServerEndpointExporter來暴露它這樣既能享受Spring的依賴注入又能保持底層控制的靈活性。2.3 項目基礎環境搭建首先創建一個標準的SpringBoot項目這里使用SpringBoot 2.7.x 3.x版本在配置上略有不同但核心邏輯一致。在pom.xml中我們只需要引入WebSocket的starter依賴。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-websocket/artifactId /dependency接下來我們需要一個核心配置類來啟用WebSocket支持。這里的關鍵是ServerEndpointExporterBean它負責將所有帶有ServerEndpoint注解的類注冊為WebSocket端點。import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.web.socket.server.standard.ServerEndpointExporter; Configuration public class WebSocketConfig { /** * 這個Bean會自動注冊使用了ServerEndpoint注解聲明的Websocket endpoint。 * 如果部署在外部容器如Tomcat中容器會自己提供這個Bean可以省略。 */ Bean public ServerEndpointExporter serverEndpointExporter() { return new ServerEndpointExporter(); } }注意如果你將項目打包成WAR包部署到獨立的Tomcat等Servlet容器容器自身會掃描和注冊ServerEndpoint此時再定義ServerEndpointExporter會導致端點被注冊兩次從而引發錯誤。這種情況下你應該移除這個Bean。3. 實現核心WebSocket端點與連接管理有了基礎框架我們來構建最核心的WebSocket服務器端點。這個類將處理所有連接的生命周期事件建立、關閉、錯誤以及消息收發。3.1 定義WebSocket端點類我們創建一個PushWebSocketEndpoint類。使用Component和ServerEndpoint注解將其聲明為一個端點。ServerEndpoint的value屬性定義了客戶端連接的URI路徑。這里有一個非常重要的點WebSocket端點的每個連接都會創建一個新的端點實例。這意味著你不能在類成員變量中直接保存連接狀態如Session因為它們是實例級別的。我們必須使用靜態的ConcurrentHashMap來在全局管理所有連接。import javax.websocket.*; import javax.websocket.server.PathParam; import javax.websocket.server.ServerEndpoint; import org.springframework.stereotype.Component; import java.io.IOException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; Component ServerEndpoint(/ws/push/{userId}) // 在路徑中攜帶用戶ID用于初步標識 public class PushWebSocketEndpoint { // 靜態變量用來記錄當前在線連接數。注意它是線程安全的。 private static final AtomicInteger ONLINE_COUNT new AtomicInteger(0); // 用來存放每個客戶端對應的WebSocketSession對象Key為用戶ID private static final ConcurrentHashMapString, Session SESSION_POOL new ConcurrentHashMap(); // 與某個客戶端的連接會話需要通過它來給客戶端發送數據 private Session session; // 當前連接的用戶ID private String userId; /** * 連接建立成功調用的方法 * param session 可選的參數。session為與某個客戶端的連接會話需要通過它來給客戶端發送數據 * param userId 路徑參數從連接URI中獲取 */ OnOpen public void onOpen(Session session, PathParam(userId) String userId) { this.session session; this.userId userId; // 將當前session存入全局map if (SESSION_POOL.putIfAbsent(userId, session) null) { // 如果之前不存在該用戶的連接則在線數加1 int cnt ONLINE_COUNT.incrementAndGet(); log.info(有新連接加入用戶ID{}當前在線人數為{}, userId, cnt); } else { // 如果用戶已存在連接例如多標簽頁登錄可以選擇踢掉舊連接或拒絕新連接 // 這里我們選擇踢掉舊的建立新的 Session oldSession SESSION_POOL.get(userId); if (oldSession ! null oldSession.isOpen()) { try { oldSession.close(new CloseReason(CloseReason.CloseCodes.NORMAL_CLOSURE, 新連接建立舊連接被關閉)); } catch (IOException e) { log.error(關閉舊連接異常, e); } } SESSION_POOL.put(userId, session); log.info(用戶ID{} 已存在連接已替換為新連接, userId); } // 連接建立后可以主動發送一條歡迎消息 sendMessage(session, 連接WebSocket服務器成功); } // 省略其他方法... }3.2 處理連接關閉與異常連接關閉是常態我們必須妥善處理及時清理資源避免內存泄漏。/** * 連接關閉調用的方法 */ OnClose public void onClose() { if (this.userId ! null SESSION_POOL.remove(this.userId, this.session)) { // 從map中成功移除當前session后在線數減1 int cnt ONLINE_COUNT.decrementAndGet(); log.info(有一連接關閉用戶ID{}當前在線人數為{}, this.userId, cnt); } // 可以在這里觸發一些業務邏輯比如通知該用戶的好友“該用戶已下線” } /** * 發生錯誤時調用 * param session * param error */ OnError public void onError(Session session, Throwable error) { log.error(WebSocket發生錯誤用戶ID{}, this.userId, error); // 通常錯誤也會導致連接關閉onClose方法會被調用所以這里主要做日志記錄。 }3.3 實現消息發送工具方法我們需要一個公共的、線程安全的發送消息方法。因為WebSocket的Session.getBasicRemote().sendText()方法是同步的在并發下可能有問題我們應該使用AsyncRemote進行異步發送并處理可能的異常。/** * 發送消息給指定用戶 * param userId 用戶ID * param message 消息內容 */ public static void sendMessageToUser(String userId, String message) { Session targetSession SESSION_POOL.get(userId); if (targetSession ! null targetSession.isOpen()) { sendMessage(targetSession, message); } else { log.warn(用戶ID{} 不在線或連接已關閉消息發送失敗: {}, userId, message); // 這里可以結合業務將消息存入數據庫或消息隊列待用戶上線后推送 } } /** * 群發消息給所有在線用戶 * param message 消息內容 */ public static void broadcastMessage(String message) { SESSION_POOL.forEach((uid, session) - { if (session.isOpen()) { sendMessage(session, message); } }); } /** * 內部使用的異步發送方法封裝了異常處理 * param session 目標會話 * param message 消息內容 */ private static void sendMessage(Session session, String message) { try { // 使用異步發送避免阻塞業務線程 session.getAsyncRemote().sendText(message); } catch (Exception e) { log.error(發送WebSocket消息失敗Session ID: {}, session.getId(), e); } }4. 用戶身份校驗從路徑參數到Token鑒權在OnOpen方法中我們通過路徑參數{userId}拿到了用戶標識。但這存在嚴重的安全風險任何知道URL格式的人都可以偽裝成其他用戶建立連接。因此路徑參數僅用于初步路由和標識絕不能作為身份驗證的依據。真正的身份校驗應該在連接建立時的握手階段完成。WebSocket握手是基于HTTP的我們可以在連接URI中攜帶Token如JWT并在服務端進行驗證。4.1 客戶端連接時攜帶Token前端連接時不能簡單地用new WebSocket(“ws://localhost:8080/ws/push/123”)。更安全的做法是將Token放在查詢參數中。// 前端示例 const userId 123; const token eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9...; // 你的JWT Token const ws new WebSocket(ws://localhost:8080/ws/push?token${token}userId${userId});4.2 服務端握手攔截與校驗在JSR-356中我們可以實現一個ServerEndpointConfig.Configurator來攔截握手過程。通過重寫modifyHandshake方法我們可以拿到HTTP請求的HandshakeRequest對象從中提取Token并進行驗證。首先創建一個配置器import javax.servlet.http.HttpServletRequest; import javax.websocket.HandshakeResponse; import javax.websocket.server.HandshakeRequest; import javax.websocket.server.ServerEndpointConfig; import java.util.List; import java.util.Map; public class TokenHandshakeConfigurator extends ServerEndpointConfig.Configurator { Override public void modifyHandshake(ServerEndpointConfig sec, HandshakeRequest request, HandshakeResponse response) { // 獲取HTTP Servlet請求對象 MapString, Object userProperties sec.getUserProperties(); HttpServletRequest httpServletRequest (HttpServletRequest) request.getHttpSession(); userProperties.put(HttpServletRequest.class.getName(), httpServletRequest); // 從請求參數中獲取token和userId MapString, ListString parameters request.getParameterMap(); ListString tokenList parameters.get(token); ListString userIdList parameters.get(userId); if (tokenList ! null !tokenList.isEmpty() userIdList ! null !userIdList.isEmpty()) { String token tokenList.get(0); String userId userIdList.get(0); // 進行Token驗證這里需要你實現自己的JWT解析和驗證邏輯 boolean isValid validateToken(token, userId); if (isValid) { // 驗證通過將userId放入用戶屬性供OnOpen方法使用 userProperties.put(userId, userId); userProperties.put(token, token); } else { // 驗證失敗可以在這里拋出異常阻止連接建立 throw new IllegalArgumentException(Token驗證失敗); } } else { throw new IllegalArgumentException(連接參數缺失); } } private boolean validateToken(String token, String userId) { // 實現你的JWT驗證邏輯例如使用jjwt庫 // 校驗簽名、過期時間并確認token中的userId與傳入的一致 // 返回true/false // 這里只是一個示例實際需要完整實現 try { // 偽代碼: Jwts.parser().setSigningKey(key).parseClaimsJws(token); // 從claims中取出userId進行比對 return true; // 假設驗證成功 } catch (Exception e) { return false; } } }然后修改我們的端點注解指定使用這個配置器ServerEndpoint(value /ws/push, configurator TokenHandshakeConfigurator.class) public class PushWebSocketEndpoint { // ... OnOpen public void onOpen(Session session, EndpointConfig config) { // 從配置中獲取驗證通過的userId this.userId (String) config.getUserProperties().get(userId); this.session session; // ... 后續連接管理邏輯 } // ... }實操心得Token校驗一定要在握手階段完成。如果在OnMessage方法里才校驗攻擊者已經建立了連接會消耗你的服務器資源。握手階段失敗連接根本不會建立這是最經濟的安全防線。5. 心跳機制PING-PONG實現與連接健康度管理WebSocket連接可能因為網絡波動、代理超時、客戶端崩潰等原因無聲無息地斷開。心跳機制Heartbeat就是用來檢測連接是否依然存活的生命線。其原理是服務端定期向客戶端發送一個PING幀一種特殊類型的WebSocket控制幀客戶端收到后必須回復一個PONG幀。5.1 為什么需要自己實現心跳你可能聽說過WebSocket協議本身有PING/PONG幀。但遺憾的是JSR-356Java WebSocket API并沒有向應用層暴露主動發送PING幀的接口。Session對象的getBasicRemote().sendPing()方法并不存在。底層容器如Tomcat可能會自動處理PING/PONG但這對于應用層是透明的我們無法依賴它來主動探測并處理死連接。因此我們需要在應用層模擬心跳。通常有兩種方式業務消息充當心跳客戶端定期發送一條特定的業務消息如{type:heartbeat}服務端收到后回復。這種方式簡單但混淆了業務和保活邏輯。獨立的PING-PONG協議服務端定期發送PING消息普通文本消息客戶端約定收到后回復PONG消息。我們采用這種方式因為它更清晰。5.2 服務端心跳調度器我們利用Spring的ScheduledExecutorService或Scheduled注解創建一個定時任務遍歷所有連接發送PING并檢查超時。首先定義一個心跳管理類import org.springframework.scheduling.annotation.Async; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Component; import javax.websocket.Session; import java.io.IOException; import java.util.Date; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; Component EnableScheduling public class WebSocketHeartbeatScheduler { // 記錄每個session最后一次收到PONG的時間 private static final MapString, Long LAST_PONG_TIME new ConcurrentHashMap(); // 心跳間隔毫秒 private static final long HEARTBEAT_INTERVAL 30000; // 30秒 // 超時時間毫秒超過此時間未收到PONG則認為連接死亡 private static final long HEARTBEAT_TIMEOUT 60000; // 60秒 /** * 定時發送PING并檢查超時連接 * fixedRate單位是毫秒 */ Scheduled(fixedRate HEARTBEAT_INTERVAL) Async // 使用異步執行避免阻塞調度線程 public void heartbeatCheck() { long now System.currentTimeMillis(); PushWebSocketEndpoint.SESSION_POOL.forEach((userId, session) - { if (session.isOpen()) { try { // 發送PING消息應用層 session.getAsyncRemote().sendText(PING); // 檢查是否超時 Long lastPongTime LAST_PONG_TIME.get(session.getId()); if (lastPongTime ! null (now - lastPongTime) HEARTBEAT_TIMEOUT) { log.warn(用戶ID{} 心跳超時即將關閉連接, userId); session.close(); } } catch (IOException e) { log.error(發送PING或關閉連接失敗用戶ID{}, userId, e); } } }); } /** * 更新收到PONG的時間由端點類在收到PONG消息時調用 * param sessionId WebSocket Session ID */ public static void updatePongTime(String sessionId) { LAST_PONG_TIME.put(sessionId, System.currentTimeMillis()); } /** * 連接關閉時清理記錄 * param sessionId */ public static void removePongRecord(String sessionId) { LAST_PONG_TIME.remove(sessionId); } }5.3 端點類處理PONG消息修改PushWebSocketEndpoint類增加對PONG消息的處理并在連接關閉時清理記錄。Component ServerEndpoint(value /ws/push, configurator TokenHandshakeConfigurator.class) public class PushWebSocketEndpoint { // ... 其他成員變量和方法 /** * 收到客戶端消息后調用的方法 * param message 客戶端發送過來的消息 * param session 可選的參數 */ OnMessage public void onMessage(String message, Session session) { log.debug(收到來自用戶ID{} 的消息: {}, this.userId, message); // 處理心跳回復 if (PONG.equalsIgnoreCase(message.trim())) { WebSocketHeartbeatScheduler.updatePongTime(session.getId()); log.debug(收到用戶ID{} 的心跳回復, this.userId); return; // 心跳消息不進入業務處理 } // 這里是處理其他業務消息的邏輯... // processBusinessMessage(message); } OnClose public void onClose() { // ... 原有的清理邏輯 WebSocketHeartbeatScheduler.removePongRecord(this.session.getId()); // ... } }5.4 客戶端心跳響應前端也需要相應配合在收到服務端的“PING”消息后立刻回復“PONG”。// 前端WebSocket事件監聽 ws.onmessage function(event) { const msg event.data; if (msg PING) { // 立即回復PONG ws.send(PONG); return; } // 處理其他業務消息... console.log(收到業務消息:, msg); };踩坑記錄心跳超時時間HEARTBEAT_TIMEOUT不能設置得太短。因為網絡延遲、客戶端GC暫停都可能導致PONG回復慢。通常設置為心跳間隔的2-3倍是比較合理的。另外一定要在連接關閉時清理LAST_PONG_TIME記錄否則這個Map會一直增長造成內存泄漏。6. 用戶分組與定向消息廣播簡單的全局廣播broadcastMessage在很多場景下并不適用。比如我們只想給“北京地區的用戶”或者“購買了A產品的用戶”發送通知。這就需要分組功能。6.1 設計分組數據結構我們需要一個高效的數據結構來維護“組”和“組內用戶”的關系。考慮到并發性我們繼續使用ConcurrentHashMap。import org.springframework.stereotype.Component; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArraySet; Component public class WebSocketGroupManager { // key: 組名 (例如: “group_admin”, “city_beijing”), value: 該組下的用戶ID集合 private static final ConcurrentHashMapString, SetString GROUP_MEMBERS new ConcurrentHashMap(); /** * 將用戶加入指定組 * param groupName 組名 * param userId 用戶ID */ public void joinGroup(String groupName, String userId) { // computeIfAbsent 是線程安全的如果組不存在則創建一個新的CopyOnWriteArraySet SetString userIds GROUP_MEMBERS.computeIfAbsent(groupName, k - new CopyOnWriteArraySet()); userIds.add(userId); log.info(用戶ID{} 加入組: {}, userId, groupName); } /** * 將用戶移出指定組 * param groupName 組名 * param userId 用戶ID */ public void leaveGroup(String groupName, String userId) { SetString userIds GROUP_MEMBERS.get(groupName); if (userIds ! null) { userIds.remove(userId); log.info(用戶ID{} 離開組: {}, userId, groupName); // 如果組空了可以選擇移除這個組避免內存浪費 if (userIds.isEmpty()) { GROUP_MEMBERS.remove(groupName); } } } /** * 用戶斷開連接時將其從所有組中移除 * param userId 用戶ID */ public void removeUserFromAllGroups(String userId) { GROUP_MEMBERS.forEach((groupName, userIds) - { if (userIds.remove(userId)) { log.info(用戶ID{} 從組 {} 中移除連接斷開, userId, groupName); } if (userIds.isEmpty()) { GROUP_MEMBERS.remove(groupName); } }); } /** * 向指定組的所有在線用戶發送消息 * param groupName 組名 * param message 消息內容 */ public void sendMessageToGroup(String groupName, String message) { SetString userIds GROUP_MEMBERS.get(groupName); if (userIds ! null !userIds.isEmpty()) { userIds.forEach(userId - { // 復用之前寫的單發方法 PushWebSocketEndpoint.sendMessageToUser(userId, message); }); log.info(向組 {} 發送消息組內成員數: {}, groupName, userIds.size()); } else { log.warn(組 {} 不存在或為空消息未發送: {}, groupName, message); } } /** * 獲取某個組的所有在線成員ID快照 * param groupName * return */ public SetString getGroupMembers(String groupName) { SetString members GROUP_MEMBERS.get(groupName); return members null ? new CopyOnWriteArraySet() : new CopyOnWriteArraySet(members); } }6.2 在連接生命周期中管理分組用戶通常在連接建立后通過發送一條“加入組”的指令來訂閱某個分組。我們在PushWebSocketEndpoint中處理這個消息。Component ServerEndpoint(value /ws/push, configurator TokenHandshakeConfigurator.class) public class PushWebSocketEndpoint { // ... 注入GroupManager Autowired private static WebSocketGroupManager groupManager; // 注意這里需要特殊處理靜態注入 // 解決ServerEndpoint類中Autowired靜態成員注入為null的問題 // 通過一個非靜態的setter方法將Spring容器中的Bean賦值給靜態變量 private static WebSocketGroupManager staticGroupManager; Autowired public void setGroupManager(WebSocketGroupManager groupManager) { PushWebSocketEndpoint.staticGroupManager groupManager; } OnMessage public void onMessage(String message, Session session) { // ... 心跳處理邏輯 // 處理業務消息這里假設消息是JSON格式 try { // 使用簡單的JSON解析實際項目建議用Jackson/Gson // 假設消息格式: {type: join_group, group: admin} if (message.contains(\type\: \join_group\)) { // 解析groupName String groupName ...; // 從message中解析出組名 staticGroupManager.joinGroup(groupName, this.userId); sendMessage(this.session, 已成功加入組: groupName); } else if (message.contains(\type\: \leave_group\)) { // 離開組 String groupName ...; staticGroupManager.leaveGroup(groupName, this.userId); sendMessage(this.session, 已離開組: groupName); } else { // 其他業務消息... } } catch (Exception e) { log.error(處理消息失敗用戶ID{}, 消息內容{}, this.userId, message, e); sendMessage(this.session, 消息處理錯誤: e.getMessage()); } } OnClose public void onClose() { // ... 原有的清理邏輯 // 用戶斷開時從所有組中移除 if (staticGroupManager ! null this.userId ! null) { staticGroupManager.removeUserFromAllGroups(this.userId); } // ... } }技術細節ServerEndpoint是由WebSocket容器管理的不是Spring Bean所以其內部無法直接使用Autowired注入Spring管理的Bean。我們通過一個靜態變量和setter方法巧妙地解決了這個問題。這是一種常見的模式。6.3 業務層調用分組廣播現在在任何Spring管理的Bean如Service、Controller中你都可以輕松地向特定組發送消息了。Service public class NotificationService { Autowired private WebSocketGroupManager groupManager; public void notifyAdmins(String message) { groupManager.sendMessageToGroup(group_admin, message); } public void notifyUsersInCity(String city, String message) { String groupName city_ city; groupManager.sendMessageToGroup(groupName, message); } }7. 性能優化與生產環境考量當連接數上升到幾千甚至上萬時最初的簡單實現可能會遇到性能瓶頸。這里分享幾個關鍵的優化點。7.1 連接Session存儲的優化我們之前用ConcurrentHashMapString, Session存儲連接。當需要廣播時會遍歷整個Map。如果連接數巨大例如10萬這個遍歷操作會非常耗時并且會阻塞心跳調度器的線程。優化方案按組存儲Session引用我們可以在WebSocketGroupManager中不僅存儲用戶ID還直接存儲其對應的Session弱引用WeakReferenceSession。這樣在發送組播消息時可以直接獲取到組內的Session集合進行發送無需遍歷全局Map。但要注意處理Session關閉后弱引用被GC的情況需要定期清理無效的引用。更優方案引入本地緩存或分布式方案對于單機可以使用Caffeine或Guava Cache來緩存Session信息并設置合適的過期策略。對于集群環境Session不能存在單機內存中必須引入外部存儲如Redis并配合廣播機制如Redis Pub/Sub來通知集群內所有節點進行消息推送。這會復雜很多通常需要引入Spring的WebSocketMessageBroker和STOMP over RabbitMQ/Kafka。7.2 心跳檢查的優化我們之前的heartbeatCheck方法是遍歷所有Session。當連接數很大時這個循環本身會成為性能熱點并且發送PING是IO操作在單線程中順序執行會非常慢。優化方案分桶與異步化分桶將所有的Session分散到多個“桶”Bucket中每個桶由一個獨立的線程或定時任務負責心跳檢查。這樣可以并行處理提高效率。批量異步發送使用CompletableFuture或反應式編程模型將“發送PING”這個IO操作批量異步執行避免阻塞心跳檢查線程。// 偽代碼分桶心跳檢查思路 Component public class OptimizedHeartbeatScheduler { private static final int BUCKET_COUNT 10; private ListConcurrentHashMapString, Session sessionBuckets; PostConstruct public void init() { // 初始化10個桶 sessionBuckets new ArrayList(BUCKET_COUNT); for (int i 0; i BUCKET_COUNT; i) { sessionBuckets.add(new ConcurrentHashMap()); } } // 根據sessionId的hash值決定放入哪個桶 public void addSession(Session session, String userId) { int bucketIndex Math.abs(userId.hashCode()) % BUCKET_COUNT; sessionBuckets.get(bucketIndex).put(userId, session); } // 啟動10個定時任務每個任務負責一個桶 Scheduled(fixedRate 30000) public void heartbeatBucket0() { checkBucket(0); } // ... 為其他桶也定義類似的任務或者用一個任務循環處理所有桶但使用線程池 }7.3 消息推送的可靠性保證我們的sendMessageToUser方法在用戶不在線時只是打印了警告。在生產環境中這通常不夠。我們需要一個“離線消息”機制。簡單方案持久化到數據庫當發送消息時如果目標用戶不在線將消息存入數據庫的一張offline_message表中包含userIdcontentcreateTime等字段。當用戶重新建立WebSocket連接后在OnOpen方法中查詢該用戶的離線消息并推送然后刪除或標記已發送。進階方案消息隊列對于高并發、高可靠的場景應該引入消息隊列如RocketMQ, Kafka。業務系統將推送事件發送到MQ由一個獨立的推送服務消費MQ該服務負責維護WebSocket連接和發送。這樣實現了業務與推送的解耦并且可以利用MQ的持久化、重試等特性保證消息不丟失。7.4 連接數限制與拒絕服務防護不加限制地允許連接可能導致資源耗盡。我們需要在TokenHandshakeConfigurator的modifyHandshake方法中增加一些防護邏輯。public class TokenHandshakeConfigurator extends ServerEndpointConfig.Configurator { private static final int MAX_CONNECTIONS_PER_IP 50; // 每個IP最大連接數 private static final ConcurrentHashMapString, AtomicInteger IP_CONNECTION_COUNT new ConcurrentHashMap(); Override public void modifyHandshake(ServerEndpointConfig sec, HandshakeRequest request, HandshakeResponse response) { // ... Token驗證邏輯 // 獲取客戶端IP String clientIp getClientIp(request); // 檢查IP連接數 AtomicInteger count IP_CONNECTION_COUNT.computeIfAbsent(clientIp, k - new AtomicInteger(0)); if (count.incrementAndGet() MAX_CONNECTIONS_PER_IP) { count.decrementAndGet(); // 恢復計數 throw new IllegalArgumentException(連接數超限); } // 將IP和計數器引用存入用戶屬性方便連接關閉時遞減 sec.getUserProperties().put(clientIp, clientIp); sec.getUserProperties().put(ipCounter, count); } // 在連接關閉的監聽器中遞減計數需要在Endpoint中獲取并調用 public static void decrementIpCount(EndpointConfig config) { String clientIp (String) config.getUserProperties().get(clientIp); AtomicInteger counter (AtomicInteger) config.getUserProperties().get(ipCounter); if (counter ! null) { counter.decrementAndGet(); if (counter.get() 0) { IP_CONNECTION_COUNT.remove(clientIp, counter); } } } private String getClientIp(HandshakeRequest request) { // 從請求頭中獲取真實IP注意處理代理如X-Forwarded-For // 這里是簡化版 MapString, ListString headers request.getHeaders(); ListString ipHeaders headers.get(X-Forwarded-For); if (ipHeaders ! null !ipHeaders.isEmpty()) { return ipHeaders.get(0).split(,)[0].trim(); } // 否則從HttpServletRequest中獲取 HttpServletRequest req (HttpServletRequest) request.getHttpSession(); return req.getRemoteAddr(); } }然后在PushWebSocketEndpoint的OnClose方法中調用TokenHandshakeConfigurator.decrementIpCount(this.session.getUserProperties())。8. 前端集成示例與常見問題排查服務端準備好了前端如何對接這里給出一個精簡但完整的Vue 3組件示例并附上幾個我踩過的坑。8.1 Vue 3組件示例template div p連接狀態: {{ status }}/p button clickconnect :disabledisConnected連接/button button clickdisconnect :disabled!isConnected斷開/button button clickjoinAdminGroup :disabled!isConnected加入管理員組/button ul li v-for(msg, index) in messages :keyindex{{ msg }}/li /ul /div /template script setup import { ref, onUnmounted } from vue; const ws ref(null); const status ref(未連接); const isConnected ref(false); const messages ref([]); const connect () { const userId user_123; const token your_jwt_token_here; // 應從登錄狀態獲取 const wsUrl ws://${location.host}/ws/push?token${token}userId${userId}; ws.value new WebSocket(wsUrl); ws.value.onopen () { status.value 已連接; isConnected.value true; messages.value.push(WebSocket連接已建立); }; ws.value.onmessage (event) { const msg event.data; if (msg PING) { ws.value.send(PONG); console.log(已回復PONG); return; } messages.value.push(收到: ${msg}); }; ws.value.onerror (error) { console.error(WebSocket錯誤:, error); status.value 連接錯誤; }; ws.value.onclose () { status.value 已斷開; isConnected.value false; messages.value.push(WebSocket連接已關閉); }; }; const disconnect () { if (ws.value) { ws.value.close(); ws.value null; } }; const joinAdminGroup () { if (ws.value ws.value.readyState WebSocket.OPEN) { const joinCmd JSON.stringify({ type: join_group, group: group_admin }); ws.value.send(joinCmd); messages.value.push(已發送加入管理員組請求); } }; // 組件卸載時自動斷開連接 onUnmounted(() { disconnect(); }); /script8.2 常見問題與排查清單連接失敗返回404檢查端點路徑確認前端連接的URL/ws/push與后端ServerEndpoint注解中的value完全一致。檢查配置Bean確認ServerEndpointExporterBean已正確配置。如果項目是SpringBoot內嵌容器必須有這個Bean。檢查跨域如果前端與后端域名/端口不同需要配置CORS。對于WebSocketCORS在握手階段生效。你可以在TokenHandshakeConfigurator的modifyHandshake方法中添加響應頭response.getHeaders().put(“Access-Control-Allow-Origin”, request.getHeaders().get(“Origin”)); 注意生產環境要嚴格限制Origin。連接建立后立刻斷開檢查Token校驗在TokenHandshakeConfigurator的modifyHandshake中拋出的任何異常都會導致握手失敗連接關閉。查看服務器日志確認Token驗證邏輯無誤。檢查Nginx等代理配置如果你使用了Nginx反向代理必須配置其支持WebSocket。關鍵配置如下location /ws/ { proxy_pass http://backend_server; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header Host $host; proxy_read_timeout 3600s; # 長連接超時時間 }心跳不工作連接一段時間后自動斷開確認心跳邏輯被執行在服務端heartbeatCheck方法內打日志看定時任務是否正常觸發。確認PING/PONG消息格式服務端發的是純文本”PING”前端判斷的是event.data ‘PING’。務必注意大小寫和可能的空格。建議前后端統一使用JSON格式如{“type”: “heartbeat”, “data”: “ping”}。檢查防火墻/代理超時很多網絡設備如阿里云SLB會對空閑連接設置超時通常60秒。你的心跳間隔必須小于這個超時時間。建議設置為30秒。發送消息時出現IllegalStateException: The remote endpoint was in state [TEXT_FULL_WRITING]原因同一個Session上前一個異步發送操作還沒完成又觸發了新的發送導致狀態沖突。解決確保你的發送方法是線程安全的。我們之前使用的session.getAsyncRemote().sendText()是線程安全的但如果你在多個線程中同時調用同一個Session的發送方法仍有小概率出錯。更穩妥的做法是使用同步發送隊列。可以為每個Session維護一個消息隊列用一個單線程池依次發送。對于吞吐要求不極高的推送場景同步發送session.getBasicRemote().sendText()在簡單加鎖后反而更穩定。內存泄漏連接數只增不減檢查OnClose和OnError方法確保在所有連接關閉的路徑上正常關閉、異常關閉、心跳超時強制關閉都從SESSION_POOL和GROUP_MEMBERS等全局容器中移除了對應的Session和用戶信息。使用弱引用或定期清理如前所述考慮使用WeakReferenceSession或者定期掃描SESSION_POOL移除已經!session.isOpen()的死連接。這套從連接管理、安全校驗、心跳保活到分組廣播的WebSocket實現方案經過多個中等流量項目的驗證穩定性和擴展性都不錯。它最大的價值在于清晰地將各個關注點連接、安全、健康、路由分離代碼結構一目了然后續無論是加監控、改持久化方案還是接入消息隊列都有清晰的切入點可以操作。在實際部署時記得根據壓測結果調整線程池、心跳參數和JVM內存設置特別是SESSION_POOL的規模要做好預估。