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

Springboot Websocket Stomp 消息訂閱推送

 更新時間:2021年07月09日 11:23:22   作者:代碼大師麥克勞瑞  
本文主要介紹了Springboot Websocket Stomp 消息訂閱推送,文中通過示例代碼介紹的非常詳細,對大家的學習或者工作具有一定的參考學習價值,需要的朋友們下面隨著小編來一起學習學習吧

需求背景

閑話不扯,直奔主題。需要和web前端建立長鏈接,互相實時通訊,因此想到了websocket,后面隨著需求的變更,需要用戶訂閱主題,實現(xiàn)消息的精準推送,發(fā)布訂閱等,則想到了STOMP(Simple Text-Orientated Messaging Protocol) 面向消息的簡單文本協(xié)議。

websocket協(xié)議

想到了之前寫的一個websocket長鏈接的demo,也貼上代碼供大家參考。

pom文件
直接引入spring-boot-starter-websocket即可。

    	<dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-websocket</artifactId>
        </dependency>

聲明websocket endpoint

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.web.socket.server.standard.ServerEndpointExporter;

/**
 * @ClassName WebSocketConfig
 * @Author scott
 * @Date 2021/6/16
 * @Version V1.0
 **/
@Configuration
public class WebSocketConfig {

    /**
     * 注入一個ServerEndpointExporter,該Bean會自動注冊使用@ServerEndpoint注解申明的websocket endpoint
     */
    @Bean
    public ServerEndpointExporter serverEndpointExporter() {
        return new ServerEndpointExporter();
    }

}

websocket實現(xiàn)類,其中通過注解監(jiān)聽了各種事件,實現(xiàn)了推送消息等相關(guān)邏輯

import com.google.common.cache.Cache;
import com.google.common.cache.CacheBuilder;
import com.ruoyi.common.core.domain.AjaxResult;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;

import javax.websocket.*;
import javax.websocket.server.PathParam;
import javax.websocket.server.ServerEndpoint;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

/**
 * @ClassName: DataTypePushWebSocket
 * @Author: scott
 * @Date: 2021/6/16
**/
@ServerEndpoint(value = "/ws/dataType/push/{token}")
@Component
public class DataTypePushWebSocket {

    private static final Logger log = LoggerFactory.getLogger(DataTypePushWebSocket.class);

    /**
     * 記錄當前在線連接數(shù)
     */
    private static AtomicInteger onlineCount = new AtomicInteger(0);

    private static Cache<String, Session> SESSION_CACHE = CacheBuilder.newBuilder()
            .initialCapacity(10)
            .maximumSize(300)
            .expireAfterWrite(10, TimeUnit.MINUTES)
            .build();

    /**
     * 連接建立成功調(diào)用的方法
     */
    @OnOpen
    public void onOpen(Session session, @PathParam("token")String token) {
        String sessionId = session.getId();
        onlineCount.incrementAndGet(); // 在線數(shù)加1
        this.sendMessage("sessionId:" + sessionId +",已經(jīng)和server建立連接", session);
        SESSION_CACHE.put(sessionId,session);
        log.info("有新連接加入:{},當前在線連接數(shù)為:{}", session.getId(), onlineCount.get());
    }

    /**
     * 連接關(guān)閉調(diào)用的方法
     */
    @OnClose
    public void onClose(Session session,@PathParam("token")String token) {
        onlineCount.decrementAndGet(); // 在線數(shù)減1
        SESSION_CACHE.invalidate(session.getId());
        log.info("有一連接關(guān)閉:{},當前在線連接數(shù)為:{}", session.getId(), onlineCount.get());
    }

    /**
     * 收到客戶端消息后調(diào)用的方法
     *
     * @param message 客戶端發(fā)送過來的消息
     */
    @OnMessage
    public void onMessage(String message, Session session,@PathParam("token")String token) {
        log.info("服務端收到客戶端[{}]的消息:{}", session.getId(), message);
        this.sendMessage("服務端已收到推送消息:" + message, session);
    }

    @OnError
    public void onError(Session session, Throwable error) {
        log.error("發(fā)生錯誤");
        error.printStackTrace();
    }

    /**
     * 服務端發(fā)送消息給客戶端
     */
    private static void sendMessage(String message, Session toSession) {
        try {
            log.info("服務端給客戶端[{}]發(fā)送消息{}", toSession.getId(), message);
            toSession.getBasicRemote().sendText(message);
        } catch (Exception e) {
            log.error("服務端發(fā)送消息給客戶端失敗:{}", e);
        }
    }

    public static AjaxResult sendMessage(String message, String sessionId){
        Session session = SESSION_CACHE.getIfPresent(sessionId);
        if(Objects.isNull(session)){
            return AjaxResult.error("token已失效");
        }
        sendMessage(message,session);
        return AjaxResult.success();
    }

