IT数码 购物 网址 头条 软件 日历 阅读 图书馆
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
图片批量下载器
↓批量下载图片,美女图库↓
图片自动播放器
↓图片自动播放器↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁
 
   -> 网络协议 -> springboot使用redisTemplate+websocket实现集群消息的发布订阅 -> 正文阅读

[网络协议]springboot使用redisTemplate+websocket实现集群消息的发布订阅

本文主要介绍小项目中如何通过redisTemplate+websocket实现后端向前端推送消息

一、websocket配置类

@Component
public class WebsocketConfig {
    /**
     * ServerEndpointExporter 作用
     * <p>
     * 这个Bean会自动注册使用@ServerEndpoint注解声明的websocket endpoint
     *
     * @return
     */
    @Bean
    public ServerEndpointExporter serverEndpointExporter() {
        return new ServerEndpointExporter();
    }
}

二、websocket接口类

@Component
@ServerEndpoint(value = "/websocket")
public class WebSocket {

    private static final Logger logger = LoggerFactory.getLogger(WebSocket.class);

    // 每个客户端都会有相应的session,服务端可以发送相关消息
    private Session session;

    // 此处定义静态变量,以在其他方法中获取到所有连接
    public static CopyOnWriteArraySet<WebSocket> wbSockets = new CopyOnWriteArraySet<WebSocket>();

    /**
     * 建立连接。
     * 建立连接时入参为session
     */
    @OnOpen
    public void onOpen(Session session) {
        this.session = session;
        wbSockets.add(this); //将此对象存入集合中以在之后广播用,如果要实现一对一订阅,则类型对应为Map。由于这里广播就可以了随意用Set
        logger.info("New session insert,sessionId is {}", session.getId());
    }

    /**
     * 关闭连接
     */
    @OnClose
    public void onClose() {
        wbSockets.remove(this);//将socket对象从集合中移除,以便广播时不发送次连接。如果不移除会报错(需要测试)
        logger.info("A session insert,sessionId is {}", session.getId());
    }

    /**
     * 接收消息
     */
    @OnMessage
    public void OnMessage(String message, Session session) {
        logger.info("A message insert,sessionId is {}", session.getId());
    }

    /**
     * 发送消息
     */
    public void sendMessage(Object message) throws Exception {
        //遍历客户端
        for (WebSocket webSocket : wbSockets) {
            logger.info("websocket info is {}", IJsonUtil.obj2Json(message));
            //服务器主动推送
            webSocket.session.getBasicRemote().sendObject(message);
        }
    }

    /**
     * 发送消息
     */
    public void sendMessage(String message) throws Exception {
        //遍历客户端
        for (WebSocket webSocket : wbSockets) {
            logger.info("websocket info is {}", message);
            //服务器主动推送
            webSocket.session.getBasicRemote().sendText(message);
        }
    }

    /**
     * 发送消息(用于发送给指定客户端消息)
     */
    public void sendMessage(String sessionId, String message) throws Exception {
        Session session = null;
        WebSocket tempWebSocket = null;
        logger.info("websocket info is {}", IJsonUtil.obj2Json(message));
        //遍历客户端
        for (WebSocket webSocket : wbSockets) {
            if (webSocket.session.getId().equals(sessionId)) {
                tempWebSocket = webSocket;
                session = webSocket.session;
                break;
            }
        }
        if (session != null) {
            tempWebSocket.session.getBasicRemote().sendText(message);
        } else {
            logger.info("没有找到你指定ID的会话:{}", sessionId);
        }
    }
}

三、redis消息接收处理类

@Component
public class RedisListenerHandle {

    private static final Logger logger = LoggerFactory.getLogger(RedisListenerHandle.class);

    @Autowired
    private WebSocket webSocket;

    /**
     * @param message
     * @description 接收到消息的业务处理
     * @create 2020-06-23 17:08:21
     */
    public void receiveMessage(String message) {
        logger.info("RedisListenerHandle get message info:{}", message);
        try {
            webSocket.sendMessage(message);
            logger.info("RedisListenerHandle push message success:{}", InetAddress.getLocalHost().getHostAddress());
        } catch (Exception e) {
            logger.info("RedisListenerHandle push message exception", e);
        }
    }
}

