Spring?Boot?中實(shí)現(xiàn)?WebSocket?集群思路詳解
在 Spring Boot 中實(shí)現(xiàn) WebSocket 集群處理,核心挑戰(zhàn)是解決多節(jié)點(diǎn)間的 WebSocket 會(huì)話共享和消息同步問(wèn)題(因?yàn)閱喂?jié)點(diǎn) WebSocket 的會(huì)話是內(nèi)存存儲(chǔ)的,多節(jié)點(diǎn)間無(wú)法直接通信)。以下是完整的實(shí)現(xiàn)方案:
一、核心思路
- 會(huì)話共享:使用分布式存儲(chǔ)(如 Redis)存儲(chǔ)所有節(jié)點(diǎn)的 WebSocket 會(huì)話信息(會(huì)話 ID、用戶標(biāo)識(shí)、節(jié)點(diǎn)標(biāo)識(shí)等)。
- 消息轉(zhuǎn)發(fā):當(dāng)某節(jié)點(diǎn)需要向用戶發(fā)送消息時(shí),先通過(guò) Redis 判斷用戶連接在哪個(gè)節(jié)點(diǎn),再通過(guò)消息隊(duì)列(如 Redis Pub/Sub)將消息轉(zhuǎn)發(fā)到目標(biāo)節(jié)點(diǎn),由目標(biāo)節(jié)點(diǎn)推送消息給用戶。
- 集群感知:每個(gè)節(jié)點(diǎn)啟動(dòng)時(shí)注冊(cè)自身信息,下線時(shí)注銷(xiāo),確保消息能正確路由。
二、技術(shù)選型
- WebSocket 框架:Spring WebSocket(基于 JSR-356 標(biāo)準(zhǔn))
- 分布式存儲(chǔ):Redis(存儲(chǔ)會(huì)話映射、節(jié)點(diǎn)信息)
- 消息隊(duì)列:Redis Pub/Sub(節(jié)點(diǎn)間消息轉(zhuǎn)發(fā))
- 依賴管理:Spring Boot Starter WebSocket + Spring Data Redis
三、實(shí)現(xiàn)步驟
1. 引入依賴
在pom.xml中添加以下依賴:
<!-- Spring WebSocket -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
<!-- Redis (用于會(huì)話存儲(chǔ)和消息轉(zhuǎn)發(fā)) -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- 連接池 (可選,優(yōu)化Redis性能) -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
</dependency>2. 配置 Redis
在application.yml中配置 Redis 連接:
spring:
redis:
host: localhost # Redis服務(wù)器地址
port: 6379 # Redis端口
password: # 密碼(如有)
lettuce:
pool:
max-active: 8 # 最大連接數(shù)
max-idle: 8 # 最大空閑連接
min-idle: 2 # 最小空閑連接3. WebSocket 核心配置
實(shí)現(xiàn) WebSocket 配置類(lèi),注冊(cè)端點(diǎn)和消息處理器:
import org.springframework.context.annotation.Bean;
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;
import org.springframework.web.socket.server.HandshakeInterceptor;
@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
private final WebSocketHandler webSocketHandler;
private final HandshakeInterceptor handshakeInterceptor;
// 注入自定義處理器和攔截器
public WebSocketConfig(WebSocketHandler webSocketHandler, HandshakeInterceptor handshakeInterceptor) {
this.webSocketHandler = webSocketHandler;
this.handshakeInterceptor = handshakeInterceptor;
}
@Override
public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
// 注冊(cè)WebSocket端點(diǎn),允許跨域
registry.addHandler(webSocketHandler, "/ws")
.addInterceptors(handshakeInterceptor)
.setAllowedOrigins("*"); // 生產(chǎn)環(huán)境需指定具體域名
}
}4. 會(huì)話管理與集群同步
需要解決 3 個(gè)核心問(wèn)題:會(huì)話注冊(cè)、消息路由、節(jié)點(diǎn)間消息轉(zhuǎn)發(fā)。
4.1 定義常量(Redis Key)
public class RedisKeyConstants {
// 存儲(chǔ)用戶-節(jié)點(diǎn)映射:key=userId,value=nodeId
public static final String USER_NODE_MAP = "websocket:user:node";
// 存儲(chǔ)節(jié)點(diǎn)信息:key=nodeId,value=節(jié)點(diǎn)信息(如IP:端口)
public static final String NODE_INFO = "websocket:node:info";
// Redis Pub/Sub頻道(用于節(jié)點(diǎn)間消息轉(zhuǎn)發(fā))
public static final String MSG_CHANNEL = "websocket:msg:channel";
}4.2 握手?jǐn)r截器(記錄會(huì)話與用戶映射)
在連接建立時(shí),將用戶 ID 與當(dāng)前節(jié)點(diǎn) ID 關(guān)聯(lián)并存儲(chǔ)到 Redis:
import org.springframework.http.server.ServerHttpRequest;
import org.springframework.http.server.ServerHttpResponse;
import org.springframework.http.server.ServletServerHttpRequest;
import org.springframework.web.socket.WebSocketHandler;
import org.springframework.web.socket.server.HandshakeInterceptor;
import org.springframework.data.redis.core.StringRedisTemplate;
import javax.servlet.http.HttpSession;
import java.util.Map;
import java.util.UUID;
@Component
public class WebSocketHandshakeInterceptor implements HandshakeInterceptor {
private final StringRedisTemplate redisTemplate;
private final String nodeId; // 當(dāng)前節(jié)點(diǎn)唯一標(biāo)識(shí)(如UUID)
public WebSocketHandshakeInterceptor(StringRedisTemplate redisTemplate) {
this.redisTemplate = redisTemplate;
this.nodeId = UUID.randomUUID().toString(); // 節(jié)點(diǎn)啟動(dòng)時(shí)生成唯一ID
}
@Override
public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response,
WebSocketHandler wsHandler, Map<String, Object> attributes) throws Exception {
// 從請(qǐng)求參數(shù)或Session中獲取用戶ID(根據(jù)業(yè)務(wù)場(chǎng)景調(diào)整)
ServletServerHttpRequest servletRequest = (ServletServerHttpRequest) request;
HttpSession session = servletRequest.getServletRequest().getSession();
String userId = (String) session.getAttribute("userId"); // 假設(shè)用戶已登錄,Session中存userId
if (userId != null) {
// 將用戶與當(dāng)前節(jié)點(diǎn)關(guān)聯(lián)(存儲(chǔ)到Redis)
redisTemplate.opsForValue().set(
RedisKeyConstants.USER_NODE_MAP + ":" + userId,
nodeId
);
// 注冊(cè)當(dāng)前節(jié)點(diǎn)信息(可選,用于監(jiān)控)
redisTemplate.opsForValue().set(
RedisKeyConstants.NODE_INFO + ":" + nodeId,
request.getLocalAddress().toString() // 節(jié)點(diǎn)地址
);
attributes.put("userId", userId); // 傳遞userId到處理器
return true;
}
return false; // 未登錄用戶拒絕連接
}
@Override
public void afterHandshake(ServerHttpRequest request, ServerHttpResponse response,
WebSocketHandler wsHandler, Exception exception) {}
}4.3 WebSocket 消息處理器
處理消息接收,并通過(guò) Redis Pub/Sub 轉(zhuǎn)發(fā)跨節(jié)點(diǎn)消息:
import org.springframework.data.redis.core.StringRedisTemplate;
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 com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Component
public class WebSocketHandler extends TextWebSocketHandler {
// 本地會(huì)話緩存(當(dāng)前節(jié)點(diǎn)的WebSocket會(huì)話)
private final Map<String, WebSocketSession> localSessions = new ConcurrentHashMap<>();
private final StringRedisTemplate redisTemplate;
private final ObjectMapper objectMapper;
private final String nodeId; // 當(dāng)前節(jié)點(diǎn)ID(與攔截器一致)
public WebSocketHandler(StringRedisTemplate redisTemplate, ObjectMapper objectMapper,
@Value("${websocket.node.id}") String nodeId) {
this.redisTemplate = redisTemplate;
this.objectMapper = objectMapper;
this.nodeId = nodeId;
// 訂閱Redis頻道,接收其他節(jié)點(diǎn)的消息
subscribeToRedisChannel();
}
// 連接建立時(shí),緩存本地會(huì)話
@Override
public void afterConnectionEstablished(WebSocketSession session) throws Exception {
String userId = (String) session.getAttributes().get("userId");
localSessions.put(userId, session);
}
// 接收客戶端消息
@Override
protected void handleTextMessage(WebSocketSession session, TextMessage message) throws Exception {
String userId = (String) session.getAttributes().get("userId");
String payload = message.getPayload(); // 客戶端發(fā)送的消息
// 解析消息(假設(shè)消息格式:{"targetUserId": "xxx", "content": "xxx"})
Map<String, String> msg = objectMapper.readValue(payload, Map.class);
String targetUserId = msg.get("targetUserId");
String content = msg.get("content");
// 發(fā)送消息給目標(biāo)用戶(跨節(jié)點(diǎn)則轉(zhuǎn)發(fā))
sendMessageToUser(targetUserId, content);
}
// 發(fā)送消息給指定用戶
public void sendMessageToUser(String targetUserId, String content) throws IOException {
// 1. 從Redis獲取目標(biāo)用戶所在節(jié)點(diǎn)
String targetNodeId = redisTemplate.opsForValue().get(
RedisKeyConstants.USER_NODE_MAP + ":" + targetUserId
);
if (targetNodeId == null) {
throw new RuntimeException("用戶未連接");
}
// 2. 若目標(biāo)節(jié)點(diǎn)是當(dāng)前節(jié)點(diǎn),直接發(fā)送;否則通過(guò)Redis轉(zhuǎn)發(fā)
if (targetNodeId.equals(nodeId)) {
WebSocketSession session = localSessions.get(targetUserId);
if (session != null && session.isOpen()) {
session.sendMessage(new TextMessage(content));
}
} else {
// 構(gòu)造跨節(jié)點(diǎn)消息(包含目標(biāo)用戶、內(nèi)容、目標(biāo)節(jié)點(diǎn))
Map<String, String> crossMsg = new HashMap<>();
crossMsg.put("targetUserId", targetUserId);
crossMsg.put("content", content);
crossMsg.put("targetNodeId", targetNodeId);
// 發(fā)布到Redis頻道
redisTemplate.convertAndSend(
RedisKeyConstants.MSG_CHANNEL,
objectMapper.writeValueAsString(crossMsg)
);
}
}
// 訂閱Redis頻道,接收其他節(jié)點(diǎn)的消息并轉(zhuǎn)發(fā)給本地用戶
private void subscribeToRedisChannel() {
redisTemplate.getConnectionFactory().getConnection().subscribe(
message -> {
try {
String payload = new String(message.getBody());
Map<String, String> crossMsg = objectMapper.readValue(payload, Map.class);
// 僅處理目標(biāo)節(jié)點(diǎn)為當(dāng)前節(jié)點(diǎn)的消息
if (crossMsg.get("targetNodeId").equals(nodeId)) {
String targetUserId = crossMsg.get("targetUserId");
String content = crossMsg.get("content");
WebSocketSession session = localSessions.get(targetUserId);
if (session != null && session.isOpen()) {
session.sendMessage(new TextMessage(content));
}
}
} catch (Exception e) {
e.printStackTrace();
}
},
RedisKeyConstants.MSG_CHANNEL.getBytes()
);
}
// 連接關(guān)閉時(shí),清理會(huì)話和Redis映射
@Override
public void afterConnectionClosed(WebSocketSession session, CloseStatus status) throws Exception {
String userId = (String) session.getAttributes().get("userId");
localSessions.remove(userId);
// 刪除Redis中的用戶-節(jié)點(diǎn)映射
redisTemplate.delete(RedisKeyConstants.USER_NODE_MAP + ":" + userId);
}
}4.4 節(jié)點(diǎn) ID 配置
在application.yml中配置節(jié)點(diǎn) ID(也可啟動(dòng)時(shí)自動(dòng)生成,確保唯一):
websocket:
node:
id: ${HOSTNAME:node-${random.uuid}} # 優(yōu)先使用主機(jī)名,否則隨機(jī)生成四、集群測(cè)試
- 啟動(dòng)多個(gè) Spring Boot 實(shí)例(模擬集群節(jié)點(diǎn)),確保都連接到同一個(gè) Redis。
- 客戶端 1 連接節(jié)點(diǎn) A,客戶端 2 連接節(jié)點(diǎn) B。
- 客戶端 1 發(fā)送消息給客戶端 2,消息流程:節(jié)點(diǎn) A 接收消息,查詢 Redis 發(fā)現(xiàn)客戶端 2 在節(jié)點(diǎn) B。節(jié)點(diǎn) A 通過(guò) Redis Pub/Sub 發(fā)布消息到MSG_CHANNEL。節(jié)點(diǎn) B 訂閱了該頻道,接收消息并推送給客戶端 2。
五、優(yōu)化建議
- 會(huì)話過(guò)期清理:給 Redis 中的USER_NODE_MAP設(shè)置過(guò)期時(shí)間(如 30 分鐘),并通過(guò) WebSocket 心跳機(jī)制續(xù)約。
- 消息可靠性:Redis Pub/Sub 是 fire-and-forget 模式,若需可靠消息,可替換為 RabbitMQ 或 Kafka。
- 負(fù)載均衡:前端通過(guò) Nginx 等負(fù)載均衡器連接 WebSocket 集群(需配置proxy_set_header Upgrade $http_upgrade;等參數(shù)支持 WebSocket)。
- 監(jiān)控:通過(guò) Redis 的NODE_INFO鍵監(jiān)控節(jié)點(diǎn)狀態(tài),實(shí)現(xiàn)故障轉(zhuǎn)移
到此這篇關(guān)于Spring Boot 中實(shí)現(xiàn) WebSocket 集群處理的文章就介紹到這了,更多相關(guān)Spring Boot WebSocket 集群內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
- springboot基于Redis發(fā)布訂閱集群下WebSocket的解決方案
- springboot websocket集群(stomp協(xié)議)連接時(shí)候傳遞參數(shù)
- Spring Boot集成WebSocket項(xiàng)目實(shí)戰(zhàn)的示例代碼
- SpringBoot分布式WebSocket的實(shí)現(xiàn)指南
- Java?springBoot初步使用websocket的代碼示例
- SpringBoot實(shí)現(xiàn)websocket服務(wù)端及客戶端的詳細(xì)過(guò)程
- vue+SpringBoot使用WebSocket方式
相關(guān)文章
基于Java編寫(xiě)一個(gè)html轉(zhuǎn)pdf的工具類(lèi)
這篇文章主要為大家詳細(xì)介紹了如何基于Java編寫(xiě)一個(gè)html轉(zhuǎn)pdf的工具類(lèi),文中的示例代碼講解詳細(xì),感興趣的小伙伴可以跟隨小編一起了解下2025-09-09
Java中DecimalFormat用法及符號(hào)含義
DecimalFormat是NumberFormat的一個(gè)具體子類(lèi),用于格式化十進(jìn)制數(shù)字。這篇文章介紹了DecimalFormat的用法及符號(hào)含義,需要的朋友可以收藏下,方便下次瀏覽觀看2021-12-12
Java中Map與對(duì)象之間互相轉(zhuǎn)換的幾種常用方式
在Java中將對(duì)象和Map相互轉(zhuǎn)換是常見(jiàn)的操作,可以通過(guò)不同的方式實(shí)現(xiàn)這種轉(zhuǎn)換,下面這篇文章主要給大家介紹了關(guān)于Java中Map與對(duì)象之間互相轉(zhuǎn)換的幾種常用方式,需要的朋友可以參考下2024-01-01
Spring?Boot基于?JWT?優(yōu)化?Spring?Security?無(wú)狀態(tài)登錄實(shí)戰(zhàn)指南
本文介紹如何使用JWT優(yōu)化SpringSecurity實(shí)現(xiàn)無(wú)狀態(tài)登錄,提高接口安全性,并通過(guò)實(shí)際操作步驟展示了如何配置JWT參數(shù)、實(shí)現(xiàn)JWT登錄接口、認(rèn)證過(guò)濾器等,感興趣的朋友跟隨小編一起看看吧2025-11-11
java并發(fā)中DelayQueue延遲隊(duì)列原理剖析
DelayQueue隊(duì)列是一個(gè)延遲隊(duì)列,本文將結(jié)合實(shí)例代碼,詳細(xì)的介紹DelayQueue延遲隊(duì)列的源碼分析,感興趣的小伙伴們可以參考一下2021-06-06
Java中MapStruct 映射過(guò)程中忽略某個(gè)字段的實(shí)現(xiàn)
在MapStruct中,如果你想要在映射過(guò)程中忽略某個(gè)字段,可以使用 @Mapping注解的ignore屬性,本文就來(lái)介紹一下Java中MapStruct映射過(guò)程中忽略某個(gè)字段的實(shí)現(xiàn),感興趣的可以了解一下2025-11-11
Java中的信號(hào)量Semaphore詳細(xì)解讀
這篇文章主要介紹了Java中的信號(hào)量Semaphore詳細(xì)解讀,Java信號(hào)量機(jī)制可以用來(lái)保證線程互斥,創(chuàng)建Semaphore對(duì)象傳入一個(gè)整形參數(shù),類(lèi)似于公共資源,需要的朋友可以參考下2023-11-11