    public static AjaxResult sendBroadcast(String message){
        long size = SESSION_CACHE.size();
        if(size <=0){
            return AjaxResult.error("當前沒有在線客戶端,無法推送消息");
        }
        ConcurrentMap<String, Session> sessionConcurrentMap = SESSION_CACHE.asMap();
        Set<String> keys = sessionConcurrentMap.keySet();
        for (String key : keys) {
            Session session = SESSION_CACHE.getIfPresent(key);
            DataTypePushWebSocket.sendMessage(message,session);
        }

        return AjaxResult.success();

    }

}

至此websocket服務端代碼已經(jīng)完成。

stomp協(xié)議

前端代碼.這個是在某個vue工程中寫的js,各位大佬自己動手改改即可。其中Settings.wsPath是后端定義的ws地址例如ws://localhost:9003/ws

import Stomp from 'stompjs'
import Settings from '@/settings.js'

export default {
  // 是否啟用日志 默認啟用
  debug:true,
  // 客戶端連接信息
  stompClient:{},
  // 初始化
  init(callBack){
    this.stompClient = Stomp.client(Settings.wsPath)
    this.stompClient.hasDebug = this.debug
    this.stompClient.connect({},suce =>{
      this.console("連接成功,信息如下 ↓")
      this.console(this.stompClient)
      if(callBack){
        callBack()
      }
    },err => {
      if(err) {
        this.console("連接失敗,信息如下 ↓")
        this.console(err)
      }
    })
  },
  // 訂閱
  sub(address,callBack){
    if(!this.stompClient.connected){
      this.console("沒有連接,無法訂閱")
      return
    }
    // 生成 id
    let timestamp= new Date().getTime() + address
    this.console("訂閱成功 -> "+address)
    this.stompClient.subscribe(address,message => {
      this.console(address+" 訂閱消息通知,信息如下 ↓")
      this.console(message)
      let data = message.body
      callBack(data)
    },{
      id: timestamp
    })
  },
  unSub(address){
    if(!this.stompClient.connected){
      this.console("沒有連接,無法取消訂閱 -> "+address)
      return
    }
    let id = ""
    for(let item in this.stompClient.subscriptions){
      if(item.endsWith(address)){
        id = item
        break
      }
    }
    this.stompClient.unsubscribe(id)
    this.console("取消訂閱成功 -> id:"+ id + " address:"+address)
  },
  // 斷開連接
  disconnect(callBack){
    if(!this.stompClient.connected){
      this.console("沒有連接,無法斷開連接")
      return
    }
    this.stompClient.disconnect(() =>{
      console.log("斷開成功")
      if(callBack){
        callBack()
      }
    })
  },
  // 單位 秒
  reconnect(time){
    setInterval(() =>{
      if(!this.stompClient.connected){
        this.console("重新連接中...")
        this.init()
      }
    },time * 1000)
  },
  console(msg){
    if(this.debug){
      console.log(msg)
    }
  },
  // 向訂閱發(fā)送消息
  send(address,msg) {
    this.stompClient.send(address,{},msg)
  }
}

后端stomp config,里面都有注釋,寫的很詳細,并且我加入了和前端的心跳ping pong。

package com.cn.scott.config;

import org.springframework.context.annotation.Configuration;
import org.springframework.messaging.simp.config.MessageBrokerRegistry;
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
import org.springframework.web.socket.config.annotation.EnableWebSocketMessageBroker;
import org.springframework.web.socket.config.annotation.StompEndpointRegistry;
import org.springframework.web.socket.config.annotation.WebSocketMessageBrokerConfigurer;

/**
 * @ClassName: WebSocketStompConfig
 * @Author: scott
 * @Date: 2021/7/8
**/
@Configuration
@EnableWebSocketMessageBroker
public class WebSocketStompConfig implements WebSocketMessageBrokerConfigurer {

    private static long HEART_BEAT=10000;

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        //允許使用socketJs方式訪問,訪問點為webSocket,允許跨域
        //在網(wǎng)頁上我們就可以通過這個鏈接
        //ws://127.0.0.1:port/ws來和服務器的WebSocket連接
        registry.addEndpoint("/ws").setAllowedOrigins("*");
    }

    @Override
    public void configureMessageBroker(MessageBrokerRegistry registry) {
        ThreadPoolTaskScheduler te = new ThreadPoolTaskScheduler();
        te.setPoolSize(1);
        te.setThreadNamePrefix("wss-heartbeat-thread-");
        te.initialize();
        //基于內(nèi)存的STOMP消息代理來代替mq的消息代理
        //訂閱Broker名稱,/user代表點對點即發(fā)指定用戶,/topic代表發(fā)布廣播即群發(fā)
        //setHeartbeatValue 設(shè)置心跳及心跳時間
        registry.enableSimpleBroker("/user", "/topic").setHeartbeatValue(new long[]{HEART_BEAT,HEART_BEAT}).setTaskScheduler(te);
        //點對點使用的訂閱前綴,不設(shè)置的話,默認也是/user/
        registry.setUserDestinationPrefix("/user/");
    }
}