四、消息监听类

@Component
public class RedisListenerBean {
    /**
     * redis消息监听器容器
     * 可以添加多个监听不同话题的redis监听器,只需要把消息监听器和相应的消息订阅处理器绑定,该消息监听器
     * 通过反射技术调用消息订阅处理器的相关方法进行一些业务处理
     * @param connectionFactory
     * @return
     */
    @Bean
    public RedisMessageListenerContainer container(RedisConnectionFactory connectionFactory, MessageListenerAdapter listenerAdapter) {
        RedisMessageListenerContainer container = new RedisMessageListenerContainer();
        container.setConnectionFactory(connectionFactory);
        container.addMessageListener(listenerAdapter, new PatternTopic(CommonConstants.REDIS_CHANNEL));
        return container;
    }

    @Bean
    public MessageListenerAdapter messageListenerAdapter(RedisListenerHandle redisListenerHandle){
        return new MessageListenerAdapter(redisListenerHandle,"receiveMessage");
    }
}

再通过redisTemplate.convertAndSend()即可实现消息推送

总结:

1、为什么websocket不支持集群?

????????假如有A,B两台应用服务器,用户A与服务器A建立了连接,用户B与服务器B建立了连接,当要把某一消息发送给A,B两个用户时,如果是应用A在处理发送请求,那么用户A可以接收到消息,用户B则不能收到。如果是应用B在处理发送请求,那么用户B可以接收到消息,用户A则不能收到。这是因为websocket只能与一台服务器建立连接,而与哪台服务器建立连接又是通过nginx做的负载均衡。

2、redis实现websocket集群原理?

? ? ? ? 通过上面那个问题我们知道websocket只能与一台应用服务器建立连接,那么要实现消息推送给所有用户,无非就两种方案:

? ? ? ? a、服务器之间共享websocket连接

? ? ? ? b、让所有服务应用都发送消息

? ? ? ? 那么这里采用的则是第二种方案。先让A、B两台服务器都监听同一个topic,然后无论是A还是B服务处理发送请求,都把消息放到redis中,然后A、B监听到消息,都会处理消息,再通过websocket发送给各自连接的用户。

3、替代方案

? ? ? ? 上述所说的可以替换成rabbitmq,如果是并发较高的可以替换成rocketmq、kafka等吞吐量更好的消息推送服务。当然并发高的时候换成rocketmq或者kafka也需要考虑redis能不能承受的了。but,前端也有更好的推送技术,比如以mqtt协议实现推送的activeMq等等。

  网络协议 最新文章
使用Easyswoole 搭建简单的Websoket服务
常见的数据通信方式有哪些?
Openssl 1024bit RSA算法---公私钥获取和处
HTTPS协议的密钥交换流程
《小白WEB安全入门》03. 漏洞篇
HttpRunner4.x 安装与使用
2021-07-04
手写RPC学习笔记
K8S高可用版本部署
mySQL计算IP地址范围
上一篇文章      下一篇文章      查看所有文章
加:2021-11-23 12:43:57  更:2021-11-23 12:45:55 
 
开发: C++知识库 Java知识库 JavaScript Python PHP知识库 人工智能 区块链 大数据 移动开发 嵌入式 开发工具 数据结构与算法 开发测试 游戏开发 网络协议 系统运维
教程: HTML教程 CSS教程 JavaScript教程 Go语言教程 JQuery教程 VUE教程 VUE3教程 Bootstrap教程 SQL数据库教程 C语言教程 C++教程 Java教程 Python教程 Python3教程 C#教程
数码: 电脑 笔记本 显卡 显示器 固态硬盘 硬盘 耳机 手机 iphone vivo oppo 小米 华为 单反 装机 图拉丁

360图书馆 购物 三丰科技 阅读网 日历 万年历 2025年1日历 -2025/1/6 19:01:21-

图片自动播放器
↓图片自动播放器↓
TxT小说阅读器
↓语音阅读,小说下载,古典文学↓
一键清除垃圾
↓轻轻一点,清除系统垃圾↓
图片批量下载器
↓批量下载图片,美女图库↓
  网站联系: qq:121756557 email:121756557@qq.com  IT数码