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

java中消息推送功能的關(guān)鍵實(shí)現(xiàn)代碼

 更新時間:2025年09月30日 09:18:21   作者:hqxstudying  
在日常開發(fā)中,消息推送是非常典型的業(yè)務(wù)需求,下面這篇文章主要介紹了java中消息推送功能的關(guān)鍵實(shí)現(xiàn)代碼,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下

前言

在 Java 項目中實(shí)現(xiàn)消息推送服務(wù),后端核心需要解決實(shí)時性、可靠性、擴(kuò)展性(集群支持)和連接管理四大問題。以下從技術(shù)選型、核心架構(gòu)、關(guān)鍵實(shí)現(xiàn)和注意事項四個方面展開說明:

一、技術(shù)選型:根據(jù)場景選協(xié)議

消息推送的核心是服務(wù)器主動向客戶端發(fā)送數(shù)據(jù),需根據(jù)實(shí)時性、客戶端類型(Web/APP/ 物聯(lián)網(wǎng)設(shè)備)選擇合適的通信協(xié)議:

協(xié)議 / 方案適用場景優(yōu)勢劣勢Java 技術(shù)棧支持
WebSocketWeb 端實(shí)時通信(如聊天、通知)全雙工、低延遲、長連接部分老舊瀏覽器不支持Spring WebSocket、Netty、Tomcat 原生
MQTT物聯(lián)網(wǎng)設(shè)備(低帶寬、不穩(wěn)定網(wǎng)絡(luò))輕量、支持 QoS(消息質(zhì)量等級)需額外部署 MQTT broker(如 EMQX)Eclipse Paho、Spring Integration
長輪詢(Long Polling)兼容性要求高的場景(如老瀏覽器)實(shí)現(xiàn)簡單、兼容性好延遲較高、服務(wù)器資源消耗大Servlet + 異步處理
Server-Sent Events (SSE)服務(wù)器單向推送(如實(shí)時日志)輕量、僅服務(wù)器向客戶端推送不支持客戶端向服務(wù)器發(fā)送數(shù)據(jù)Spring WebFlux、原生 Servlet

二、核心架構(gòu):后端模塊設(shè)計

無論選擇哪種協(xié)議,消息推送服務(wù)的后端架構(gòu)通常包含以下核心模塊:

┌─────────────────┐    ┌─────────────────┐    ┌─────────────────┐
│   業(yè)務(wù)系統(tǒng)API   │───>│   消息路由模塊   │───>│   連接管理模塊   │───> 客戶端
└─────────────────┘    └─────────────────┘    └─────────────────┘
        │                       │                       │
        ▼                       ▼                       ▼
┌─────────────────┐    ┌─────────────────┐    ┌─────────────────┐
│  離線消息存儲    │<───│  集群通信模塊    │<───│  會話管理模塊    │
└─────────────────┘    └─────────────────┘    └─────────────────┘
  • 業(yè)務(wù)系統(tǒng) API:提供接口給業(yè)務(wù)系統(tǒng)(如訂單系統(tǒng)、通知系統(tǒng))調(diào)用,觸發(fā)消息推送(例如:pushMessage(String userId, String content))。

  • 會話管理模塊:維護(hù)客戶端與服務(wù)器的連接會話(Session),核心是用戶 ID 與連接的映射關(guān)系。

    • 單機(jī):用ConcurrentHashMap<String, Session>存儲(key 為用戶 ID,value 為連接會話)。
    • 集群:用 Redis 存儲分布式會話(需序列化 Session 信息,或僅存儲 “用戶 ID - 節(jié)點(diǎn) IP” 映射)。
  • 連接管理模塊:處理客戶端的連接建立、斷開、心跳檢測。

    • 例如 WebSocket 的onOpen()(建立連接時綁定用戶 ID 與 Session)、onClose()(移除映射)、onError()(異常處理)。
  • 消息路由模塊:根據(jù)目標(biāo)用戶 ID,找到對應(yīng)的連接會話并發(fā)送消息。

    • 單機(jī):直接從本地ConcurrentHashMap獲取 Session 發(fā)送。
    • 集群:通過 Redis Pub/Sub 廣播消息,所有節(jié)點(diǎn)收到后檢查本地是否有目標(biāo)用戶的連接,有則發(fā)送。
  • 離線消息模塊:當(dāng)用戶不在線時,將消息暫存(如 MySQL、MongoDB 或 Kafka),待用戶上線后拉取。

  • 集群通信模塊:解決多節(jié)點(diǎn)部署時的消息同步問題(如 Redis Pub/Sub、RabbitMQ 的 Fanout 交換機(jī))。

