From 5549170855c433b45957d2142aefc54ba66ed3e0 Mon Sep 17 00:00:00 2001
From: Administrator <admin>
Date: Tue, 27 Apr 2021 20:57:02 +0800
Subject: [PATCH] 1.设备详情接口修改,新增deptName字段返回 2.system 模块新增tcp,upd测试服务
---
blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/Server.java | 86 ++++++++++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java | 4
blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/handle/UdpServerHandler.java | 51 ++++++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/mapper/EquipmentMapper.xml | 15 +
blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/server/UdpServer.java | 129 ++++++++++++++++
blade-service/blade-system/src/main/java/org/springblade/system/SystemApplication.java | 12 +
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/vo/EquipmentVOS.java | 4
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/controller/EquipmentController.java | 24 ++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/entity/Equipment.java | 4
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyTcpServer/TcpServerHandler.java | 4
blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/ServerHandler.java | 104 +++++++++++++
11 files changed, 422 insertions(+), 15 deletions(-)
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 fe2271e..05dd031 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.nettyUdpServer.server.UdpServer;
import org.springblade.jfpt.webscoket.WebSocketServer;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@@ -44,7 +45,8 @@
Server server=new Server(8088);
TcpServer tcpServer = new TcpServer(8099);
WebSocketServer webSocketServer= new WebSocketServer(9034);
- //UdpServer udpServer=new UdpServer(8099);
+// UdpServer udpServer = new UdpServer(9534);
+// Server serverTwo = new Server(9534);
}
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/controller/EquipmentController.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/controller/EquipmentController.java
index 7bbc177..8332a9e 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/controller/EquipmentController.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/controller/EquipmentController.java
@@ -80,8 +80,25 @@
response.setHeader("Access-Control-Allow-Methods", "POST, GET, OPTIONS, DELETE");
response.setHeader("Access-Control-Allow-Credentials", "true");
/*Equipment detail = equipmentService.getOne(Condition.getQueryWrapper(equipment));*/
-
- return R.data(equipmentService.selectInfo(equipment));
+ EquipmentVOS equipmentVOS = equipmentService.selectInfo(equipment);
+ //获取部门数据
+ List<DeptVo> deptVos = deptFeignApi.selDeptList();
+ if (null!=equipmentVOS.getDeptId()){
+ String[] split =equipmentVOS.getDeptId().split(",");
+ List<String> list = Arrays.asList(split);
+ StringBuffer deptNameBuiffer = new StringBuffer();
+ //数据匹配封装
+ for (String deptId:list) {
+ for (DeptVo deptVo:deptVos) {
+ if (deptId.equals(deptVo.getId().toString())){
+ deptNameBuiffer.append(deptVo.getDeptName()).append("/");
+ }
+ }
+ }
+ //封装部门名称数据
+ equipmentVOS.setDeptName(deptNameBuiffer.substring(0,deptNameBuiffer.length()-1));
+ }
+ return R.data(equipmentVOS);
}
/**
@@ -137,10 +154,7 @@
}else{
pages.get(j).setOnlineStatus("0");
}
-
}
-
-
return R.data(pages);
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/entity/Equipment.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/entity/Equipment.java
index c35ab2f..64d63f5 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/entity/Equipment.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/entity/Equipment.java
@@ -136,5 +136,9 @@
private String twoContacts;
@TableField("ThreeContacts")
private String ThreeContacts;
+ /**
+ * 设备厂商名称
+ */
+ private String manufacturers;
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/mapper/EquipmentMapper.xml b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/mapper/EquipmentMapper.xml
index 442e7b7..587f6ab 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/mapper/EquipmentMapper.xml
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/mapper/EquipmentMapper.xml
@@ -30,6 +30,7 @@
<result column="state" property="state"/>
<result column="reason" property="reason"/>
<result column="jtype" property="jtype"/>
+ <result column="manufacturers" property="manufacturers"/>
</resultMap>
<resultMap id="treeNodeResultMap" type="org.springblade.core.tool.node.TreeNode">
@@ -142,7 +143,7 @@
equipment.oneContacts,equipment.onePhone,equipment.oneId,
equipment.twoContacts,equipment.twoPhone,equipment.twoId,
equipment.ThreeContacts,equipment.ThreePhone,equipment.threeId,
- equipment.deptId,
+ equipment.deptId,equipment.manufacturers,
equipment.addvcd,equipment.pay,equipment.expireTime from sys_equipment equipment
left join sys_catalog catalog on catalog.eqid = equipment.deviceNumber
LEFT JOIN (
@@ -181,6 +182,9 @@
<if test="equipment.city!=null">
and city like concat('%',#{equipment.city},'%')
</if>
+ <if test="equipment.manufacturers!=null and equipment.manufacturers!=''">
+ and manufacturers = #{equipment.manufacturers}
+ </if>
<if test="equipment.devicestate!=null">
and devicestate =#{equipment.devicestate}
@@ -210,7 +214,7 @@
<!--新增-->
<insert id="s" parameterType="org.springblade.jfpt.equipment.entity.Equipment">
- INSERT INTO sys_equipment (deviceName,deviceNumber,deviceType,ownership,ownerId,devicestate,stime,coordinate,parent_id,jd,wd,addvcd,street,dtype) VALUES (#{deviceName},#{deviceNumber},#{deviceType},#{ownership},#{yname},#{devicestate},#{stime},POINT(#{jd},#{wd}),#{parentId},#{jd},#{wd},#{addvcd},#{streeName},#{dtype})
+ INSERT INTO sys_equipment (deviceName,deviceNumber,deviceType,ownership,ownerId,devicestate,stime,coordinate,parent_id,jd,wd,addvcd,street,dtype,manufacturers) VALUES (#{deviceName},#{deviceNumber},#{deviceType},#{ownership},#{yname},#{devicestate},#{stime},POINT(#{jd},#{wd}),#{parentId},#{jd},#{wd},#{addvcd},#{streeName},#{dtype},#{manufacturers})
</insert>
<!--新增-->
@@ -348,8 +352,9 @@
<update id="updateOne" parameterType="org.springblade.jfpt.equipment.entity.Equipment">
update sys_equipment SET deviceName=#{deviceName},deviceNumber=#{deviceNumber},deviceType=#{deviceType},
- ownership=#{ownership},ownerId=#{yname},devicestate=#{devicestate},stime=#{stime},
- coordinate=POINT(#{jd},#{wd}),parent_id=#{parentId},jd=#{jd},wd=#{wd},addvcd=#{addvnm},street=#{streeName},dtype=#{dtype} where id=#{id}
+ ownership=#{ownership},ownerId=#{yname},devicestate=#{devicestate},stime=#{stime},coordinate=POINT(#{jd},#{wd}),
+ parent_id=#{parentId},jd=#{jd},wd=#{wd},addvcd=#{addvnm},street=#{streeName},dtype=#{dtype},manufacturers=#{manufacturers}
+ where id=#{id}
</update>
<update id="updataType">
@@ -365,7 +370,7 @@
SELECT
dept.deptId,dept.id,dept.deviceName,dept.deviceNumber,dept.deviceType,dept.ownership,
dept.ownerId,dept.devicestate,dept.stime,dept.district,dept.province,dept.city,
- ST_ASTEXT(dept.coordinate) as coordinate,dept.dtype,dept.expireTime,dept.heartbeat, dept.state,dept.reason,
+ ST_ASTEXT(dept.coordinate) as coordinate,dept.dtype,dept.manufacturers,dept.expireTime,dept.heartbeat, dept.state,dept.reason,
dept.parent_id,dept.jd,dept.wd,dept.street,u.yname,d.addvnm,dept.addvcd,dept.serialNumber,dept.channelNumber,
dept.oneContacts,dept.onePhone,dept.oneId,
dept.twoContacts,dept.twoPhone,dept.twoId,
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/vo/EquipmentVOS.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/vo/EquipmentVOS.java
index 2c6eefd..8519d23 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/vo/EquipmentVOS.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/equipment/vo/EquipmentVOS.java
@@ -77,5 +77,9 @@
private String bType;
private String addvnm;
private String streeName;
+ /**
+ * 组织机构名称
+ */
+ private String deptName;
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyTcpServer/TcpServerHandler.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyTcpServer/TcpServerHandler.java
index 23ad48f..a36ea5e 100644
--- a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyTcpServer/TcpServerHandler.java
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyTcpServer/TcpServerHandler.java
@@ -74,14 +74,12 @@
buf.readBytes(req);
//手动释放内存
ReferenceCountUtil.release(info);
-
String body = new String(req, "UTF-8");
String content = body;
System.out.println("接收客户端数据:" + body);
- System.out.println("content.length() = " + content.length());
//获取字符串的长度
int length = content.length();
- //当字符串长度大于200时,采集体温数据,小于则采集报警数据
+ //当字符串长度大于200时,采集体温数据
//只有数据长度大于300时,采集体温数据 200-300 之间为心跳数据,当前不采集
if (length>250){
//插入数据
diff --git a/blade-service/blade-system/src/main/java/org/springblade/system/SystemApplication.java b/blade-service/blade-system/src/main/java/org/springblade/system/SystemApplication.java
index 25ec919..23acd8e 100644
--- a/blade-service/blade-system/src/main/java/org/springblade/system/SystemApplication.java
+++ b/blade-service/blade-system/src/main/java/org/springblade/system/SystemApplication.java
@@ -19,6 +19,9 @@
import org.springblade.core.cloud.feign.EnableBladeFeign;
import org.springblade.core.launch.BladeApplication;
import org.springblade.core.launch.constant.AppConstant;
+import org.springblade.system.nettyServer.Server;
+import org.springblade.system.nettyUdpServer.server.UdpServer;
+import org.springframework.boot.CommandLineRunner;
import org.springframework.cloud.client.SpringCloudApplication;
/**
@@ -27,11 +30,18 @@
*/
@EnableBladeFeign
@SpringCloudApplication
-public class SystemApplication {
+public class SystemApplication implements CommandLineRunner {
public static void main(String[] args) {
BladeApplication.run(AppConstant.APPLICATION_SYSTEM_NAME, SystemApplication.class, args);
}
+ @Override
+ public void run(String... args) throws Exception {
+ UdpServer udpServer = new UdpServer(9534);
+ Server server = new Server(9534);
+ UdpServer udpServer8011 = new UdpServer(8011);
+ Server server8011 = new Server(8011);
+ }
}
diff --git a/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/Server.java b/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/Server.java
new file mode 100644
index 0000000..9e09b64
--- /dev/null
+++ b/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/Server.java
@@ -0,0 +1,86 @@
+package org.springblade.system.nettyServer;
+
+import io.netty.bootstrap.ServerBootstrap;
+import io.netty.channel.*;
+import io.netty.channel.nio.NioEventLoopGroup;
+import io.netty.channel.socket.ServerSocketChannel;
+import io.netty.channel.socket.SocketChannel;
+import io.netty.channel.socket.nio.NioServerSocketChannel;
+
+
+public class Server {
+ private int port;
+ private ServerSocketChannel serverSocketChannel;
+
+ public Server(int port){
+ this.port = port;
+ bind();
+ }
+
+ private void bind() {
+ Thread thread = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ //服务端要建立两个group,一个负责接收客户端的连接,一个负责处理数据传输
+ //连接处理group
+ EventLoopGroup boss = new NioEventLoopGroup();
+ //事件处理group
+ EventLoopGroup worker = new NioEventLoopGroup();
+ ServerBootstrap bootstrap = new ServerBootstrap();
+ // 绑定处理group
+ bootstrap.group(boss, worker).channel(NioServerSocketChannel.class)
+ //保持连接数
+ .option(ChannelOption.SO_BACKLOG, 1024)
+ //有数据立即发送
+ .option(ChannelOption.TCP_NODELAY, true)
+ //保持连接
+ .childOption(ChannelOption.SO_KEEPALIVE, true)
+ //处理新连接
+ .childHandler(new ChannelInitializer<SocketChannel>() {
+ @Override
+ protected void initChannel(SocketChannel sc) throws Exception {
+ // 增加任务处理
+ ChannelPipeline p = sc.pipeline();
+ p.addLast(
+// //使用了netty自带的编码器和解码器
+// new StringDecoder(),
+// new StringEncoder(),
+ //心跳检测,读超时,写超时,读写超时
+ //new IdleStateHandler(5, 0, 0, TimeUnit.SECONDS),
+ //自定义的处理器
+ new ServerHandler());
+ }
+ });
+
+ //绑定端口,同步等待成功
+ ChannelFuture future;
+ try {
+ future = bootstrap.bind(port).sync();
+ if (future.isSuccess()) {
+ serverSocketChannel = (ServerSocketChannel) future.channel();
+ System.out.println("服务端启动成功,端口:"+port);
+ } else {
+ System.out.println("服务端启动失败!");
+ }
+
+ //等待服务监听端口关闭,就是由于这里会将线程阻塞,导致无法发送信息,所以我这里开了线程
+ future.channel().closeFuture().sync();
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ finally {
+ //优雅地退出,释放线程池资源
+ boss.shutdownGracefully();
+ worker.shutdownGracefully();
+ }
+ }
+ });
+ thread.start();
+ }
+
+ public void sendMessage(Object msg){
+ if(serverSocketChannel != null){
+ serverSocketChannel.writeAndFlush(msg);
+ }
+ }
+}
diff --git a/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/ServerHandler.java b/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/ServerHandler.java
new file mode 100644
index 0000000..b40a9b4
--- /dev/null
+++ b/blade-service/blade-system/src/main/java/org/springblade/system/nettyServer/ServerHandler.java
@@ -0,0 +1,104 @@
+package org.springblade.system.nettyServer;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.channel.Channel;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.ChannelInboundHandlerAdapter;
+import org.springframework.stereotype.Component;
+import javax.annotation.PostConstruct;
+import java.text.SimpleDateFormat;
+import java.util.Date;
+import java.util.concurrent.ConcurrentHashMap;
+
+@Component
+public class ServerHandler extends ChannelInboundHandlerAdapter {
+
+ private String reg_LA = "LA[d]{8}[d|A-F]{12}[d|A-F]{8}[d|A-Z]{2}[d|A-F]{4}[x2A][d|A-F]{6}[#@]";
+ private String reg_LB = "LB[\\d|A-F]{12}[\\x2A].*[#@]";
+ private String reg_LB2 = "LB[\\d|A-F]{6}[\\x2A].*[#@]";
+ private String reg_LD = "LD[d]{8}[d|A-F]{12}:[A-Z]{4}[\\x2A].*[#@]";
+
+
+ private ConcurrentHashMap<String, Channel> sessionChannelMap = new ConcurrentHashMap<String, Channel>();
+
+ private static ServerHandler serverHandler;
+
+ @PostConstruct
+ public void init() {
+ serverHandler = this;
+ }
+
+ /**
+ * 客户端与服务端创建连接的时候调用
+ */
+ @Override
+ public void channelActive(ChannelHandlerContext ctx) throws Exception {
+ System.out.println("CTX:" + ctx.channel());
+ System.out.println("客户端与服务端连接开始...");
+// ChannelId id = ctx.channel().id();
+// String s = id.toString();
+ //NettyConfig.group.add(ctx.channel());
+ }
+
+ /**
+ * 客户端与服务端断开连接时调用
+ */
+ @Override
+ public void channelInactive(ChannelHandlerContext ctx) throws Exception {
+ System.out.println("客户端与服务端连接关闭...");
+ }
+
+ /**
+ * 服务端接收客户端发送过来的数据结束之后调用
+ */
+ @Override
+ public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
+ ctx.flush();
+ System.out.println("信息接收完毕...");
+ }
+
+ /**
+ * 工程出现异常的时候调用
+ */
+ @Override
+ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
+ cause.printStackTrace();
+ ctx.close();
+ }
+
+ /**
+ * 服务端处理客户端websocket请求的核心方法,这里接收了客户端发来的信息
+ */
+ @Override
+ public void channelRead(ChannelHandlerContext channelHandlerContext, Object info) throws Exception {
+ long startTime1 = System.currentTimeMillis();
+ Date day = new Date();
+ SimpleDateFormat df = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
+ System.out.println("接收到了:" + info);
+ ByteBuf buf = (ByteBuf) info;
+ byte[] req = new byte[buf.readableBytes()];
+ buf.readBytes(req);
+ String body = new String(req, "UTF-8");
+ String content = body;
+ System.out.println("【TCP】>>>>>> 接收客户端数据:" + body);
+ String state="系统故障"; //设备故障状态
+ String reason=null; //故障原因
+ //获取字符串的长度
+ int length = content.length();
+ //当字符串长度大于200时,采集体温数据,小于则采集报警数据
+ System.out.println("tcp报文长度 = " + length);
+
+// channelHandlerContext.writeAndFlush(Unpooled.copiedBuffer("LEOK#@!", CharsetUtil.UTF_8));
+ //服务端使用这个就能向 每个连接上来的客户端群发消息
+ //NettyConfig.group.writeAndFlush(info);
+// Iterator<Channel> iterator = NettyConfig.group.iterator();
+// System.out.println("iterator:"+iterator);
+// while(iterator.hasNext()){
+// //打印出所有客户端的远程地址
+// System.out.println((iterator.next()).remoteAddress());
+//
+// }
+ }
+
+
+}
diff --git a/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/handle/UdpServerHandler.java b/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/handle/UdpServerHandler.java
new file mode 100644
index 0000000..762c465
--- /dev/null
+++ b/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/handle/UdpServerHandler.java
@@ -0,0 +1,51 @@
+package org.springblade.system.nettyUdpServer.handle;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.SimpleChannelInboundHandler;
+import io.netty.channel.socket.DatagramPacket;
+import io.netty.util.CharsetUtil;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Component;
+
+/**
+ * updHandler udp 服务端数据接收处理
+ */
+@Component
+public class UdpServerHandler extends SimpleChannelInboundHandler<DatagramPacket> {
+ private Logger logger = LoggerFactory.getLogger(this.getClass());
+
+ @Override
+ protected void channelRead0(ChannelHandlerContext channelHandlerContext, DatagramPacket datagramPacket) throws Exception {
+ // 读取收到的数据
+ ByteBuf buf = (ByteBuf) datagramPacket.copy().content();
+ byte[] req = new byte[buf.readableBytes()];
+ buf.readBytes(req);
+ String body = new String(req, CharsetUtil.UTF_8);
+ System.out.println("【UDP】>>>>>> 收到客户端的数据:"+body);
+ System.out.println("upd报文长度 = " + body.length());
+
+ // 回复一条信息给客户端
+// channelHandlerContext.writeAndFlush(new DatagramPacket(
+// Unpooled.copiedBuffer("Hello,我是Server,我的时间戳是"+System.currentTimeMillis()
+// , CharsetUtil.UTF_8)
+// , datagramPacket.sender())).sync();
+ }
+
+ //捕获异常
+ @Override
+ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause)throws Exception {
+ //logger.log(Level.INFO, "AuthServerInitHandler exceptionCaught");
+ logger.error("UdpServerHandler exceptionCaught"+cause.getMessage());
+ System.out.println("UdpServerHandler exceptionCaught"+cause.getMessage());
+ cause.printStackTrace();
+ ctx.close();
+ }
+
+ //消息没有结束的时候触发
+ @Override
+ public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
+ ctx.flush();
+ }
+}
diff --git a/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/server/UdpServer.java b/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/server/UdpServer.java
new file mode 100644
index 0000000..33b4e41
--- /dev/null
+++ b/blade-service/blade-system/src/main/java/org/springblade/system/nettyUdpServer/server/UdpServer.java
@@ -0,0 +1,129 @@
+package org.springblade.system.nettyUdpServer.server;
+
+import io.netty.bootstrap.Bootstrap;
+import io.netty.channel.*;
+import io.netty.channel.nio.NioEventLoopGroup;
+import io.netty.channel.socket.nio.NioDatagramChannel;
+import io.netty.handler.codec.json.JsonObjectDecoder;
+import io.netty.handler.codec.string.StringDecoder;
+import io.netty.handler.codec.string.StringEncoder;
+import io.netty.util.concurrent.DefaultEventExecutorGroup;
+import io.netty.util.concurrent.EventExecutorGroup;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springblade.system.nettyUdpServer.handle.UdpServerHandler;
+
+import java.net.InetAddress;
+import java.util.concurrent.Executors;
+
+/**
+ * updServer udp服务端
+ */
+public class UdpServer {
+ private static Logger LOG = LoggerFactory.getLogger(UdpServer.class);
+
+ public int port;
+
+ private Channel channel;
+
+ public UdpServer(int port) {
+ this.port = port;
+ run();
+ }
+
+ /**
+ * 不用
+ */
+ public void init() {
+ final NioEventLoopGroup group = new NioEventLoopGroup();
+ try {
+ Bootstrap b = new Bootstrap();
+ b.group(group).channel(NioDatagramChannel.class)
+ .option(ChannelOption.SO_BROADCAST, true)
+ .handler(new ChannelInitializer<NioDatagramChannel>() {
+ @Override
+ public void initChannel(NioDatagramChannel ch) throws Exception {
+ ChannelPipeline p = ch.pipeline();
+ p.addLast("decoder", new JsonObjectDecoder());
+ p.addLast(new UdpServerHandler());
+ }
+ });
+ InetAddress address = InetAddress.getLocalHost();
+ channel = b.bind(address, port).sync().channel();
+ Executors.newSingleThreadExecutor().execute(new Runnable() {
+ @Override
+ public void run() {
+ try {
+ channel.closeFuture().sync();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ } finally {
+ group.shutdownGracefully();
+ }
+ }
+ });
+ LOG.info("UDP服务器启动, host:{},port:{}", address.getHostAddress(), port);
+ } catch (Exception e) {
+ e.printStackTrace();
+ }
+ }
+
+
+ //给管道抽象出接口,给Channel更多的能力和配置,例如Channel的状态,参数,IO操作
+ //使用ChannelPipeline实现自定义IO
+ //Channel channel;
+ public void run() {
+ Thread thread = new Thread(new Runnable() {
+ @Override
+ public void run() {
+ //InetSocketAddress socketAddress = new InetSocketAddress("s16s652780.51mypc.cn", 21403);
+ //启动服务
+ EventLoopGroup workerGroup = new NioEventLoopGroup();
+ //优化使用的线程
+ final EventExecutorGroup group = new DefaultEventExecutorGroup(16);
+
+ try {
+ Bootstrap b = new Bootstrap();//udp不能使用ServerBootstrap
+ b.group(workerGroup).channel(NioDatagramChannel.class)//设置UDP通道
+ //设置udp的管道工厂
+ .handler(new ChannelInitializer<NioDatagramChannel>() {
+ //NioDatagramChannel标志着是UDP格式的
+ @Override
+ protected void initChannel(NioDatagramChannel ch)
+ throws Exception {
+ // TODO Auto-generated method stub
+ //创建一个执行Handler的容器
+ ChannelPipeline pipeline = ch.pipeline();
+ pipeline.addLast(new StringDecoder());
+ pipeline.addLast(new StringEncoder());
+ //执行具体的处理器
+ pipeline.addLast(group, "handler", new UdpServerHandler());//消息处理器
+ }
+
+ })//初始化处理器
+ //true / false 多播模式(UDP适用),可以向多个主机发送消息
+ .option(ChannelOption.SO_BROADCAST, true)// 支持广播
+ .option(ChannelOption.SO_RCVBUF, 2048 * 1024)// 设置UDP读缓冲区为2M
+ .option(ChannelOption.SO_SNDBUF, 1024 * 1024);// 设置UDP写缓冲区为1M
+
+ // 绑定端口,开始接收进来的连接
+ ChannelFuture f = b.bind(port).sync();
+ //获取channel通道
+ //channel=f.channel();
+ System.out.println("UDP Server 启动,端口:"+port);
+ // 等待服务器 socket 关闭 。
+ // 这不会发生,可以优雅地关闭服务器。
+ f.channel().closeFuture().sync();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ } finally {
+ // 优雅退出 释放线程池资源
+ group.shutdownGracefully();
+ workerGroup.shutdownGracefully();
+ }
+ }
+ });
+ thread.start();
+ }
+
+}
--
Gitblit v1.9.3