的演進(jìn)與實踐)
前陣子有個朋友的公司做在線客服系統(tǒng)單機跑了一年多相安無事。結(jié)果有次渠道推廣爆量WebSocket 在線連接數(shù)直接沖上五位數(shù)服務(wù)開始頻繁卡頓、內(nèi)存飆高。他們第一反應(yīng)是加機器結(jié)果加了機器不但沒緩解反而冒出更詭異的 bug用戶明明在線消息卻經(jīng)常推送不到發(fā)一條消息別人要隔好幾秒才收到甚至干脆收不到。問題就出在 WebSocket 的本質(zhì)上——它是一條有狀態(tài)的 TCP 長連接。HTTP 請求處理完就斷開了隨便負(fù)載均衡到哪臺機器都行WebSocket 一旦握手成功這個連接就長在了某一臺節(jié)點上后續(xù)所有消息都只能由那臺節(jié)點轉(zhuǎn)發(fā)。你可以在 Redis 里存用戶的 session 數(shù)據(jù)但沒法把一條已經(jīng)建立的 TCP 連接瞬移到另一臺機器上。這就是所有 WebSocket 集群方案的出發(fā)點。這篇文章我會從單機服務(wù)的容量邊界講起再逐步拆解三種主流的集群方案Sticky Session 粘連、Redis Pub/Sub 消息路由、MQ 推送服務(wù)化架構(gòu)最后給出一套可以直接落地的代碼骨架和我在生產(chǎn)環(huán)境踩過的坑。不管你是做在線客服、消息推送、聊天室還是實時協(xié)作這套思路都能直接套用。1. WebSocket 的有狀態(tài)本質(zhì)所有集群麻煩的根源1.1 一條 TCP 長連接不是一份可以隨意路由的報文WebSocket 的握手過程和 HTTP 很像本質(zhì)是借助 HTTP Upgrade 機制完成協(xié)議升級GET /ws/chat HTTP/1.1 Host: im.example.com Upgrade: websocket Connection: Upgrade Sec-WebSocket-Key: x3JJHMbDL1EzLkh9GBhXDw Sec-WebSocket-Version: 13關(guān)鍵是這次握手發(fā)生在哪臺機器上后續(xù)一整條雙向通道就跟死在哪臺機器上。因為 WebSocket 是長連接服務(wù)端內(nèi)存里保存著這個連接的 session 對象、讀寫緩沖區(qū)、各種狀態(tài)標(biāo)記這些都不是 Redis 里存一個字符串就能搬走的東西。TCP socket 四元組綁定的是某一臺機器的某個端口數(shù)據(jù)報文只會被投遞到那個 socket 上其他節(jié)點根本沒有這個連接的任何信息。這就導(dǎo)致了一個很尷尬的現(xiàn)狀你可以在 Redis 里存用戶的登錄態(tài)沒法把連接本身同步給所有節(jié)點。負(fù)載均衡可以把 HTTP 請求均勻分發(fā)到每臺機器但對 WebSocket 來說連接一旦建立它的歸屬就固定了。1.2 HTTP 無狀態(tài)與 WebSocket 有狀態(tài)的對比HTTP 的每個請求都是獨立的服務(wù)端處理完就丟不存任何客戶端的狀態(tài)。所以 Nginx 后面掛 10 臺機器和掛 1 臺機器對業(yè)務(wù)代碼來說沒有區(qū)別隨便怎么輪詢都行。WebSocket 完全不同它在一次 TCP 連接上建立了全雙工通道而且這個通道是長久的。服務(wù)端必須要記住這個 userId 對應(yīng)哪個 session這個 session 連在哪個 socket 上否則收到消息不知道往哪里發(fā)。這就是狀態(tài)而狀態(tài)是分布式系統(tǒng)最棘手的敵人。對比維度HTTP 短連接WebSocket 長連接連接生命周期請求結(jié)束即斷開一直保持直到雙方關(guān)閉服務(wù)端是否保存連接狀態(tài)一般不保存必須保存 session 引用負(fù)載均衡策略隨便輪詢、隨機、加權(quán)連接一旦建立就不能遷移故障轉(zhuǎn)移請求重發(fā)即可連接斷開需要客戶端重連集群復(fù)雜度低天然橫向擴展高需要額外的路由機制很多人第一次做 WebSocket 集群時會下意識地按 HTTP 的思路去設(shè)計前邊掛 Nginx后邊掛一堆應(yīng)用節(jié)點結(jié)果一上線就發(fā)現(xiàn)消息發(fā)不出去然后才開始理解有狀態(tài)意味著什么。1.3 先定義業(yè)務(wù)場景再選后續(xù)方案做技術(shù)選型前我建議先想清楚業(yè)務(wù)場景因為不同的場景對集群方案的要求差別很大。我見過最典型的幾類服務(wù)端主動推送比如庫存變動通知、訂單狀態(tài)推送數(shù)據(jù)源在服務(wù)端客戶端被動接收。這類場景廣播和點對點都要用但對消息可靠性要求相對寬松。在線客服 / IM 聊天消息是雙向的用戶和客服之間的消息要準(zhǔn)確投遞不能丟順序還不能亂。這類場景對點對點路由要求很高通常要配套離線消息。聊天室 / 直播彈幕重點是廣播能力一個房間的消息要推給房間內(nèi)所有人而且量大、實時性要求高。這類場景要特別注意廣播風(fēng)暴問題。實時協(xié)同編輯比如白板、在線文檔消息頻率高、延遲敏感而且要求多端狀態(tài)一致。不同場景決定了你后面是用 Redis Pub/Sub 就夠了還是必須上 MQ甚至需要單獨做一個推送網(wǎng)關(guān)。這些決策在單機階段看不出來但等到集群階段就全是債。2. 單機服務(wù)先把容量邊界和連接管理做扎實2.1 一臺機器到底能扛多少連接很多人的第一個誤區(qū)是WebSocket 很重一臺機器扛不了多少連接。其實恰恰相反WebSocket 的協(xié)議開銷非常小真正吃資源的是每個連接占用的文件描述符、內(nèi)核 socket 緩沖區(qū)和應(yīng)用層 buffer。Linux 下 WebSocket 服務(wù)底層走的是 epoll 模型百萬并發(fā)連接在理論上是可以做到的但實際業(yè)務(wù)環(huán)境遠(yuǎn)達(dá)不到。我通常這樣估算每個空閑連接在應(yīng)用層大約占用 20KB50KB 內(nèi)存包括 session 對象、讀緩沖區(qū)、寫緩沖區(qū)。10 萬在線連接大約需要 2GB5GB 內(nèi)存這是純連接的消耗還不算業(yè)務(wù)對象。CPU 消耗主要來自心跳包的編解碼和消息的序列化空閑連接幾乎不占 CPU。所以一臺 8C16G 的機器跑 5 萬在線連接、每秒幾千條消息一般綽綽有余跑到 20 萬以上就要認(rèn)真調(diào)優(yōu)了。單機部署前一定要改幾個 Linux 內(nèi)核參數(shù)這是最容易被忽略的# 調(diào)整文件描述符上限 ulimit -n 1048576 # 內(nèi)核層面提升連接隊列長度 net.core.somaxconn 65535 net.ipv4.tcp_max_syn_backlog 65535 # 加大本地端口范圍防止大量短連接耗盡端口 net.ipv4.ip_local_port_range 1024 65535 # 加快 TIME_WAIT 回收 net.ipv4.tcp_fin_timeout 15不調(diào)文件描述符上限的話默認(rèn) 1024 的 ulimit 會直接卡死在連接數(shù)上業(yè)務(wù)代碼寫得再漂亮也沒用。2.2 連接管理的核心數(shù)據(jù)結(jié)構(gòu)單機模式下連接管理的核心就是在內(nèi)存里維護(hù)一張用戶 ID 到 WebSocketSession的映射表。Java 里用 ConcurrentHashMapGo 里用 sync.Map本質(zhì)都一樣Component public class SessionRegistry { // userId - WebSocketSession private final ConcurrentHashMapString, WebSocketSession sessions new ConcurrentHashMap(); public void register(String userId, WebSocketSession session) { sessions.put(userId, session); } public void unregister(String userId) { sessions.remove(userId); } public WebSocketSession get(String userId) { return sessions.get(userId); } public int count() { return sessions.size(); } }注意這里有個細(xì)節(jié)注冊的 key 是業(yè)務(wù)用戶 ID不是 session ID。因為你的消息投遞是面向用戶的而不是面向連接的。如果同一個用戶開多個標(biāo)簽頁、多臺設(shè)備那就需要建立 userId 到一組 session 的映射也就是 MapString, Set 廣播給這個用戶所有端。這個在 IM 場景下屬于剛需做單機的時候就要留好這個設(shè)計余地。2.3 心跳與死連接清理單機模式最容易踩的坑是連接泄漏。客戶端斷網(wǎng)、拔網(wǎng)線、電腦休眠TCP 層不一定能及時感知尤其是有中間 NAT 設(shè)備的情況下連接會一直掛在那里變成半開連接half-open。如果不做處理這些死連接會一直占著內(nèi)存和文件描述符最終把服務(wù)拖垮。解決方式就是應(yīng)用層心跳。我常用的方案是服務(wù)端每 30 秒下發(fā)一個 Ping 幀客戶端收到后回 Pong 幀服務(wù)端如果連續(xù) 3 次90 秒沒收到某個連接的 Pong就判定它已死亡主動關(guān)閉并清理 session。public void startHeartbeatCheck() { scheduledExecutor.scheduleAtFixedRate(() - { long now System.currentTimeMillis(); sessionRegistry.getAll().forEach((userId, session) - { long idleTime now - lastPongTime(userId); if (idleTime 90_000) { session.close(CloseStatus.SESSION_NOT_RELIABLE); sessionRegistry.unregister(userId); } else if (idleTime 30_000) { session.sendMessage(new PingMessage()); } }); }, 10, 10, TimeUnit.SECONDS); }這里間隔不是隨便定的。間隔太短比如 5 秒心跳包會占用大量帶寬和 CPU10 萬連接每秒光心跳就是幾萬條消息間隔太長比如 5 分鐘死連接清理不及時連接數(shù)虛高。30 秒心跳、90 秒判定死亡是我在多個項目里驗證過比較穩(wěn)的參數(shù)。2.4 單機模式最常見的坑除了心跳單機還會遇到幾個高頻問題Nginx 默認(rèn)超時如果前面掛了 Nginx 做反向代理默認(rèn) proxy_read_timeout 是 60 秒WebSocket 連接空閑超過 60 秒就會被 Nginx 掐斷。解決方式是顯式設(shè)置較大的超時時間proxy_read_timeout 3600s。在事件循環(huán)里做阻塞操作很多 WebSocket 框架是基于 Netty 或類似的事件循環(huán)模型如果你在消息處理器里直接調(diào)用遠(yuǎn)程 API、查數(shù)據(jù)庫會阻塞事件線程導(dǎo)致整個服務(wù)吞吐暴跌。正確做法是把耗時的操作丟到業(yè)務(wù)線程池或者用異步方式處理。單機依賴單點連接全在一臺機器上進(jìn)程崩潰、機器重啟所有連接瞬間全部斷開。這個無解只能靠集群解決。單機做扎實的意義在于它幫你把連接管理的細(xì)節(jié)注冊、心跳、清理、投遞都理清楚了這些邏輯在集群模式下會被復(fù)用而不會白白浪費。3. 集群第一板斧Sticky Session 粘滯與網(wǎng)關(guān)層配置3.1 ip_hash 與 cookie 粘滯的原理先說說最簡單的一種集群思路讓同一個用戶的連接總是被負(fù)載均衡到同一臺后端節(jié)點。Nginx 的 ip_hash 算法會根據(jù)客戶端 IP 做哈希同一個 IP 的請求會被分配到同一個 upstream 節(jié)點upstream ws_backend { ip_hash; server 192.168.1.10:8080; server 192.168.1.11:8080; server 192.168.1.12:8080; } server { listen 80; location /ws { proxy_pass http://ws_backend; # WebSocket 升級必需 proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; # 長連接超時放寬 proxy_read_timeout 3600s; proxy_send_timeout 3600s; } }這樣同一個客戶端 IP 的所有 WebSocket 連接都會落在同一臺節(jié)點上。對于用戶量不大、節(jié)點不多的場景這個方案能解決大部分問題而且改動量最小。還有一個變體是基于 cookie 的粘滯Nginx 會下發(fā)一個帶后端節(jié)點標(biāo)識的 cookie后續(xù)請求根據(jù) cookie 直接路由到指定節(jié)點。相比之下 cookie 方案比 ip_hash 更精確因為同一個 NAT 后面的多個用戶 IP 相同ip_hash 會把他們都砸到同一臺節(jié)點上而 cookie 方案能區(qū)分開。3.2 粘滯方案的優(yōu)勢與天花板Sticky Session 的優(yōu)勢非常明顯零業(yè)務(wù)改造不需要額外引入 Redis 或 MQ連接在哪個節(jié)點就是哪個節(jié)點點對點消息直接查本地 session 表就能發(fā)。但它的天花板也很低節(jié)點故障就是災(zāi)難如果一臺節(jié)點宕機粘滯在這臺機器上的所有連接全部斷開客戶端需要重新握手但此時 ip_hash 仍然會把它們路由到同一臺已經(jīng)宕機的節(jié)點直到 Nginx 把該節(jié)點摘除。即使摘除了其它節(jié)點上也沒有這些用戶的 session必須靠客戶端重新注冊。負(fù)載不均衡某個 IP 段用戶量大時ip_hash 可能把大量連接堆在同一臺節(jié)點上其他節(jié)點空閑。這就是加了機器反而沒效果的經(jīng)典原因之一。無法解決廣播問題要向所有用戶廣播消息時需要遍歷所有節(jié)點上的所有連接。粘滯方案里每個節(jié)點只知道自己本地的連接你仍然需要一套機制把廣播消息分發(fā)到每個節(jié)點。所以我的結(jié)論是Sticky Session 只適合作為最小可用集群方案一般在項目初期、用戶量不大、對可用性要求不高的場景使用。它不能算真正意義上的集群解決方案只能算負(fù)載均衡策略。3.3 網(wǎng)關(guān)層必須處理的細(xì)節(jié)不管用不用粘滯網(wǎng)關(guān)層Nginx的 WebSocket 配置都有幾個必須處理的細(xì)節(jié)很多人在這里踩坑location /ws { proxy_pass http://ws_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_connect_timeout 60s; proxy_read_timeout 3600s; proxy_send_timeout 3600s; proxy_buffer_size 64k; proxy_buffers 8 64k; }幾個關(guān)鍵點proxy_http_version 1.1必須設(shè)置WebSocket Upgrade 依賴 HTTP/1.1 的持久連接特性。proxy_set_header Upgrade和Connection upgrade是協(xié)議升級的關(guān)鍵不設(shè)置的話 Nginx 不會轉(zhuǎn)發(fā) Upgrade 頭WebSocket 握手直接失敗。proxy_read_timeout和proxy_send_timeout必須調(diào)大否則空閑連接會被 Nginx 掐掉。proxy_buffer_size關(guān)系到 WebSocket 幀的緩沖區(qū)大小如果消息體比較大比如超過 64KB需要同步調(diào)大 buffer否則會出現(xiàn)消息截斷或報錯。另外如果集群里走的是 HTTP/2要注意 Nginx 對 WebSocket over HTTP/2 的支持在較老版本里不完善生產(chǎn)環(huán)境建議 WebSocket 走獨立的 HTTP/1.1 監(jiān)聽端口和普通 HTTPS 業(yè)務(wù)分開。4. 集群第二板斧Redis Pub/Sub 做跨節(jié)點消息路由4.1 核心思路本地注冊表 全局路由表Sticky Session 解決了連接歸屬問題但沒有解決跨節(jié)點消息路由的問題。真正通用的做法是引入一層全局路由信息用 Redis 保存用戶 ID 落在哪個節(jié)點的映射關(guān)系節(jié)點間通過 Redis Pub/Sub 互相通信。這個方案的核心思路是每個節(jié)點在本地內(nèi)存維護(hù)自己的 session 注冊表只保存連到本節(jié)點的連接。每個節(jié)點啟動時生成一個全局唯一的 nodeId并把自己注冊到 Redis。用戶連接建立時把 userId - nodeId 的映射寫入 Redis并定期續(xù)期。節(jié)點間消息傳遞走 Redis Pub/Sub發(fā)送節(jié)點把消息發(fā)布到指定節(jié)點的頻道目標(biāo)節(jié)點的訂閱者收到后查本地 session 表并推送。這樣每個節(jié)點都不需要知道其他節(jié)點的完整連接信息只需要知道目標(biāo)用戶在哪臺節(jié)點剩下的投遞動作由目標(biāo)節(jié)點本地完成。Redis 的 key 設(shè)計我習(xí)慣這么搞ws:user:{userId} - nodeId # 全局路由表TTL 90 秒心跳續(xù)期 ws:nodes - SetnodeId # 存活節(jié)點列表 ws:msg:{nodeId} - 消息 # 點對點消息頻道發(fā)往指定節(jié)點 ws:broadcast - 消息 # 廣播消息頻道所有節(jié)點訂閱每個節(jié)點訂閱兩個頻道ws:msg:{自己的nodeId}和ws:broadcast。這樣點對點和廣播就用一套機制統(tǒng)一處理了。4.2 點對點消息的完整鏈路點對點消息的完整鏈路是這樣的業(yè)務(wù)服務(wù)想給用戶 U 推送一條消息。查詢 Redisws:user:U得到目標(biāo)節(jié)點 nodeId。如果 nodeId 就是本節(jié)點直接從本地 session 表查連接并發(fā)送。如果不是本節(jié)點把消息發(fā)布到 Redis 頻道ws:msg:{nodeId}。目標(biāo)節(jié)點的訂閱者收到消息查詢本地 session 表找到連接后發(fā)送。用 Java 偽代碼表示大概是這個樣子public void sendToUser(String userId, String payload) { String nodeId stringRedisTemplate.opsForValue().get(ws:user: userId); if (nodeId null) { // 用戶不在線走離線消息邏輯 handleOfflineMessage(userId, payload); return; } if (nodeId.equals(localNodeId)) { // 本節(jié)點直接投遞 WebSocketSession session sessionRegistry.get(userId); if (session ! null session.isOpen()) { session.sendMessage(new TextMessage(payload)); } } else { // 跨節(jié)點路由發(fā)布到目標(biāo)節(jié)點頻道 WsRouteMessage routeMsg new WsRouteMessage(userId, payload); stringRedisTemplate.convertAndSend(ws:msg: nodeId, routeMsg.toJson()); } }訂閱端統(tǒng)一處理Component public class WsMessageSubscriber extends AbstractMessageListener { Override public void onMessage(Message message, byte[] pattern) { WsRouteMessage msg JSON.parseObject(message.getBody(), WsRouteMessage.class); if (msg.isBroadcast()) { // 廣播消息遍歷本地所有連接發(fā)送 sessionRegistry.getAll().forEach((userId, session) - { if (session.isOpen()) { session.sendMessage(new TextMessage(msg.getPayload())); } }); return; } // 點對點消息查本地 session WebSocketSession session sessionRegistry.get(msg.getTargetUserId()); if (session ! null session.isOpen()) { session.sendMessage(new TextMessage(msg.getPayload())); } } }這套機制的核心好處是每個節(jié)點只保存自己的連接路由信息收斂到 Redis 里節(jié)點可以隨時水平擴展新節(jié)點上線只需要訂閱自己的頻道即可。4.3 廣播消息的完整鏈路廣播消息的處理比點對點簡單發(fā)送方直接往ws:broadcast頻道發(fā)布消息所有節(jié)點訂閱后各自往本地連接推送。但這里有一個隱蔽的問題如果廣播的接收者是某個聊天室的所有人而不是所有在線用戶那么每個節(jié)點在收到廣播后還需要判斷哪些本地連接屬于這個聊天室。這就需要在本地 session 表之外再維護(hù)一個聊天室 - 成員連接的映射表。// 房間 - userIds private final ConcurrentHashMapString, SetString roomMembers new ConcurrentHashMap();連接建立時根據(jù)客戶端帶上來的參數(shù)比如 URL query 里的 roomId把 userId 加入對應(yīng)房間連接銷毀時從所有房間移出。廣播給房間時節(jié)點拿到房間成員列表再逐個查 session 表發(fā)送。這樣廣播的范圍就被限制在目標(biāo)房間內(nèi)而不是全量廣播。4.4 Redis Pub/Sub 的邊界消息丟失與不可回溯Redis Pub/Sub 有個非常重要的特性消息不持久化。發(fā)布者把消息發(fā)出去如果此時某個訂閱者恰好不在線節(jié)點宕機、網(wǎng)絡(luò)抖動這條消息就永久丟失了。Redis 不會像 MQ 那樣幫你把消息存起來等消費者恢復(fù)后再投遞。所以 Redis Pub/Sub 方案只適用于消息實時投遞、丟了也無所謂或可以從業(yè)務(wù)側(cè)補償?shù)膱鼍?。比如在線狀態(tài)推送、心跳類通知、彈幕這類實時性消息。如果消息不能丟——比如聊天記錄、訂單通知——就必須在業(yè)務(wù)層做持久化和補償或者直接換用 MQ 方案。另外Redis Pub/Sub 的廣播是 push 模型如果某個節(jié)點處理消息過慢會導(dǎo)致該節(jié)點的訂閱者積壓甚至斷連。如果消息量非常大建議結(jié)合以下做法把廣播消息按業(yè)務(wù)維度切分到多個 channel比如ws:broadcast:room:{roomId}只有需要接收的節(jié)點才訂閱減少無效消息傳輸。在節(jié)點本地用隊列緩沖收到的消息再由獨立線程池發(fā)送避免阻塞 Redis 訂閱線程。這個是生產(chǎn)環(huán)境很重要的優(yōu)化點后面踩坑部分會細(xì)說。5. 集群第三板斧消息隊列與推送服務(wù)化架構(gòu)5.1 什么時候必須上 MQRedis Pub/Sub 有消息丟失的硬傷對于 IM、客服系統(tǒng)、交易通知這類要求不丟消息的場景就需要引入消息隊列RabbitMQ、Kafka、RocketMQ 等。我判斷是否需要上 MQ 的幾條標(biāo)準(zhǔn)消息不能丟比如用戶聊天記錄、支付結(jié)果通知丟失會引發(fā)資損或客訴。需要削峰填谷比如秒殺場景服務(wù)端瞬間產(chǎn)生大量推送消息直接打到 WebSocket 連接上會把節(jié)點打掛MQ 可以做流量緩沖。需要離線消息用戶不在線時消息要持久化等用戶上線后再補推。需要消息有序性IM 場景里同一個聊天窗口的消息必須有序Redis Pub/Sub 做不到精細(xì)化的順序保證而 MQ 可以按 key 分區(qū)保證局部有序。5.2 事件驅(qū)動設(shè)計引入 MQ 之后架構(gòu)從節(jié)點間互相路由演進(jìn)成事件驅(qū)動 獨立的推送層。整個鏈路變成業(yè)務(wù)服務(wù) --生產(chǎn)-- MQ Exchange --路由-- Queue --消費-- WebSocket 節(jié)點 --推送-- 客戶端業(yè)務(wù)服務(wù)不再直接關(guān)心目標(biāo)用戶在哪個節(jié)點它只需要把消息投遞到 MQ 對應(yīng)的隊列。WebSocket 節(jié)點作為消費者監(jiān)聽隊列拿到消息后查本地 session 表并推送。這個設(shè)計最大的好處是業(yè)務(wù)邏輯和連接管理徹底解耦。訂單服務(wù)不需要關(guān)心用戶當(dāng)前連在哪臺機器上它只管發(fā)消息WebSocket 節(jié)點只管消費和推送。節(jié)點可以隨時擴縮容對業(yè)務(wù)方完全透明。從 Redis Pub/Sub 遷移到 MQ 時之前那套路由表依然有用——節(jié)點消費到消息后仍然需要查ws:user:{userId}判斷目標(biāo)用戶是否在本節(jié)點。不過這里有個優(yōu)化點可以讓 MQ 按目標(biāo)節(jié)點做分區(qū)比如把消息路由到指定節(jié)點的專用隊列這樣每個節(jié)點只消費自己需要處理的消息減少無效消費。5.3 離線消息與消息補償有了 MQ離線消息就好處理了。用戶不在線時消息先落庫或者存儲在 Redis 里用戶重新建立 WebSocket 連接后服務(wù)端從存儲中拉取該用戶的離線消息補推。這里有一個我趟過的坑離線消息的補推不能一股腦全推。用戶斷線十分鐘可能積壓幾百條消息一次性推過去不僅客戶端渲染卡頓還會觸發(fā)大量 ACK 回執(zhí)反而把剛剛恢復(fù)的連接打掛。正確做法是分批補推比如每次推 20 條等客戶端確認(rèn)后再推下一批。另外消息補償機制要考慮冪等。客戶端收到消息后可能會回執(zhí)服務(wù)端重推時要有去重邏輯否則用戶會看到重復(fù)消息。通常用消息 ID 做冪等鍵客戶端按 ID 去重服務(wù)端按 ID 記錄已推送游標(biāo)。MQ 方案也有新的問題要處理消費者宕機恢復(fù)后從哪個 offset 開始消費超時未 ACK 的消息是否會重復(fù)投遞這些屬于 MQ 使用的基礎(chǔ)問題這里不展開但一定要在設(shè)計方案時提前想好。6. 實戰(zhàn)拆解一套可落地的 WebSocket 集群骨架6.1 整體架構(gòu)布局把前面的方案結(jié)合起來一套比較完整的 WebSocket 集群架構(gòu)長這樣接入層Nginx / SLB 做負(fù)載均衡不配置粘滯WebSocket 握手隨機分發(fā)到任意節(jié)點。因為引入路由層后連接落在哪臺節(jié)點已經(jīng)無所謂了。連接層N 個 WebSocket 應(yīng)用節(jié)點每個節(jié)點維護(hù)自己的本地 session 表并注冊到 Redis。路由層Redis 保存 userId - nodeId 映射節(jié)點間通過 Pub/Sub 通信。消息層RabbitMQ 承擔(dān)可靠消息投遞業(yè)務(wù)系統(tǒng)通過 MQ 解耦。存儲層MySQL 存聊天記錄等需要持久化的數(shù)據(jù)Redis 同時承擔(dān)路由表和熱數(shù)據(jù)緩存。這個架構(gòu)的好處是每一層都可以獨立擴展連接多了加 WebSocket 節(jié)點消息量大了擴 MQ 分區(qū)路由表性能不夠就升級 Redis 集群。6.2 連接注冊與注銷流程連接建立時的完整流程客戶端發(fā)起 WebSocket 握手Nginx 轉(zhuǎn)發(fā)到任意一個應(yīng)用節(jié)點。節(jié)點在 onOpen 回調(diào)里拿到 userId從 token 或 URL 參數(shù)解析。把 WebSocketSession 注冊到本地 session 表。把ws:user:{userId}寫入 Redis值為當(dāng)前節(jié)點 nodeIdTTL 90 秒。開啟該連接的心跳監(jiān)控。連接斷開時的清理流程客戶端主動關(guān)閉或心跳超時判定死亡后觸發(fā) onClose 回調(diào)。從本地 session 表移除該 userId。刪除 Redis 里的ws:user:{userId}先比對 nodeId 是否為本節(jié)點防止誤刪。從所有房間成員表里移除該 userId。把離線消息標(biāo)記為待補推狀態(tài)。這里有個細(xì)節(jié)刪除 Redis 路由表時要帶上 nodeId 做條件刪除。因為可能用戶剛斷線又在新節(jié)點上建立了新連接并寫入了新的路由表此時舊節(jié)點如果直接 del會把新路由信息也刪掉導(dǎo)致消息路由失敗。用 Lua 腳本或者 compare-and-delete 都能解決。6.3 節(jié)點上下線與故障轉(zhuǎn)移節(jié)點的健康檢查一般在 Redis 里做每個節(jié)點啟動時把自己的 nodeId 寫進(jìn)一個有序集合并周期性地更新心跳時間戳// 節(jié)點心跳上報每 30 秒一次 stringRedisTemplate.opsForZSet().add( ws:nodes, localNodeId, System.currentTimeMillis() );其他節(jié)點或監(jiān)控服務(wù)定期掃描這個有序集合把心跳時間超過 90 秒的 nodeId 判定為宕機節(jié)點從集合里移除并幫他做善后工作廣播節(jié)點 xxx 已下線的通知各節(jié)點清理本地可能存在的該節(jié)點的關(guān)聯(lián)信息。該節(jié)點持有的所有用戶連接被動斷開客戶端通過重連機制落到其他節(jié)點。如有必要從存儲層拉取這些用戶的會話狀態(tài)重新路由。節(jié)點故障時客戶端重連是不可避免的但一定要做重連保護(hù)否則會觸發(fā)重連風(fēng)暴。我常用的策略是指數(shù)退避 隨機抖動??蛻舳说谝淮沃剡B等 1 秒之后 2 秒、4 秒、8 秒……最大 30 秒封頂每次重連時間加一個 030% 的隨機抖動避免所有客戶端同時重連壓垮網(wǎng)關(guān)。6.4 關(guān)鍵代碼骨架最后給一個完整的 WebSocket 連接注冊 跨節(jié)點路由的最小骨架基于 Spring Boot 和 RedisServerEndpoint(/ws/{userId}) Component public class WsEndpoint { OnOpen public void onOpen(Session session, PathParam(userId) String userId) { // 1. 本地注冊 SessionRegistry.register(userId, session); // 2. 全局路由表注冊TTL 90s心跳續(xù)期 RedisUtil.set(ws:user: userId, LocalNode.getNodeId(), 90); // 3. 訂閱當(dāng)前節(jié)點的消息頻道只訂閱一次 RedisSubscriber.subscribe(ws:msg: LocalNode.getNodeId()); // 4. 啟動心跳任務(wù) HeartbeatManager.start(session, userId); } OnClose public void onClose(PathParam(userId) String userId) { SessionRegistry.unregister(userId); RedisUtil.compareAndDelete(ws:user: userId, LocalNode.getNodeId()); RoomManager.removeFromAllRooms(userId); } OnMessage public void onMessage(String message, PathParam(userId) String userId) { // 業(yè)務(wù)消息處理投遞到 MQ 或直接路由 ChatService.handleUserMessage(userId, message); } OnError public void onError(Session session, Throwable error) { // 記錄日志連接由心跳機制兜底清理 log.error(ws error, error); } }消息路由服務(wù)Service public class MessageRouter { public void sendToUser(String userId, String payload) { String targetNode RedisUtil.get(ws:user: userId); if (targetNode null) { offlineMessageStore.save(userId, payload); return; } if (targetNode.equals(LocalNode.getNodeId())) { SessionRegistry.sendToUser(userId, payload); } else { RedisUtil.publish(ws:msg: targetNode, payload); } } public void broadcastToRoom(String roomId, String payload) { RedisUtil.publish(ws:broadcast:room: roomId, payload); } }這套骨架可以直接跑通單機和集群兩種模式單機運行時所有連接都注冊到唯一的節(jié)點上sendToUser 直接走本地發(fā)送集群運行時靠 Redis 路由表自動切換到跨節(jié)點發(fā)布。業(yè)務(wù)層幾乎不需要改動。7. 生產(chǎn)環(huán)境踩坑記心跳、連接泄漏與廣播風(fēng)暴7.1 心跳間隔不當(dāng)引發(fā)的雪崩重連有次我在壓測環(huán)境發(fā)現(xiàn)一個詭異現(xiàn)象在線連接數(shù)穩(wěn)定在 5 萬左右但每分鐘都有大量連接斷開重連服務(wù)端日志里全是session closed和new connection。排查了很久才發(fā)現(xiàn)問題出在心跳參數(shù)上。當(dāng)時的配置是 10 秒發(fā)一次 Ping、30 秒判定死亡但網(wǎng)關(guān)層 Nginx 的超時設(shè)置只有 20 秒。Nginx 在超過 20 秒沒有收到任何數(shù)據(jù)時會主動關(guān)閉連接而服務(wù)端的心跳判定周期是 30 秒導(dǎo)致 Nginx 先于服務(wù)端把空閑連接關(guān)掉了??蛻舳税l(fā)現(xiàn)連接斷開后立刻重連重連風(fēng)暴把負(fù)載打得很高。這個問題的根源是網(wǎng)關(guān)超時和心跳周期不匹配。后來我把整套鏈路的心跳節(jié)奏統(tǒng)一了應(yīng)用層 30 秒發(fā)心跳Nginx 超時 3600 秒服務(wù)端 90 秒判定死亡。至此再沒出現(xiàn)過莫名重連。教訓(xùn)很簡單心跳不是一個應(yīng)用層參數(shù)而是一條完整鏈路上的協(xié)作參數(shù)??蛻舳恕⒕W(wǎng)關(guān)、應(yīng)用節(jié)點、Redis TTL 四者的超時時間必須按防火墻 Nginx 應(yīng)用判定死亡 心跳間隔的層級關(guān)系設(shè)置好任何一環(huán)不匹配都會出怪問題。7.2 連接泄漏只會讀不會關(guān)的客戶端怎么處理還有一次線上問題讓我印象很深某個老版本的 App 客戶端在弱網(wǎng)環(huán)境下不會正常發(fā)送關(guān)閉幀也不會響應(yīng) Pong。TCP 連接半死不活服務(wù)端檢測不到異常連接數(shù)一直漲最終內(nèi)存被打滿。應(yīng)用層心跳只能發(fā)現(xiàn)不響應(yīng) Pong的連接但如果你只是每 30 秒發(fā)一次 Ping且客戶端永遠(yuǎn)不會回 Pong那你最早也要等 90 秒才能清理掉它。如果 QPS 很高90 秒就能積累大量死連接。后來我做了兩個優(yōu)化在心跳檢測時除了發(fā) Ping還會檢查連接最近一次收到任何幀的時間。只要客戶端還在 TCP 層傳輸數(shù)據(jù)哪怕是垃圾幀就把它標(biāo)記為活躍超過 90 秒沒有任何數(shù)據(jù)到達(dá)的直接強殺。對客戶端主動斷開但 TCP 層沒有 FIN 的情況依賴內(nèi)核的 keepalive 來做最后兜底應(yīng)用層只負(fù)責(zé)更快的檢測。另外我要強調(diào)一點清理死連接時一定要在 finally 塊里執(zhí)行 unregister。否則連接關(guān)閉異常會導(dǎo)致 session 殘留在注冊表里造成幽靈連接這是連接泄漏最隱蔽的形態(tài)。7.3 廣播風(fēng)暴與消息重復(fù)最后一個坑來自廣播場景。做直播彈幕時一開始廣播消息直接走ws:broadcast全局頻道所有節(jié)點收到后向本地所有連接推送。等房間人數(shù)上到幾千問題就爆發(fā)了每個節(jié)點推送隊列積壓Redis 訂閱線程卡死消息延遲從毫秒級飆升到秒級。問題本質(zhì)是廣播范圍沒有做細(xì)粒度控制。全局廣播頻道會讓每個節(jié)點都收到所有消息但一個彈幕只屬于一個房間節(jié)點收到后還要去查這個房間里有哪些本地連接大部分查詢都是空轉(zhuǎn)。我把廣播頻道從全局切到了房間維度ws:broadcast:room:{roomId}節(jié)點按需訂閱用戶所在房間的頻道。在此基礎(chǔ)上發(fā)送線程也從 Redis 訂閱線程里拆了出來訂閱線程收到消息后只負(fù)責(zé)放進(jìn)本地隊列由獨立的發(fā)送線程池消費并推送給客戶端。這樣即使某個房間的消息量特別大影響的也只是這個房間的發(fā)送線程不會拖垮整個節(jié)點。消息重復(fù)的問題也值得一提。Redis Pub/Sub 本身不會重復(fù)投遞但引入 MQ 后消費者如果在發(fā)送推送后、ACK 之前宕機重啟后會重新消費這條消息導(dǎo)致用戶收到重復(fù)推送。解決方式是給每條消息生成唯一 ID節(jié)點在本地維護(hù)一個最近處理過的消息 ID 緩存收到重復(fù) ID 直接丟棄。緩存窗口不用太長30 秒即可因為 MQ 的重復(fù)消費通常發(fā)生在很短時間內(nèi)。做 WebSocket 集群這幾年我最大的體會是不要一上來就追求最復(fù)雜的架構(gòu)。單機能把連接管理、心跳、推送這些基本功做扎實是比上來就上分布式更重要的能力。集群方案的選擇也遵循這個原則——用戶量上來了先做 Sticky Session 頂一陣遇到跨節(jié)點路由需求了再加 Redis Pub/Sub消息可靠性要求高了再引入 MQ 和服務(wù)化改造。每一層方案都解決上一層的痛點但也會引入新的復(fù)雜度只有在真正需要的時候才值得接住這份復(fù)雜度。