三、關(guān)鍵實(shí)現(xiàn):以 WebSocket + Spring Boot 為例

以最常用的 Web 端實(shí)時推送為例,基于 Spring Boot + WebSocket + Redis(集群支持)實(shí)現(xiàn)核心流程:

1. 依賴配置(Maven)

<!-- WebSocket核心依賴 -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-websocket</artifactId>
</dependency>
<!-- Redis(集群通信與會話存儲) -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-redis</artifactId>
</dependency>

2. WebSocket 配置(連接建立與處理器)

@Configuration
@EnableWebSocket
public class WebSocketConfig implements WebSocketConfigurer {
    @Autowired
    private MessageHandler messageHandler; // 自定義消息處理器

    @Override
    public void registerWebSocketHandlers(WebSocketHandlerRegistry registry) {
        // 配置WebSocket端點(diǎn):客戶端通過ws://ip:port/ws?token=xxx連接
        registry.addHandler(messageHandler, "/ws")
                .setAllowedOrigins("*") // 允許跨域(生產(chǎn)環(huán)境需限制域名)
                .addInterceptors(new HandshakeInterceptor() {
                    // 握手前驗(yàn)證用戶(如Token解析用戶ID)
                    @Override
                    public boolean beforeHandshake(ServerHttpRequest request, ServerHttpResponse response, 
                                                  WebSocketHandler wsHandler, Map<String, Object> attributes) {
                        String token = ((ServletServerHttpRequest) request).getServletRequest().getParameter("token");
                        String userId = parseToken(token); // 解析Token獲取用戶ID
                        if (userId == null) {
                            response.setStatusCode(HttpStatus.UNAUTHORIZED);
                            return false; // 驗(yàn)證失敗,拒絕連接
                        }
                        attributes.put("userId", userId); // 存儲用戶ID到屬性中
                        return true;
                    }

                    @Override
                    public void afterHandshake(...) {}
                });
    }
}

3. 消息處理器(連接管理與會話綁定)

@Component
public class MessageHandler extends TextWebSocketHandler {
    // 本地會話映射:用戶ID -> WebSocket會話(單機(jī)用)
    private final Map<String, WebSocketSession> localSessions = new ConcurrentHashMap<>();
    @Autowired
    private StringRedisTemplate redisTemplate; // Redis操作模板
    @Autowired
    private OfflineMessageService offlineMessageService; // 離線消息服務(wù)

    // 連接建立時:綁定用戶ID與Session
    @Override
    public void afterConnectionEstablished(WebSocketSession session) {
        String userId = (String) session.getAttributes().get("userId");
        localSessions.put(userId, session);
        
        // 1. 拉取離線消息并推送
        List<Message> offlineMessages = offlineMessageService.getByUserId(userId);
        for (Message msg : offlineMessages) {
            sendMessage(session, msg.getContent());
        }
        offlineMessageService.deleteByUserId(userId); // 清除已推送的離線消息
        
        // 2. 集群:將用戶ID與當(dāng)前節(jié)點(diǎn)IP存入Redis(便于其他節(jié)點(diǎn)定位)
        redisTemplate.opsForValue().set("user:session:" + userId, getLocalIp());
    }

    // 連接關(guān)閉時:移除會話映射
    @Override
    public void afterConnectionClosed(WebSocketSession session, CloseStatus status) {
        String userId = (String) session.getAttributes().get("userId");
        localSessions.remove(userId);
        redisTemplate.delete("user:session:" + userId); // 集群:刪除Redis映射
    }

    // 發(fā)送消息到指定用戶(核心方法)
    public void pushToUser(String userId, String content) {
        // 1. 先檢查本地是否有該用戶的連接
        WebSocketSession session = localSessions.get(userId);
        if (session != null && session.isOpen()) {
            sendMessage(session, content);
            return;
        }

        // 2. 集群:檢查用戶是否連接在其他節(jié)點(diǎn)
        String targetNodeIp = redisTemplate.opsForValue().get("user:session:" + userId);
        if (targetNodeIp != null) {
            // 發(fā)布消息到Redis頻道,目標(biāo)節(jié)點(diǎn)訂閱后發(fā)送
            redisTemplate.convertAndSend("push:channel", 
                JSON.toJSONString(new PushMessage(userId, content)));
            return;
        }

        // 3. 用戶不在線:存入離線消息
        offlineMessageService.save(new Message(userId, content, LocalDateTime.now()));
    }

