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

Spring?Boot?中實(shí)現(xiàn)?WebSocket?集群思路詳解

 更新時(shí)間:2025年12月10日 16:44:00   作者:黃三技術(shù)java  
本文介紹了如何在SpringBoot中實(shí)現(xiàn)WebSocket集群處理,解決了會(huì)話共享和消息同步問(wèn)題,通過(guò)分布式存儲(chǔ)(Redis)和消息隊(duì)列(RedisPub/Sub)實(shí)現(xiàn)跨節(jié)點(diǎn)通信,感興趣的朋友跟隨小編一起看看吧

在 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)方案:

一、核心思路

  1. 會(huì)話共享:使用分布式存儲(chǔ)(如 Redis)存儲(chǔ)所有節(jié)點(diǎn)的 WebSocket 會(huì)話信息(會(huì)話 ID、用戶標(biāo)識(shí)、節(jié)點(diǎn)標(biāo)識(shí)等)。
  2. 消息轉(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)推送消息給用戶。
  3. 集群感知:每個(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è)試

  1. 啟動(dòng)多個(gè) Spring Boot 實(shí)例(模擬集群節(jié)點(diǎn)),確保都連接到同一個(gè) Redis。
  2. 客戶端 1 連接節(jié)點(diǎn) A,客戶端 2 連接節(jié)點(diǎn) B。
  3. 客戶端 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)化建議

  1. 會(huì)話過(guò)期清理:給 Redis 中的USER_NODE_MAP設(shè)置過(guò)期時(shí)間(如 30 分鐘),并通過(guò) WebSocket 心跳機(jī)制續(xù)約。
  2. 消息可靠性:Redis Pub/Sub 是 fire-and-forget 模式,若需可靠消息,可替換為 RabbitMQ 或 Kafka。
  3. 負(fù)載均衡:前端通過(guò) Nginx 等負(fù)載均衡器連接 WebSocket 集群(需配置proxy_set_header Upgrade $http_upgrade;等參數(shù)支持 WebSocket)。
  4. 監(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)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • 一文讓你了解透徹Java中的IO模型

    一文讓你了解透徹Java中的IO模型

    本文只是說(shuō)明了IO模型,讓你了解IO模型是什么,怎么區(qū)分IO模型,以及分析了Java中的三種IO模型,本文是純理論知識(shí),看完之后會(huì)讓你對(duì)IO有更加深刻的理解,感興趣的同學(xué)可以參考一下
    2023-05-05
  • 基于Java編寫(xiě)一個(gè)html轉(zhuǎn)pdf的工具類(lèi)

    基于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)含義

    Java中DecimalFormat用法及符號(hào)含義

    DecimalFormat是NumberFormat的一個(gè)具體子類(lèi),用于格式化十進(jìn)制數(shù)字。這篇文章介紹了DecimalFormat的用法及符號(hào)含義,需要的朋友可以收藏下,方便下次瀏覽觀看
    2021-12-12
  • Java中Map與對(duì)象之間互相轉(zhuǎn)換的幾種常用方式

    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)指南

    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ì)列原理剖析

    java并發(fā)中DelayQueue延遲隊(duì)列原理剖析

    DelayQueue隊(duì)列是一個(gè)延遲隊(duì)列,本文將結(jié)合實(shí)例代碼,詳細(xì)的介紹DelayQueue延遲隊(duì)列的源碼分析,感興趣的小伙伴們可以參考一下
    2021-06-06
  • Java中MapStruct 映射過(guò)程中忽略某個(gè)字段的實(shí)現(xiàn)

    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)量Semaphore詳細(xì)解讀,Java信號(hào)量機(jī)制可以用來(lái)保證線程互斥,創(chuàng)建Semaphore對(duì)象傳入一個(gè)整形參數(shù),類(lèi)似于公共資源,需要的朋友可以參考下
    2023-11-11
  • 簡(jiǎn)單總結(jié)單例模式的4種寫(xiě)法

    簡(jiǎn)單總結(jié)單例模式的4種寫(xiě)法

    今天帶大家學(xué)習(xí)java的相關(guān)知識(shí),文章圍繞著單例模式的4種寫(xiě)法展開(kāi),文中有非常詳細(xì)的介紹及代碼示例,需要的朋友可以參考下
    2021-06-06
  • Java SerialVersionUID作用詳解

    Java SerialVersionUID作用詳解

    這篇文章主要介紹了Java SerialVersionUID作用詳解,本篇文章通過(guò)簡(jiǎn)要的案例,講解了該項(xiàng)技術(shù)的了解與使用,以下就是詳細(xì)內(nèi)容,需要的朋友可以參考下
    2021-08-08

最新評(píng)論

耿马| 晋中市| 红桥区| 葵青区| 营山县| 辽阳市| 中江县| 应用必备| 晴隆县| 延吉市| 青海省| 嫩江县| 澄迈县| 四子王旗| 谷城县| 裕民县| 芮城县| 黄浦区| 开封市| 渑池县| 洛隆县| 蓝田县| 佛教| 宁陵县| 康乐县| 略阳县| 岗巴县| 理塘县| 监利县| 达州市| 滁州市| 颍上县| 元江| 策勒县| 贵阳市| 龙门县| 库车县| 许昌县| 丰台区| 赤壁市| 谢通门县|