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