    // 實(shí)際發(fā)送消息(封裝異常處理)
    private void sendMessage(WebSocketSession session, String content) {
        try {
            session.sendMessage(new TextMessage(content));
        } catch (IOException e) {
            log.error("消息發(fā)送失敗", e);
        }
    }
}

4. 集群支持:Redis Pub/Sub 訂閱

@Component
public class RedisMessageListener {
    @Autowired
    private MessageHandler messageHandler;

    @Bean
    public RedisMessageListenerContainer container(RedisConnectionFactory connectionFactory) {
        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
        container.setConnectionFactory(connectionFactory);
        // 訂閱推送頻道,接收其他節(jié)點(diǎn)的消息
        container.addMessageListener((message, pattern) -> {
            String json = new String(message.getBody());
            PushMessage pushMsg = JSON.parseObject(json, PushMessage.class);
            // 調(diào)用本地消息處理器發(fā)送(此時用戶應(yīng)在當(dāng)前節(jié)點(diǎn))
            messageHandler.pushToUser(pushMsg.getUserId(), pushMsg.getContent());
        }, new ChannelTopic("push:channel"));
        return container;
    }
}

5. 業(yè)務(wù)調(diào)用接口

@RestController
@RequestMapping("/api/message")
public class MessageController {
    @Autowired
    private MessageHandler messageHandler;

    // 業(yè)務(wù)系統(tǒng)調(diào)用此接口觸發(fā)推送
    @PostMapping("/push")
    public Result push(@RequestBody PushRequest request) {
        messageHandler.pushToUser(request.getUserId(), request.getContent());
        return Result.success();
    }
}

四、注意事項

  • 安全性

    • 連接建立時必須驗(yàn)證用戶身份(如 Token、Session),防止未授權(quán)連接。
    • 生產(chǎn)環(huán)境需限制 WebSocket 的跨域來源(setAllowedOrigins不要用*)。
  • 可靠性

    • 實(shí)現(xiàn)消息確認(rèn)機(jī)制(客戶端收到消息后回復(fù) ACK,服務(wù)器未收到則重試)。
    • 離線消息需持久化(建議用 Kafka 或數(shù)據(jù)庫,支持消息過期清理)。
  • 性能與擴(kuò)展

    • 長連接數(shù)量大時,用 Netty 替代 Tomcat 原生 WebSocket(Netty 的 NIO 模型更高效)。
    • 集群部署時,通過 Redis 或 MQ 實(shí)現(xiàn)消息廣播,避免 “消息孤島”。
  • 監(jiān)控與運(yùn)維

    • 監(jiān)控連接數(shù)、消息發(fā)送成功率、節(jié)點(diǎn)負(fù)載(如用 Prometheus + Grafana)。
    • 實(shí)現(xiàn)連接心跳檢測(定期發(fā)送 ping 幀,超時未響應(yīng)則主動斷開)。

總結(jié)

后端消息推送服務(wù)的核心是 **“連接管理”+“消息路由”**,小規(guī)模場景可用 Spring WebSocket + 本地會話;中大規(guī)模集群需結(jié)合 Redis 實(shí)現(xiàn)分布式會話與跨節(jié)點(diǎn)通信;物聯(lián)網(wǎng)場景優(yōu)先選 MQTT 協(xié)議。根據(jù)實(shí)時性和規(guī)模需求,可逐步迭代優(yōu)化(從單機(jī)到集群,從基礎(chǔ)功能到可靠性保障)。