后端stomp協(xié)議接受、訂閱等動作通知

package com.cn.scott.ws;

import com.alibaba.fastjson.JSON;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.messaging.handler.annotation.DestinationVariable;
import org.springframework.messaging.handler.annotation.MessageMapping;
import org.springframework.messaging.simp.SimpMessagingTemplate;
import org.springframework.messaging.simp.annotation.SubscribeMapping;
import org.springframework.web.bind.annotation.RestController;

/**
 * @ClassName StompSocketHandler
 * @Author scott
 * @Date 2021/6/30
 * @Version V1.0
 **/
@RestController
public class StompSocketHandler {

    @Autowired
    private SimpMessagingTemplate simpMessagingTemplate;

    /**
    * @MethodName: subscribeMapping
     * @Description: 訂閱成功通知
     * @Param: [id]
     * @Return: void
     * @Author: scott
     * @Date: 2021/6/30
    **/
    @SubscribeMapping("/user/{id}/listener")
    public void subscribeMapping(@DestinationVariable("id") final long id) {
        System.out.println(">>>>>>用戶:"+id +",已訂閱");
        SubscribeMsg param = new SubscribeMsg(id,String.format("用戶【%s】已訂閱成功", id));
        sendToUser(param);
    }


    /**
    * @MethodName: test
     * @Description: 接收訂閱topic消息
     * @Param: [id, msg]
     * @Return: void
     * @Author: scott
     * @Date: 2021/6/30
    **/
    @MessageMapping(value = "/user/{id}/listener")
    public void UserSubListener(@DestinationVariable long  id, String msg) {
        System.out.println("收到客戶端:" +id+",的消息");
        SubscribeMsg param = new SubscribeMsg(id,String.format("已收到用戶【%s】發(fā)送消息【%s】", id,msg));
        sendToUser(param);
    }
    
     @GetMapping("/refresh/{userId}")
    public void refresh(@PathVariable Long userId, String msg) {
        StompSocketHandler.SubscribeMsg param = new StompSocketHandler.SubscribeMsg(userId,String.format("服務端向用戶【%s】發(fā)送消息【%s】", userId,msg));
        sendToUser(param);
    }

    /**
    * @MethodName: sendToUser
     * @Description: 推送消息給訂閱用戶
     * @Param: [userId]
     * @Return: void
     * @Author: scott
     * @Date: 2021/6/30
    **/
    public void sendToUser(SubscribeMsg screenChangeMsg){
        //這里可以控制權(quán)限等。。。
        simpMessagingTemplate.convertAndSendToUser(screenChangeMsg.getUserId().toString(),"/listener", JSON.toJSONString(screenChangeMsg));
    }

    /**
    * @MethodName: sendBroadCast
     * @Description: 發(fā)送廣播,需要用戶事先訂閱廣播
     * @Param: [topic, msg]
     * @Return: void
     * @Author: scott
     * @Date: 2021/6/30
    **/
    public void sendBroadCast(String topic,String msg){
        simpMessagingTemplate.convertAndSend(topic,msg);
    }


    /**
     * @ClassName: SubMsg
     * @Author: scott
     * @Date: 2021/6/30
    **/
    public static class SubscribeMsg {
        private Long userId;
        private String msg;
        public SubscribeMsg(Long UserId, String msg){
            this.userId = UserId;
            this.msg = msg;
        }
        public Long getUserId() {
            return userId;
        }
        public String getMsg() {
            return msg;
        }
    }
}

連接展示

建立連接成功,這里可以看出是基于websocket協(xié)議

在這里插入圖片描述

連接信息

在這里插入圖片描述

ping pong

在這里插入圖片描述

調(diào)用接口向訂閱用戶1發(fā)送消息,http://localhost:9003/refresh/1?msg=HelloStomp,可以在客戶端控制臺查看已經(jīng)收到了消息。這個時候不同用戶通過自己的userId可以區(qū)分訂閱的主題,可以做到通過userId精準的往客戶端推送消息。

在這里插入圖片描述

還記得我們在后端配置的時候還指定了廣播的訂閱主題/topic,這時我們前端通過js只要訂閱了這個主題,那么后端在像這個主題推送消息時,所有訂閱的客戶端都能收到,感興趣的小伙伴可以自己試試,api我都寫好了。

在這里插入圖片描述

至此,實戰(zhàn)完畢,喜歡的小伙伴麻煩關(guān)注加點贊。

springboot + stomp后端源碼地址:https://gitee.com/ErGouGeSiBaKe/stomp-server

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

