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