到此這篇關(guān)于java中消息推送功能的文章就介紹到這了,更多相關(guān)java消息推送功能內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • IDEA不能自動導(dǎo)包的問題及解決方案

    IDEA不能自動導(dǎo)包的問題及解決方案

    這篇文章主要介紹了IDEA不能自動導(dǎo)包的問題及解決方案,具有很好的參考價值,希望對大家有所幫助,如有錯誤或未考慮完全的地方,望不吝賜教
    2026-03-03
  • 持久層ORM框架Hibernate框架的使用及搭建方式

    持久層ORM框架Hibernate框架的使用及搭建方式

    Hibernate是一個開放源代碼的對象關(guān)系映射框架,它對JDBC進(jìn)行了非常輕量級的對象封裝,使得Java程序員可以隨心所欲的使用對象編程思維來操縱數(shù)據(jù)庫,本文重點(diǎn)給大家介紹持久層ORM框架Hibernate框架的使用及搭建方式,感興趣的朋友一起看看吧
    2021-11-11
  • Java操作FTP實(shí)現(xiàn)上傳下載功能

    Java操作FTP實(shí)現(xiàn)上傳下載功能

    這篇文章主要為大家詳細(xì)介紹了Java如何通過操作FTP實(shí)現(xiàn)上傳下載的功能,文中的示例代碼講解詳細(xì),對我們學(xué)習(xí)Java有一定幫助,需要的可以參考一下
    2022-11-11
  • Java?Web開發(fā)中的分頁與參數(shù)校驗(yàn)舉例詳解

    Java?Web開發(fā)中的分頁與參數(shù)校驗(yàn)舉例詳解

    這篇文章主要介紹了JavaWeb開發(fā)中的分頁設(shè)計和參數(shù)校驗(yàn),分頁設(shè)計通過分頁查詢參數(shù)優(yōu)化查詢性能,文中通過代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2025-02-02
  • springboot項目中全局設(shè)置用UTC+8

    springboot項目中全局設(shè)置用UTC+8

    本文主要介紹了springboot項目中全局設(shè)置用UTC+8,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2023-04-04
  • mybatis 使用jdbc.properties文件設(shè)置不起作用的解決方法

    mybatis 使用jdbc.properties文件設(shè)置不起作用的解決方法

    這篇文章主要介紹了mybatis 使用jdbc.properties文件設(shè)置不起作用的解決方法,需要的朋友可以參考下
    2018-03-03
  • java TreeMap源碼解析詳解

    java TreeMap源碼解析詳解

    這篇文章主要介紹了java TreeMap源碼解析詳解的相關(guān)資料,需要的朋友可以參考下
    2017-04-04
  • Java 超詳細(xì)講解數(shù)據(jù)結(jié)構(gòu)的應(yīng)用

    Java 超詳細(xì)講解數(shù)據(jù)結(jié)構(gòu)的應(yīng)用

    數(shù)據(jù)結(jié)構(gòu)是計算機(jī)存儲、組織數(shù)據(jù)的方式。數(shù)據(jù)結(jié)構(gòu)是指相互之間存在一種或多種特定關(guān)系的數(shù)據(jù)元素的集合,讓我們一起來了解數(shù)據(jù)結(jié)構(gòu)是如何應(yīng)用的
    2022-04-04
  • 使用Java編寫一個輸出九九口訣乘法表的程序

    使用Java編寫一個輸出九九口訣乘法表的程序

    在學(xué)習(xí)編程的過程中,編寫簡單的程序來實(shí)現(xiàn)基本的數(shù)學(xué)運(yùn)算是一個很好的練習(xí),本文將介紹如何使用Java語言編寫一個程序,用于輸出9*9的乘法口訣表,需要的朋友可以參考下
    2026-01-01
  • 將JSON字符串?dāng)?shù)組轉(zhuǎn)對象集合方法步驟

    將JSON字符串?dāng)?shù)組轉(zhuǎn)對象集合方法步驟

    這篇文章主要給大家介紹了關(guān)于將JSON字符串?dāng)?shù)組轉(zhuǎn)對象集合的方法步驟,文中通過代碼示例介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考借鑒價值,需要的朋友可以參考下
    2023-08-08

最新評論

洪泽县| 弥渡县| 兰西县| 平利县| 海安县| 竹山县| 孝昌县| 咸宁市| 上栗县| 永靖县| 富源县| 雅江县| 岳普湖县| 黄山市| 苏尼特左旗| 个旧市| 大石桥市| 日照市| 高邑县| 晋宁县| 石狮市| 长子县| 北海市| 红原县| 巴青县| 贵阳市| 鹤壁市| 达拉特旗| 宝清县| 怀柔区| 交口县| 原平市| 慈溪市| 利津县| 成安县| 益阳市| 横峰县| 杭州市| 汽车| 通渭县| 台山市|