From 77d745a7bbe9e36e958aceb9b51958e697d2f054 Mon Sep 17 00:00:00 2001
From: tangzy <tangzy123456>
Date: Tue, 23 Mar 2021 13:56:37 +0800
Subject: [PATCH] 1.webscoket
---
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/ChannelSupervise.java | 29 ++++
/dev/null | 75 ----------
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketServer.java | 58 ++++++++
blade-service/blade-jfpts/pom.xml | 9 -
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketHandler.java | 171 ++++++++++++++++++++++++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/index.html | 54 +++++++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java | 3
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/impl/SuserServiceImpl.java | 6
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.java | 2
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/ISuserService.java | 2
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.xml | 4
11 files changed, 329 insertions(+), 84 deletions(-)
diff --git a/blade-service/blade-jfpts/pom.xml b/blade-service/blade-jfpts/pom.xml
index 6cea5f3..03c30bf 100644
--- a/blade-service/blade-jfpts/pom.xml
+++ b/blade-service/blade-jfpts/pom.xml
@@ -59,13 +59,8 @@
<artifactId>blade-starter-oss-qiniu</artifactId>
</dependency>
<dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-websocket</artifactId>
- <version>2.0.4.RELEASE</version>
- </dependency>
- <dependency>
- <groupId>org.springframework.boot</groupId>
- <artifactId>spring-boot-starter-thymeleaf</artifactId>
+ <groupId>io.netty</groupId>
+ <artifactId>netty-all</artifactId>
</dependency>
</dependencies>
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java
index 1a6344c..93bd4ea 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java
@@ -20,6 +20,7 @@
import org.springblade.core.launch.BladeApplication;
import org.springblade.jfpt.nettyServer.Server;
import org.springblade.jfpt.nettyTcpServer.TcpServer;
+import org.springblade.jfpt.webscoket.WebSocketServer;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@@ -41,8 +42,8 @@
public void run(String... args) throws Exception {
Server server=new Server(8088);
TcpServer tcpServer = new TcpServer(8099);
+ WebSocketServer webSocketServer= new WebSocketServer(9034);
//UdpServer udpServer=new UdpServer(8099);
}
-
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.java
index 22a3160..42c91b0 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.java
@@ -44,6 +44,8 @@
List<SuserVOs> selectName();
String selectNames(String ud);
void updateOne(Integer id, String ynumber, String yname, String phone, String address, String type, String dj, String information, String jd, String wd, String createtime, String addvcd);
+ void updateUser(String online_status,String account);
+
void s(String ynumber, String yname, String phone, String address, String type, String dj, String information, String jd, String wd, String createtime, String addvcd);
/**
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.xml b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.xml
index 11f0df4..f46ac29 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.xml
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/mapper/SuserMapper.xml
@@ -118,4 +118,8 @@
select count(*) from sys_suser where is_delete=0
</select>
+ <update id="updateUser">
+ update jfpth.blade_user SET online_status=#{online_status} where account=#{account}
+ </update>
+
</mapper>
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/ISuserService.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/ISuserService.java
index 5b88878..c4d783a 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/ISuserService.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/ISuserService.java
@@ -45,7 +45,7 @@
String selectNames(String ud);
void updateOne(Integer id, String ynumber, String yname, String phone, String address, String type, String dj, String information, String jd, String wd, String createtime, String addvcd);
void s(String ynumber, String yname, String phone, String address, String type, String dj, String information, String jd, String wd, String createtime, String addvcd);
-
+ void updateUser(String online_status,String account);
/**
* 查询业主总数
* @return
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/impl/SuserServiceImpl.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/impl/SuserServiceImpl.java
index f1ca418..8249edb 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/impl/SuserServiceImpl.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/suser/service/impl/SuserServiceImpl.java
@@ -71,6 +71,12 @@
baseMapper.s(ynumber, yname, phone, address, type, dj, information, jd, wd,createtime,addvcd);
}
+ @Override
+ public void updateUser(String online_status, String account) {
+ baseMapper.updateUser(online_status, account);
+ }
+
+
/**
* 查询业主总数
* @return
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/ChannelSupervise.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/ChannelSupervise.java
new file mode 100644
index 0000000..b0c581d
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/ChannelSupervise.java
@@ -0,0 +1,29 @@
+package org.springblade.jfpt.webscoket;
+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){
+ GlobalGroup.add(channel);
+ ChannelMap.put(channel.id().asShortText(),channel.id());
+ }
+ public static void removeChannel(Channel channel){
+ GlobalGroup.remove(channel);
+ ChannelMap.remove(channel.id().asShortText());
+ }
+ public static Channel findChannel(String id){
+ return GlobalGroup.find(ChannelMap.get(id));
+ }
+ public static void send2All(TextWebSocketFrame tws){
+ GlobalGroup.writeAndFlush(tws);
+ }
+
+}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocket.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocket.java
deleted file mode 100644
index 48a4010..0000000
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocket.java
+++ /dev/null
@@ -1,75 +0,0 @@
-package org.springblade.jfpt.webscoket;
-
-import org.springframework.stereotype.Component;
-
-import javax.websocket.*;
-import javax.websocket.server.PathParam;
-import javax.websocket.server.ServerEndpoint;
-import java.io.IOException;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-@ServerEndpoint("/websocket/{pageCode}")
-@Component
-public class WebSocket {
-
- private static final String loggerName=WebSocket.class.getName();
- //concurrent包的线程安全Set,用来存放每个客户端对应的MyWebSocket对象。若要实现服务端与单一客户端通信的话,可以使用Map来存放,其中Key可以为用户标识
- public static Map<String, List<Session>> electricSocketMap = new ConcurrentHashMap<String, List<Session>>();
-
- /**
- * 连接建立成功调用的方法
- *
- * @param session 可选的参数。session为与某个客户端的连接会话,需要通过它来给客户端发送数据
- */
- @OnOpen
- public void onOpen(@PathParam("pageCode") String pageCode, Session session) {
- List<Session> sessions = electricSocketMap.get(pageCode);
- if(null==sessions){
- List<Session> sessionList = new ArrayList<>();
- sessionList.add(session);
- electricSocketMap.put(pageCode,sessionList);
- }else{
- sessions.add(session);
- }
- }
-
- /**
- * 连接关闭调用的方法
- */
- @OnClose
- public void onClose(@PathParam("pageCode") String pageCode,Session session) {
- if (electricSocketMap.containsKey(pageCode)){
- electricSocketMap.get(pageCode).remove(session);
- }
- }
-
- /**
- * 收到客户端消息后调用的方法
- *
- * @param message 客户端发送过来的消息
- * @param session 可选的参数
- */
- @OnMessage
- public void onMessage(String message, Session session) {
- System.out.println("websocket received message:"+message);
- try {
- session.getBasicRemote().sendText("这是推送测试数据!您刚发送的消息是:"+message);
- } catch (IOException e) {
- e.printStackTrace();
- }
- }
-
- /**
- * 发生错误时调用
- *
- * @param session
- * @param error
- */
- @OnError
- public void onError(Session session, Throwable error) {
- System.out.println("发生错误");;
- }
-}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketHandler.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketHandler.java
new file mode 100644
index 0000000..2cbdd7f
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketHandler.java
@@ -0,0 +1,171 @@
+package org.springblade.jfpt.webscoket;
+
+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.CloseWebSocketFrame;
+import io.netty.handler.codec.http.websocketx.PingWebSocketFrame;
+import io.netty.handler.codec.http.websocketx.PongWebSocketFrame;
+import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;
+import io.netty.handler.codec.http.websocketx.WebSocketFrame;
+import io.netty.handler.codec.http.websocketx.WebSocketServerHandshaker;
+import io.netty.handler.codec.http.websocketx.WebSocketServerHandshakerFactory;
+import io.netty.util.CharsetUtil;
+import java.util.Date;
+
+import org.springblade.jfpt.suser.service.ISuserService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Component;
+
+import javax.annotation.PostConstruct;
+
+@Component
+public class WebSocketHandler extends SimpleChannelInboundHandler<Object> {
+
+ private WebSocketServerHandshaker handshaker;
+
+ private String on=null;
+
+ @Autowired
+ private ISuserService suserService;
+
+ private static WebSocketHandler webSocketHandler;
+
+ @PostConstruct
+ public void init() {
+ webSocketHandler = this;
+ }
+
+ @Override
+ protected void channelRead0(ChannelHandlerContext ctx, Object msg) throws Exception {
+ if (msg instanceof FullHttpRequest) {
+ //以http请求形式接入,但是走的是websocket
+ handleHttpRequest(ctx, (FullHttpRequest) msg);
+ } else if (msg instanceof WebSocketFrame) {
+ //处理websocket客户端的消息
+ handlerWebSocketFrame(ctx, (WebSocketFrame) msg);
+ }
+ }
+
+ @Override
+ public void channelActive(ChannelHandlerContext ctx) throws Exception {
+ //添加连接
+ System.out.println("客户端加入连接:" + ctx.channel());
+ //ChannelSupervise.addChannel(ctx.channel());
+ }
+
+ @Override
+ public void channelInactive(ChannelHandlerContext ctx) throws Exception {
+ //断开连接
+ System.out.println("客户端断开连接:" + ctx.channel());
+ ChannelSupervise.removeChannel(ctx.channel());
+ }
+
+ @Override
+ public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
+ ctx.flush();
+ }
+
+ private void handlerWebSocketFrame(ChannelHandlerContext ctx, WebSocketFrame frame) throws NumberFormatException, Exception {
+ // 判断是否关闭链路的指令
+ 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)) {
+ System.out.println("本例程仅支持文本消息,不支持二进制消息");
+ throw new UnsupportedOperationException(String.format(
+ "%s frame types not supported", frame.getClass().getName()));
+ }
+ // 返回应答消息
+ String request = ((TextWebSocketFrame) frame).text();
+ String substring1 = request.substring(0, 2);
+ //用户在线状态
+ if(substring1.equals("登录")){
+ String online="1";
+ String substring = request.substring(0, request.length());
+ System.out.println(substring);
+ String[] split = substring.split(",");//以逗号分割
+ this.on=split[1];
+ webSocketHandler.suserService.updateUser(online,this.on);
+ }
+ //用户离线状态
+ else if (substring1.equals("退出")){
+ String online="0";
+ String substring = request.substring(0, request.length());
+ System.out.println(substring);
+ String[] split = substring.split(",");//以逗号分割
+ this.on=split[1];
+ webSocketHandler.suserService.updateUser(online,this.on);
+ webSocketHandler.channelInactive(ctx);
+ }
+ TextWebSocketFrame tws = new TextWebSocketFrame(new Date().toString()
+ + ctx.channel().id() + ":" + ctx.channel());
+ // 返回【谁发的发给谁】
+
+ ctx.channel().writeAndFlush(tws);
+ // 群发
+ // ChannelSupervise.send2All(tws);
+ }
+
+ /**
+ * 唯一的一次http请求,用于创建websocket
+ *
+ * @throws InterruptedException
+ */
+ private void handleHttpRequest(final ChannelHandlerContext ctx,
+ FullHttpRequest req) throws InterruptedException {
+ //要求Upgrade为websocket,过滤掉get/Post
+ if (!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:9034/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);
+// }
+ }
+ @Override
+ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
+ cause.printStackTrace();
+ ctx.close();
+ }
+
+}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketServer.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketServer.java
new file mode 100644
index 0000000..9dada6a
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/WebSocketServer.java
@@ -0,0 +1,58 @@
+package org.springblade.jfpt.webscoket;
+import io.netty.bootstrap.ServerBootstrap;
+import io.netty.channel.ChannelFuture;
+import io.netty.channel.ChannelInitializer;
+import io.netty.channel.ChannelOption;
+import io.netty.channel.EventLoopGroup;
+import io.netty.channel.nio.NioEventLoopGroup;
+import io.netty.channel.socket.SocketChannel;
+import io.netty.channel.socket.nio.NioServerSocketChannel;
+import io.netty.handler.codec.http.HttpObjectAggregator;
+import io.netty.handler.codec.http.HttpServerCodec;
+import io.netty.handler.logging.LogLevel;
+import io.netty.handler.logging.LoggingHandler;
+import io.netty.handler.stream.ChunkedWriteHandler;
+public class WebSocketServer {
+ private int port = 9034;
+
+ public WebSocketServer(int port) {
+ bind(port);
+ }
+
+ public void bind(int port) {
+ Thread thread = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ EventLoopGroup bossGroup = new NioEventLoopGroup();
+ EventLoopGroup workerGroup = new NioEventLoopGroup();
+ try {
+ ServerBootstrap serverBootstrap = new ServerBootstrap();
+ serverBootstrap.group(bossGroup, workerGroup).channel(NioServerSocketChannel.class)
+ .handler(new LoggingHandler(LogLevel.INFO))
+ //保持连接
+ .childOption(ChannelOption.SO_KEEPALIVE, true)
+ .childHandler(new ChannelInitializer<SocketChannel>() {
+ @Override
+ protected 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 WebSocketHandler());//自定义的业务handler
+ }
+ });
+
+ ChannelFuture channelFuture = serverBootstrap.bind(port).sync();
+ System.out.println("WebSocketServer启动成功");
+ channelFuture.channel().closeFuture().sync();
+ } catch (Exception e) {
+ e.printStackTrace();
+ } finally {
+ bossGroup.shutdownGracefully();
+ workerGroup.shutdownGracefully();
+ }
+ }
+ });
+ thread.start();
+ }
+}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/index.html b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/index.html
new file mode 100644
index 0000000..73735f9
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/webscoket/index.html
@@ -0,0 +1,54 @@
+<html>
+<head>
+ <meta http-equiv="Content-Type" content="text/html; charset = utf-8"/>
+ <title>WebSocket客户端</title>
+ <script type="text/javascript">
+ var socket;
+ if(!window.WebSocket){
+ window.WebSocket = window.MozWebSocket;
+ }
+
+ if(window.WebSocket){
+ socket = new WebSocket("ws://localhost:9034/websocket");
+ socket.onmessage = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value += event.data + "\r\n";
+ };
+
+ socket.onopen = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value = "你当前的浏览器支持WebSocket,请进行后续操作\r\n";
+ };
+
+ socket.onclose = function(event){
+ var ta = document.getElementById('responseContent');
+ ta.value = "";
+ ta.value = "WebSocket连接已经关闭\r\n";
+ };
+ }else{
+ alert("您的浏览器不支持WebSocket");
+ }
+
+ function send(message){
+ if(!window.WebSocket){
+ return;
+ }
+ if(socket.readyState == WebSocket.OPEN){
+ socket.send(message);
+ }else{
+ alert("WebSocket连接没有建立成功!!");
+ }
+ }
+ </script>
+</head>
+<body>
+<form onSubmit="return false;">
+ <input type = "text" name = "message" value = ""/>
+ <br/><br/>
+ <input type = "button" value = "发送WebSocket请求消息" onClick = "send(this.form.message.value)"/>
+ <hr color="red"/>
+ <h2>客户端接收到服务端返回的应答消息</h2>
+ <textarea id = "responseContent" style = "width:1024px; height:300px"></textarea>
+</form>
+</body>
+</html>
--
Gitblit v1.9.3