最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Spring Boot 2.7 + JDK 8 實(shí)現(xiàn) WebSocket 集群分布式部署方案(基于 Redis Pub/Sub 方案)

 更新時(shí)間:2026年03月30日 09:13:18   作者:weixin_42502300  
本文詳細(xì)介紹了在SpringBoot2.7+JDK8環(huán)境下,WebSocketSession無法直接序列化存儲(chǔ)到Redis的問題,并提供了行業(yè)標(biāo)準(zhǔn)的集群方案,本文介紹的非常詳細(xì),感興趣的朋友跟隨小編一起看看吧

Spring Boot 2.7 + JDK 8 環(huán)境下,WebSocketSession 無法直接序列化存儲(chǔ)到 Redis(它是與服務(wù)器節(jié)點(diǎn)綁定的TCP連接對(duì)象,跨JVM/跨節(jié)點(diǎn)無法復(fù)用)。

核心解決方案

行業(yè)標(biāo)準(zhǔn)集群方案:本地內(nèi)存管理會(huì)話 + Redis 發(fā)布/訂閱(Pub/Sub)廣播消息

  1. 每個(gè)節(jié)點(diǎn)保留本地會(huì)話存儲(chǔ)(你原有的ConcurrentHashMap完全保留,負(fù)責(zé)管理當(dāng)前節(jié)點(diǎn)的連接);
  2. 發(fā)送消息時(shí),通過 Redis Pub/Sub 將消息廣播到所有集群節(jié)點(diǎn);
  3. 所有節(jié)點(diǎn)監(jiān)聽Redis消息,收到后給本地的目標(biāo)用戶推送消息。

該方案無需序列化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";
}

四、改造會(huì)話管理類(核心)

將原靜態(tài)工具類改為 Spring Bean,注入RedisTemplate,保留本地會(huì)話管理,新增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;
/**
 * 移動(dòng)端WebSocket 在線用戶管理(支持Redis集群)
 */
