利用netty实现websocket ;redis的订阅发布websocket相结合

2024-08-22 05:44

本文主要是介绍利用netty实现websocket ;redis的订阅发布websocket相结合,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

由于Http协议是无状态的,每一次请求只能响应一次,下次请求需要重新连接。
如果客户端请求一个服务端资源,需要实时监服务端执行状态(比如导出大数据量时需要前端监控导出状态),这个时候不断请求连接浪费资源。可以通过WebSocket建立一个长连接,实现客户端与服务端双向交流。

websocket服务器

public class NioWebSocketServer {private final Logger logger=Logger.getLogger(this.getClass());private void init(){logger.info("正在启动websocket服务器");NioEventLoopGroup boss=new NioEventLoopGroup();NioEventLoopGroup work=new NioEventLoopGroup();try {ServerBootstrap bootstrap=new ServerBootstrap();bootstrap.group(boss,work);bootstrap.channel(NioServerSocketChannel.class);bootstrap.childHandler(new NioWebSocketChannelInitializer());Channel channel = bootstrap.bind(8083).sync().channel();logger.info("webSocket服务器启动成功:"+channel);channel.closeFuture().sync();} catch (InterruptedException e) {e.printStackTrace();logger.info("运行出错:"+e);}finally {boss.shutdownGracefully();work.shutdownGracefully();logger.info("websocket服务器已关闭");}}public static void main(String[] args) {new NioWebSocketServer().init();}
}

ChannelInitializer

