Spring Boot 2.7 + JDK 8 實現(xiàn) WebSocket 集群分布式部署方案(基于 Redis Pub/Sub 方案)
在 Spring Boot 2.7 + JDK 8 環(huán)境下,WebSocketSession 無法直接序列化存儲到 Redis(它是與服務器節(jié)點綁定的TCP連接對象,跨JVM/跨節(jié)點無法復用)。
核心解決方案
行業(yè)標準集群方案:本地內存管理會話 + Redis 發(fā)布/訂閱(Pub/Sub)廣播消息
- 每個節(jié)點保留本地會話存儲(你原有的
ConcurrentHashMap完全保留,負責管理當前節(jié)點的連接); - 發(fā)送消息時,通過 Redis Pub/Sub 將消息廣播到所有集群節(jié)點;
- 所有節(jié)點監(jiān)聽Redis消息,收到后給本地的目標用戶推送消息。
該方案無需序列化WebSocketSession,完美支持集群分布式部署。
完整改造步驟
一、引入依賴
pom.xml 添加 Redis 依賴(Spring Data Redis 適配 Spring Boot 2.7):
<!-- Redis 核心依賴 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>
<!-- Redis 連接池 -->
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-pool2</artifactId>
</dependency>二、配置Redis連接
application.yml 配置Redis:
spring:
redis:
host: 127.0.0.1
port: 6379
password: # 有密碼填寫
database: 0
lettuce:
pool:
max-active: 8
max-idle: 8
min-idle: 0
max-wait: -1ms三、定義Redis常量
創(chuàng)建常量類,統(tǒng)一管理WebSocket消息頻道:
package com.example.demo.config.websocket;
/**
* WebSocket Redis 常量
*/
public interface WebSocketRedisConstants {
/**
* WebSocket 消息發(fā)布訂閱頻道
*/
String WEBSOCKET_MESSAGE_CHANNEL = "websocket:message:channel";
}四、改造會話管理類(核心)
將原靜態(tài)工具類改為 Spring Bean,注入RedisTemplate,保留本地會話管理,新增Redis消息廣播邏輯:
package com.example.demo.config.websocket;
import com.example.demo.config.LoginUser;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import org.springframework.web.socket.TextMessage;
import org.springframework.web.socket.WebSocketSession;
import java.util.*;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.CopyOnWriteArraySet;
/**
* 移動端WebSocket 在線用戶管理(支持Redis集群)
*/
@Slf4j
@Component // 改為Spring Bean,支持注入Redis
public class MobileWebSocketUserHolder {
/**
* 【本地存儲】在線用戶會話(保留原邏輯,僅管理當前節(jié)點連接)
*/
private static final Map<String, Set<WebSocketSession>> ONLINE_USER_MAP = new ConcurrentHashMap<>();
@Autowired
private StringRedisTemplate redisTemplate;
private final ObjectMapper objectMapper = new ObjectMapper();
// ====================== 原綁定/解綁邏輯 完全保留 ======================
/**
* 綁定用戶與WebSocket會話(連接成功時調用)
*/
public void bindSession(LoginUser user, WebSocketSession session) {
if (user == null || user.getUserId() == null) {
return;
}
String userId = user.getUserId();
ONLINE_USER_MAP.computeIfAbsent(userId, k -> new CopyOnWriteArraySet<>()).add(session);
log.info("用戶{}綁定WebSocket會話,當前節(jié)點在線用戶數(shù):{}", userId, ONLINE_USER_MAP.size());
}
/**
* 解綁用戶與WebSocket會話(連接關閉時調用)
*/
public void unbindSession(LoginUser user, WebSocketSession session) {
if (user == null || user.getUserId() == null) {
return;
}
String userId = user.getUserId();
Set<WebSocketSession> sessions = ONLINE_USER_MAP.get(userId);
if (sessions != null) {
sessions.remove(session);
if (sessions.isEmpty()) {
ONLINE_USER_MAP.remove(userId);
}
}
log.info("用戶{}解綁WebSocket會話,當前節(jié)點在線用戶數(shù):{}", userId, ONLINE_USER_MAP.size());
}
// ====================== 消息發(fā)送:本地發(fā)送 + Redis廣播 ======================
/**
* 給指定用戶發(fā)送消息(集群模式)
*/
public void sendMessageToUser(String userId, String message) {
// 1. 當前節(jié)點直接發(fā)送消息
sendLocalMessage(userId, message);
// 2. 發(fā)布消息到Redis,廣播給所有集群節(jié)點
try {
Map<String, String> msgMap = new HashMap<>(2);
msgMap.put("userId", userId);
msgMap.put("message", message);
String redisMsg = objectMapper.writeValueAsString(msgMap);
redisTemplate.convertAndSend(WebSocketRedisConstants.WEBSOCKET_MESSAGE_CHANNEL, redisMsg);
} catch (JsonProcessingException e) {
log.error("Redis消息序列化失敗", e);
}
}
/**
* 給所有用戶廣播消息(集群模式)
*/
public void sendMessageToAll(String message) {
// 1. 當前節(jié)點廣播
ONLINE_USER_MAP.keySet().forEach(userId -> sendLocalMessage(userId, message));
// 2. Redis廣播所有節(jié)點
try {
Map<String, String> msgMap = new HashMap<>(2);
msgMap.put("userId", "ALL");
msgMap.put("message", message);
String redisMsg = objectMapper.writeValueAsString(msgMap);
redisTemplate.convertAndSend(WebSocketRedisConstants.WEBSOCKET_MESSAGE_CHANNEL, redisMsg);
} catch (JsonProcessingException e) {
log.error("Redis廣播消息序列化失敗", e);
}
}
// ====================== 本地消息發(fā)送(私有方法) ======================
/**
* 僅給【當前節(jié)點】的用戶發(fā)送消息
*/
private void sendLocalMessage(String userId, String message) {
if (userId == null || !ONLINE_USER_MAP.containsKey(userId)) {
return;
}
Set<WebSocketSession> sessions = ONLINE_USER_MAP.get(userId);
for (WebSocketSession session : sessions) {
try {
if (session.isOpen()) {
session.sendMessage(new TextMessage(message));
}
} catch (Exception e) {
log.error("本地給用戶{}發(fā)送消息失敗", userId, e);
}
}
}
// ====================== 原查詢邏輯 完全保留 ======================
public Set<String> getOnlineUsers() {
return new HashSet<>(ONLINE_USER_MAP.keySet());
}
public boolean isOnline(String userId) {
return userId != null && ONLINE_USER_MAP.containsKey(userId);
}
}五、創(chuàng)建Redis消息監(jiān)聽器
監(jiān)聽Redis頻道,接收廣播消息并本地推送:
package com.example.demo.config.websocket;
import com.fasterxml.jackson.core.type.TypeReference;
import com.fasterxml.jackson.databind.ObjectMapper;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.connection.Message;
import org.springframework.data.redis.connection.MessageListener;
import org.springframework.stereotype.Component;
import java.util.Map;
/**
* Redis WebSocket 消息監(jiān)聽器
*/
@Slf4j
@Component
public class WebSocketRedisListener implements MessageListener {
@Autowired
private MobileWebSocketUserHolder webSocketUserHolder;
private final ObjectMapper objectMapper = new ObjectMapper();
@Override
public void onMessage(Message message, byte[] pattern) {
try {
// 解析Redis消息
String msgBody = new String(message.getBody());
Map<String, String> msgMap = objectMapper.readValue(msgBody, new TypeReference<Map<String, String>>() {});
String userId = msgMap.get("userId");
String content = msgMap.get("message");
// 本地發(fā)送消息
if ("ALL".equals(userId)) {
webSocketUserHolder.getOnlineUsers().forEach(uid -> webSocketUserHolder.sendLocalMessage(uid, content));
} else {
webSocketUserHolder.sendLocalMessage(userId, content);
}
} catch (Exception e) {
log.error("處理Redis WebSocket消息失敗", e);
}
}
}六、配置Redis發(fā)布/訂閱
注冊Redis監(jiān)聽容器,綁定頻道和監(jiān)聽器:
package com.example.demo.config;
import com.example.demo.config.websocket.WebSocketRedisConstants;
import com.example.demo.config.websocket.WebSocketRedisListener;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.redis.connection.RedisConnectionFactory;
import org.springframework.data.redis.listener.PatternTopic;
import org.springframework.data.redis.listener.RedisMessageListenerContainer;
import org.springframework.data.redis.listener.adapter.MessageListenerAdapter;
@Configuration
public class RedisConfig {
/**
* 注冊Redis消息監(jiān)聽器
*/
@Bean
public MessageListenerAdapter webSocketListenerAdapter(WebSocketRedisListener listener) {
return new MessageListenerAdapter(listener);
}
/**
* 配置Redis監(jiān)聽容器
*/
@Bean
public RedisMessageListenerContainer redisMessageListenerContainer(
RedisConnectionFactory connectionFactory,
MessageListenerAdapter webSocketListenerAdapter) {
RedisMessageListenerContainer container = new RedisMessageListenerContainer();
container.setConnectionFactory(connectionFactory);
// 綁定監(jiān)聽頻道
container.addMessageListener(webSocketListenerAdapter,
new PatternTopic(WebSocketRedisConstants.WEBSOCKET_MESSAGE_CHANNEL));
return container;
}
}七、修改WebSocket處理器(調用處適配)
由于MobileWebSocketUserHolder改為了Spring Bean,你需要在WebSocket處理器中注入使用,而非直接靜態(tài)調用:
// 原代碼(靜態(tài)調用) MobileWebSocketUserHolder.bindSession(user, session); // 改造后(注入調用) @Autowired private MobileWebSocketUserHolder webSocketUserHolder; webSocketUserHolder.bindSession(user, session);
方案原理說明
- 本地會話管理:每個節(jié)點只管理自己的
WebSocketSession(內存存儲,性能最高),不跨節(jié)點共享; - Redis廣播:任意節(jié)點調用發(fā)送消息接口時,會先給本地用戶發(fā)消息,再通過Redis Pub/Sub把消息發(fā)給所有集群節(jié)點;
- 集群推送:所有節(jié)點監(jiān)聽Redis消息,收到后給本地的目標用戶推送消息,實現(xiàn)全集群消息觸達。
集群部署注意事項
- 用戶認證一致性:集群所有節(jié)點的登錄認證邏輯必須一致(用戶ID生成規(guī)則相同);
- Redis 高可用:生產環(huán)境使用Redis集群/哨兵模式,避免單點故障;
- Session 共享非必須:本方案不需要共享WebSocketSession,這是最輕量化、最高效的集群方案;
- 心跳/重連:保留原有的WebSocket心跳機制,客戶端斷開后自動重連到任意集群節(jié)點即可。
總結
- 核心方案:本地內存管理會話 + Redis Pub/Sub 廣播消息(Spring Boot WebSocket集群標準方案);
- 無需序列化:規(guī)避了
WebSocketSession無法存儲Redis的問題; - 兼容原有邏輯:90%代碼復用,僅改造消息發(fā)送邏輯;
- 生產可用:支持多節(jié)點集群部署,無狀態(tài)、高可用。
到此這篇關于Spring Boot 2.7 + JDK 8 實現(xiàn) WebSocket 集群分布式部署方案(基于 Redis Pub/Sub 方案)的文章就介紹到這了,更多相關Spring Boot JDK WebSocket 集群內容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關文章希望大家以后多多支持腳本之家!
- SpringBoot集成WebSocket的兩種方式(JDK內置版和Spring封裝版)
- Spring?Boot?項目與JDK、Mybatis版本兼容對應關系表及問題記錄
- SpringBoot快速接入OpenAI大模型的方法(JDK8)
- 關于JDK8升級17及springboot?2.x升級3.x詳細指南
- Java搭建一個springboot3.4.1項目?JDK21的詳細過程
- SpringBoot配置開發(fā)環(huán)境的詳細步驟(JDK、Maven、IDEA等)
- IDEA無法創(chuàng)建JDK1.8版本的Springboot項目問題解決(2種方法)
- 使用Spring Initializr創(chuàng)建Spring Boot項目沒有JDK1.8的解決辦法
相關文章
Mybatis resultType返回結果為null的問題排查方式
這篇文章主要介紹了Mybatis resultType返回結果為null的問題排查方式,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2022-03-03
SpringBoot如何通過devtools實現(xiàn)熱部署
這篇文章主要介紹了SpringBoot如何通過devtools實現(xiàn)熱部署,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下2019-11-11
Springmvc如何實現(xiàn)向前臺傳遞數(shù)據(jù)
這篇文章主要介紹了Springmvc如何實現(xiàn)向前臺傳遞數(shù)據(jù),文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友可以參考下2020-07-07
MyBatis-Plus詳解(環(huán)境搭建、關聯(lián)操作)
MyBatis-Plus 是一個 MyBatis 的增強工具,在 MyBatis 的基礎上只做增強不做改變,為簡化開發(fā)、提高效率而生,今天通過本文給大家介紹MyBatis-Plus環(huán)境搭建及關聯(lián)操作,需要的朋友參考下吧2022-09-09
反射機制:getDeclaredField和getField的區(qū)別說明
這篇文章主要介紹了反射機制:getDeclaredField和getField的區(qū)別說明,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教2021-06-06