@Slf4j
@Component // 改為Spring Bean,支持注入Redis
public class MobileWebSocketUserHolder {
    /**
     * 【本地存儲(chǔ)】在線用戶會(huì)話(保留原邏輯,僅管理當(dāng)前節(jié)點(diǎn)連接)
     */
    private static final Map<String, Set<WebSocketSession>> ONLINE_USER_MAP = new ConcurrentHashMap<>();
    @Autowired
    private StringRedisTemplate redisTemplate;
    private final ObjectMapper objectMapper = new ObjectMapper();
    // ====================== 原綁定/解綁邏輯 完全保留 ======================
    /**
     * 綁定用戶與WebSocket會(huì)話(連接成功時(shí)調(diào)用)
     */
    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會(huì)話,當(dāng)前節(jié)點(diǎn)在線用戶數(shù):{}", userId, ONLINE_USER_MAP.size());
    }
    /**
     * 解綁用戶與WebSocket會(huì)話(連接關(guān)閉時(shí)調(diào)用)
     */
    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會(huì)話,當(dāng)前節(jié)點(diǎn)在線用戶數(shù):{}", userId, ONLINE_USER_MAP.size());
    }
    // ====================== 消息發(fā)送:本地發(fā)送 + Redis廣播 ======================
    /**
     * 給指定用戶發(fā)送消息(集群模式)
     */
    public void sendMessageToUser(String userId, String message) {
        // 1. 當(dāng)前節(jié)點(diǎn)直接發(fā)送消息
        sendLocalMessage(userId, message);
        // 2. 發(fā)布消息到Redis,廣播給所有集群節(jié)點(diǎn)
        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. 當(dāng)前節(jié)點(diǎn)廣播
        ONLINE_USER_MAP.keySet().forEach(userId -> sendLocalMessage(userId, message));
        // 2. Redis廣播所有節(jié)點(diǎn)
        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ā)送(私有方法) ======================
    /**
     * 僅給【當(dāng)前節(jié)點(diǎn)】的用戶發(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ā)布/訂閱

注冊(cè)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 {
    /**
     * 注冊(cè)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處理器(調(diào)用處適配)

由于MobileWebSocketUserHolder改為了Spring Bean,你需要在WebSocket處理器中注入使用,而非直接靜態(tài)調(diào)用:

// 原代碼(靜態(tài)調(diào)用)
MobileWebSocketUserHolder.bindSession(user, session);
// 改造后(注入調(diào)用)
@Autowired
private MobileWebSocketUserHolder webSocketUserHolder;
webSocketUserHolder.bindSession(user, session);

方案原理說明

  1. 本地會(huì)話管理:每個(gè)節(jié)點(diǎn)只管理自己的WebSocketSession(內(nèi)存存儲(chǔ),性能最高),不跨節(jié)點(diǎn)共享;
  2. Redis廣播:任意節(jié)點(diǎn)調(diào)用發(fā)送消息接口時(shí),會(huì)先給本地用戶發(fā)消息,再通過Redis Pub/Sub把消息發(fā)給所有集群節(jié)點(diǎn);
  3. 集群推送:所有節(jié)點(diǎn)監(jiān)聽Redis消息,收到后給本地的目標(biāo)用戶推送消息,實(shí)現(xiàn)全集群消息觸達(dá)。

集群部署注意事項(xiàng)

  1. 用戶認(rèn)證一致性:集群所有節(jié)點(diǎn)的登錄認(rèn)證邏輯必須一致(用戶ID生成規(guī)則相同);
  2. Redis 高可用:生產(chǎn)環(huán)境使用Redis集群/哨兵模式,避免單點(diǎn)故障;
  3. Session 共享非必須:本方案不需要共享WebSocketSession,這是最輕量化、最高效的集群方案;
  4. 心跳/重連:保留原有的WebSocket心跳機(jī)制,客戶端斷開后自動(dòng)重連到任意集群節(jié)點(diǎn)即可。

總結(jié)

  1. 核心方案:本地內(nèi)存管理會(huì)話 + Redis Pub/Sub 廣播消息(Spring Boot WebSocket集群標(biāo)準(zhǔn)方案);
  2. 無需序列化:規(guī)避了WebSocketSession無法存儲(chǔ)Redis的問題;
  3. 兼容原有邏輯:90%代碼復(fù)用,僅改造消息發(fā)送邏輯;
  4. 生產(chǎn)可用:支持多節(jié)點(diǎn)集群部署,無狀態(tài)、高可用。

到此這篇關(guān)于Spring Boot 2.7 + JDK 8 實(shí)現(xiàn) WebSocket 集群分布式部署方案(基于 Redis Pub/Sub 方案)的文章就介紹到這了,更多相關(guān)Spring Boot JDK WebSocket 集群內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Mybatis resultType返回結(jié)果為null的問題排查方式

    Mybatis resultType返回結(jié)果為null的問題排查方式

    這篇文章主要介紹了Mybatis resultType返回結(jié)果為null的問題排查方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-03-03
  • SpringBoot如何通過devtools實(shí)現(xiàn)熱部署

    SpringBoot如何通過devtools實(shí)現(xiàn)熱部署

    這篇文章主要介紹了SpringBoot如何通過devtools實(shí)現(xiàn)熱部署,文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-11-11
  • spring與mybatis三種整合方法

    spring與mybatis三種整合方法

    這篇文章主要介紹了spring與mybatis三種整合方法,需要的朋友可以參考下
    2017-04-04
  • Springmvc如何實(shí)現(xiàn)向前臺(tái)傳遞數(shù)據(jù)

    Springmvc如何實(shí)現(xiàn)向前臺(tái)傳遞數(shù)據(jù)

    這篇文章主要介紹了Springmvc如何實(shí)現(xiàn)向前臺(tái)傳遞數(shù)據(jù),文中通過示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-07-07
  • 純注解版spring與mybatis的整合過程

    純注解版spring與mybatis的整合過程

    這篇文章主要介紹了純注解版spring與mybatis的整合過程,本文通過示例代碼給大家介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-06-06
  • 淺談springBoot注解大全

    淺談springBoot注解大全

    本篇文章主要介紹了淺談springBoot注解大全,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧
    2018-03-03
  • springboot中關(guān)于自動(dòng)建表,無法更新字段的問題

    springboot中關(guān)于自動(dòng)建表,無法更新字段的問題

    這篇文章主要介紹了springboot中關(guān)于自動(dòng)建表,無法更新字段的問題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2022-02-02
  • Java使用POI導(dǎo)出Excel(二):多個(gè)sheet

    Java使用POI導(dǎo)出Excel(二):多個(gè)sheet

    這篇文章介紹了Java使用POI導(dǎo)出Excel的方法,文中通過示例代碼介紹的非常詳細(xì)。對(duì)大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友可以參考下
    2022-10-10
  • MyBatis-Plus詳解(環(huán)境搭建、關(guān)聯(lián)操作)

    MyBatis-Plus詳解(環(huán)境搭建、關(guān)聯(lián)操作)

    MyBatis-Plus 是一個(gè) MyBatis 的增強(qiáng)工具,在 MyBatis 的基礎(chǔ)上只做增強(qiáng)不做改變,為簡(jiǎn)化開發(fā)、提高效率而生,今天通過本文給大家介紹MyBatis-Plus環(huán)境搭建及關(guān)聯(lián)操作,需要的朋友參考下吧
    2022-09-09
  • 反射機(jī)制:getDeclaredField和getField的區(qū)別說明

    反射機(jī)制:getDeclaredField和getField的區(qū)別說明

    這篇文章主要介紹了反射機(jī)制:getDeclaredField和getField的區(qū)別說明,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-06-06

最新評(píng)論

沈丘县| 华池县| 林西县| 天全县| 伊通| 密山市| 临猗县| 浦江县| 宁河县| 长岭县| 平江县| 平泉县| 蛟河市| 交口县| 体育| 嘉黎县| 环江| 名山县| 安远县| 上蔡县| 沙洋县| 河北省| 天镇县| 岳西县| 娄底市| 甘谷县| 彭阳县| 石楼县| 克东县| 镇原县| 兴安县| 张家口市| 元江| 安康市| 巴马| 五常市| 淅川县| 九龙坡区| 绥棱县| 泌阳县| 平昌县|