import io.netty.channel.ChannelInitializer;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.codec.http.HttpObjectAggregator;
import io.netty.handler.codec.http.HttpServerCodec;
import io.netty.handler.logging.LoggingHandler;
import io.netty.handler.stream.ChunkedWriteHandler;public class NioWebSocketChannelInitializer extends ChannelInitializer<SocketChannel> {@Overrideprotected void initChannel(SocketChannel ch) {ch.pipeline().addLast("logging",new LoggingHandler("DEBUG"));//设置log监听器,并且日志级别为debug,方便观察运行流程ch.pipeline().addLast("http-codec",new HttpServerCodec());//设置解码器ch.pipeline().addLast("aggregator",new HttpObjectAggregator(65536));//聚合器,使用websocket会用到ch.pipeline().addLast("http-chunked",new ChunkedWriteHandler());//用于大数据的分区传输ch.pipeline().addLast("handler",new NioWebSocketHandler());//自定义的业务handler}
}

Hander

import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.*;
import io.netty.handler.codec.http.DefaultFullHttpResponse;
import io.netty.handler.codec.http.FullHttpRequest;
import io.netty.handler.codec.http.HttpResponseStatus;
import io.netty.handler.codec.http.HttpVersion;
import io.netty.handler.codec.http.websocketx.*;
import io.netty.util.CharsetUtil;
import org.apache.log4j.Logger;
import org.wisdom.netty.global.ChannelSupervise;import java.util.Date;import static io.netty.handler.codec.http.HttpUtil.isKeepAlive;public class NioWebSocketHandler extends SimpleChannelInboundHandler<Object> {private final Logger logger=Logger.getLogger(this.getClass());private WebSocketServerHandshaker handshaker;/**重写channelRead0方法,处理接收到的消息* @param ctx* @param msg* @throws Exception*/@Overrideprotected void channelRead0(ChannelHandlerContext ctx, Object msg) throws Exception {logger.debug("收到消息:"+msg);if (msg instanceof FullHttpRequest){//以http请求形式接入,但是走的是websockethandleHttpRequest(ctx, (FullHttpRequest) msg);}else if (msg instanceof  WebSocketFrame){//处理websocket客户端的消息handlerWebSocketFrame(ctx, (WebSocketFrame) msg);}}@Overridepublic void channelActive(ChannelHandlerContext ctx) throws Exception {//添加连接logger.debug("客户端加入连接:"+ctx.channel());ChannelSupervise.addChannel(ctx.channel());}@Overridepublic void channelInactive(ChannelHandlerContext ctx) throws Exception {//断开连接logger.debug("客户端断开连接:"+ctx.channel());ChannelSupervise.removeChannel(ctx.channel());}@Overridepublic void channelReadComplete(ChannelHandlerContext ctx) throws Exception {ctx.flush();}private void handlerWebSocketFrame(ChannelHandlerContext ctx, WebSocketFrame frame){// 判断是否关闭链路的指令if (frame instanceof CloseWebSocketFrame) {handshaker.close(ctx.channel(), (CloseWebSocketFrame) frame.retain());return;}// 判断是否ping消息if (frame instanceof PingWebSocketFrame) {ctx.channel().write(new PongWebSocketFrame(frame.content().retain()));return;}// 本例程仅支持文本消息,不支持二进制消息if (!(frame instanceof TextWebSocketFrame)) {logger.debug("本例程仅支持文本消息,不支持二进制消息");throw new UnsupportedOperationException(String.format("%s frame types not supported", frame.getClass().getName()));}// 返回应答消息String request = ((TextWebSocketFrame) frame).text();logger.debug("服务端收到:" + request);TextWebSocketFrame tws = new TextWebSocketFrame(new Date().toString()+ ctx.channel().id() + ":" + request);// 群发ChannelSupervise.send2All(tws);// 返回【谁发的发给谁】// ctx.channel().writeAndFlush(tws);}/*** 唯一的一次http请求,用于创建websocket* */private void handleHttpRequest(ChannelHandlerContext ctx,FullHttpRequest req) {//要求Upgrade为websocket,过滤掉get/Postif (!req.decoderResult().isSuccess()|| (!"websocket".equals(req.headers().get("Upgrade")))) {//若不是websocket方式,则创建BAD_REQUEST的req,返回给客户端sendHttpResponse(ctx, req, new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.BAD_REQUEST));return;}WebSocketServerHandshakerFactory wsFactory = new WebSocketServerHandshakerFactory("ws://localhost:8083/websocket", null, false);handshaker = wsFactory.newHandshaker(req);if (handshaker == null) {WebSocketServerHandshakerFactory.sendUnsupportedVersionResponse(ctx.channel());} else {handshaker.handshake(ctx.channel(), req);}}/*** 拒绝不合法的请求,并返回错误信息* */private static void sendHttpResponse(ChannelHandlerContext ctx,FullHttpRequest req, DefaultFullHttpResponse res) {// 返回应答给客户端if (res.status().code() != 200) {ByteBuf buf = Unpooled.copiedBuffer(res.status().toString(),CharsetUtil.UTF_8);res.content().writeBytes(buf);buf.release();}ChannelFuture f = ctx.channel().writeAndFlush(res);// 如果是非Keep-Alive,关闭连接if (!isKeepAlive(req) || res.status().code() != 200) {f.addListener(ChannelFutureListener.CLOSE);}}
}

存储信息

import io.netty.channel.Channel;
import io.netty.channel.ChannelId;
import io.netty.channel.group.ChannelGroup;
import io.netty.channel.group.DefaultChannelGroup;
import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
import io.netty.util.concurrent.GlobalEventExecutor;import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;public class ChannelSupervise {private   static ChannelGroup GlobalGroup=new DefaultChannelGroup(GlobalEventExecutor.INSTANCE);private  static ConcurrentMap<String, ChannelId> ChannelMap=new ConcurrentHashMap();public  static void addChannel(Channel channel){/**channel(通道是socket的连接信息;) = [id: 0x2a5ba781, L:/127.0.0.1:8083 - R:/127.0.0.1:60663]==ChannelMap=={28176dd1=28176dd1, 2a5ba781=2a5ba781}*/GlobalGroup.add(channel);ChannelMap.put(channel.id().asShortText(),channel.id());System.out.println("==channel=="+channel);System.out.println("==ChannelMap=="+ChannelMap.toString());}public static void removeChannel(Channel channel){GlobalGroup.remove(channel);ChannelMap.remove(channel.id().asShortText());System.out.println("==removeChannel=="+channel);System.out.println("==removeChannelChannelMap=="+ChannelMap.toString());}public static  Channel findChannel(String id){return GlobalGroup.find(ChannelMap.get(id));}/*** 根据channel id 进行 群发通知* @param tws*/public static void send2All(TextWebSocketFrame tws){GlobalGroup.writeAndFlush(tws);}
}

html

<!-- index.html --><!DOCTYPE html>
<html>
<head><meta charset="UTF-8"><title>WebSocket Test</title>
</head>
<body><h1>WebSocket Test</h1><div><input type="text" id="message" placeholder="Message"><button onclick="send()">Send</button></div><div id="output"></div><script>var socket = new WebSocket("ws://localhost:8083/");socket.onopen = function(event) {console.log("WebSocket opened: " + event);};socket.onmessage = function(event) {console.log("WebSocket message received: " + event.data);var output = document.getElementById("output");output.innerHTML += "<p>" + event.data + "</p>";};socket.onclose = function(event) {console.log("WebSocket closed: " + event);};function send() {var message = document.getElementById("message").value;socket.send(message);}</script>
</body>
</html>

redis 订阅

   public void sendAlarmFaultMessage(String message) {String newMessge= null;try {newMessge = new String(message.getBytes(RedisKeyConstant.UTF8), RedisKeyConstant.UTF8);} catch (UnsupportedEncodingException e) {e.printStackTrace();}//redisTemplate.convertAndSend(RedisKeyConstant.REDIS_CHANNEL, newMessge);RedissonClient redissonClient = SpringUtil.getBean(RedissonClient.class);RTopic topic = redissonClient.getTopic(RedisKeyConstant.REDIS_CHANNEL_FAULT);topic.publish(newMessge);redisTemplate.opsForList().rightPush(RedisKeyConstant.REDIS_MESSAGE_FAULT, newMessge);}

redis发布

@Beanpublic RTopic rFaultTopic(RedissonClient redissonClient) {RTopic rTopic = redissonClient.getTopic(RedisKeyConstant.REDIS_CHANNEL_FAULT);try{if(rTopic != null){rTopic.addListener(String.class, (channel, message) -> {if (channel.toString().contains(RedisKeyConstant.REDIS_CHANNEL_FAULT)) {RedisUtil.lpop(RedisKeyConstant.REDIS_MESSAGE_FAULT);log.info("Channel is already fault "+message);AlarmWebsocketService alarmWebsocketService = SpringUtil.getBean(AlarmWebsocketService.class);alarmWebsocketService.sendAllMessage(message);}});}}catch (Exception e){log.info("Error sending alarm "+ ExceptionUtils.getStackTrace(e));}return rTopic;}

这篇关于利用netty实现websocket ;redis的订阅发布websocket相结合的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1095403

相关文章

基于C++的UDP网络通信系统设计与实现详解

《基于C++的UDP网络通信系统设计与实现详解》在网络编程领域,UDP作为一种无连接的传输层协议,以其高效、低延迟的特性在实时性要求高的应用场景中占据重要地位,下面我们就来看看如何从零开始构建一个完整... 目录前言一、UDP服务器UdpServer.hpp1.1 基本框架设计1.2 初始化函数Init详解

Java中Map的五种遍历方式实现与对比

《Java中Map的五种遍历方式实现与对比》其实Map遍历藏着多种玩法,有的优雅简洁,有的性能拉满,今天咱们盘一盘这些进阶偏基础的遍历方式,告别重复又臃肿的代码,感兴趣的小伙伴可以了解下... 目录一、先搞懂:Map遍历的核心目标二、几种遍历方式的对比1. 传统EntrySet遍历(最通用)2. Lambd

springboot+redis实现订单过期(超时取消)功能的方法详解

《springboot+redis实现订单过期(超时取消)功能的方法详解》在SpringBoot中使用Redis实现订单过期(超时取消)功能,有多种成熟方案,本文为大家整理了几个详细方法,文中的示例代... 目录一、Redis键过期回调方案(推荐)1. 配置Redis监听器2. 监听键过期事件3. Redi

SpringBoot全局异常拦截与自定义错误页面实现过程解读

《SpringBoot全局异常拦截与自定义错误页面实现过程解读》本文介绍了SpringBoot中全局异常拦截与自定义错误页面的实现方法,包括异常的分类、SpringBoot默认异常处理机制、全局异常拦... 目录一、引言二、Spring Boot异常处理基础2.1 异常的分类2.2 Spring Boot默

基于SpringBoot实现分布式锁的三种方法

《基于SpringBoot实现分布式锁的三种方法》这篇文章主要为大家详细介绍了基于SpringBoot实现分布式锁的三种方法,文中的示例代码讲解详细,感兴趣的小伙伴可以跟随小编一起学习一下... 目录一、基于Redis原生命令实现分布式锁1. 基础版Redis分布式锁2. 可重入锁实现二、使用Redisso

SpringBoo WebFlux+MongoDB实现非阻塞API过程

《SpringBooWebFlux+MongoDB实现非阻塞API过程》本文介绍了如何使用SpringBootWebFlux和MongoDB实现非阻塞API,通过响应式编程提高系统的吞吐量和响应性能... 目录一、引言二、响应式编程基础2.1 响应式编程概念2.2 响应式编程的优势2.3 响应式编程相关技术

C#实现将XML数据自动化地写入Excel文件

《C#实现将XML数据自动化地写入Excel文件》在现代企业级应用中,数据处理与报表生成是核心环节,本文将深入探讨如何利用C#和一款优秀的库,将XML数据自动化地写入Excel文件,有需要的小伙伴可以... 目录理解XML数据结构与Excel的对应关系引入高效工具:使用Spire.XLS for .NETC

Nginx更新SSL证书的实现步骤

《Nginx更新SSL证书的实现步骤》本文主要介绍了Nginx更新SSL证书的实现步骤,包括下载新证书、备份旧证书、配置新证书、验证配置及遇到问题时的解决方法,感兴趣的了解一下... 目录1 下载最新的SSL证书文件2 备份旧的SSL证书文件3 配置新证书4 验证配置5 遇到的http://www.cppc

Nginx之https证书配置实现

《Nginx之https证书配置实现》本文主要介绍了Nginx之https证书配置的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来一起... 目录背景介绍为什么不能部署在 IIS 或 NAT 设备上?具体实现证书获取nginx配置扩展结果验证

SpringBoot整合 Quartz实现定时推送实战指南

《SpringBoot整合Quartz实现定时推送实战指南》文章介绍了SpringBoot中使用Quartz动态定时任务和任务持久化实现多条不确定结束时间并提前N分钟推送的方案,本文结合实例代码给大... 目录前言一、Quartz 是什么?1、核心定位:解决什么问题?2、Quartz 核心组件二、使用步骤1