相關(guān)文章

  • 聊聊注解@controller@service@component@repository的區(qū)別

    聊聊注解@controller@service@component@repository的區(qū)別

    這篇文章主要介紹了聊聊注解@controller@service@component@repository的區(qū)別,具有很好的參考價值,希望對大家有所幫助。如有錯誤或未考慮完全的地方,望不吝賜教
    2021-08-08
  • SpringBoot自動裝配之Condition深入講解

    SpringBoot自動裝配之Condition深入講解

    @Conditional表示僅當所有指定條件都匹配時,組件才有資格注冊。該@Conditional注釋可以在以下任一方式使用:作為任何@Bean方法的方法級注釋、作為任何類的直接或間接注釋的類型級別注釋@Component,包括@Configuration類、作為元注釋,目的是組成自定義構(gòu)造型注釋
    2023-01-01
  • scala 匿名函數(shù)案例詳解

    scala 匿名函數(shù)案例詳解

    Scala支持一級函數(shù),函數(shù)可以用函數(shù)文字語法表達,即(x:Int)=> x + 1,該函數(shù)可以由一個叫作函數(shù)值的對象來表示,這篇文章主要介紹了scala 匿名函數(shù)詳解,需要的朋友可以參考下
    2023-03-03
  • Seata?AT獲取數(shù)據(jù)表元數(shù)據(jù)源碼詳解

    Seata?AT獲取數(shù)據(jù)表元數(shù)據(jù)源碼詳解

    這篇文章主要為大家介紹了Seata?AT獲取數(shù)據(jù)表元數(shù)據(jù)源碼詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進步,早日升職加薪
    2022-11-11
  • Java中調(diào)用DLL動態(tài)庫的操作方法

    Java中調(diào)用DLL動態(tài)庫的操作方法

    在Java編程中,有時我們需要調(diào)用本地代碼庫,特別是Windows平臺上的DLL(動態(tài)鏈接庫),本文中,我們將詳細討論如何在Java中加載和調(diào)用DLL動態(tài)庫,并通過具體示例來展示這個過程,感興趣的朋友跟隨小編一起看看吧
    2024-03-03
  • springboot整合quartz項目使用案例

    springboot整合quartz項目使用案例

    quartz是一個定時調(diào)度的框架,就目前市場上來說,其實有比quartz更優(yōu)秀的一些定時調(diào)度框架,不但性能比quartz好,學習成本更低,而且還提供可視化操作定時任務,這篇文章主要介紹了springboot整合quartz項目使用(含完整代碼),需要的朋友可以參考下
    2023-05-05
  • Java生產(chǎn)者消費者模式實例分析

    Java生產(chǎn)者消費者模式實例分析

    這篇文章主要介紹了Java生產(chǎn)者消費者模式,結(jié)合實例形式分析了java生產(chǎn)者消費者模式的相關(guān)組成、原理及實現(xiàn)方法,需要的朋友可以參考下
    2019-03-03
  • SpringBoot3集成Swagger3的詳細教程

    SpringBoot3集成Swagger3的詳細教程

    Swagger 3(OpenAPI 3.0)提供了更加強大和靈活的API文檔生成能力,本教程將指導您如何在Spring Boot 3項目中集成Swagger3,并使用Knife4j作為UI界面,需要的朋友可以參考下
    2024-03-03
  • MapStruct處理Java中實體與模型間不匹配屬性轉(zhuǎn)換的方法

    MapStruct處理Java中實體與模型間不匹配屬性轉(zhuǎn)換的方法

    今天小編就為大家分享一篇關(guān)于MapStruct處理Java中實體與模型間不匹配屬性轉(zhuǎn)換的方法,小編覺得內(nèi)容挺不錯的,現(xiàn)在分享給大家,具有很好的參考價值,需要的朋友一起跟隨小編來看看吧
    2019-03-03
  • 解決java main函數(shù)中的args數(shù)組傳值問題

    解決java main函數(shù)中的args數(shù)組傳值問題

    這篇文章主要介紹了解決java main函數(shù)中的args數(shù)組傳值問題,具有很好的參考價值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2021-02-02

最新評論

年辖:市辖区| 鹤壁市| 绵竹市| 额尔古纳市| 灌阳县| 开封市| 海城市| 台北市| 禄丰县| 方城县| 乾安县| 哈巴河县| 安塞县| 瑞金市| 黔南| 铜梁县| 徐州市| 炉霍县| 广宗县| 卫辉市| 桂平市| 隆尧县| 如皋市| 朔州市| 武陟县| 乡城县| 大邑县| 梅河口市| 石楼县| 茌平县| 正安县| 盐津县| 红桥区| 哈尔滨市| 黎平县| 鄂温| 桦川县| 上饶市| 高尔夫| 句容市| 金山区|