From 6a65cca452ce1ec3d6132698c91b12eccdabca57 Mon Sep 17 00:00:00 2001
From: 钟日健 <5689795+arsn@user.noreply.gitee.com>
Date: Wed, 24 Feb 2021 20:31:14 +0800
Subject: [PATCH] udpServer 数据接收
---
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/JfptApplication.java | 8 +-
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServer.java | 123 ++++++++++++++++++++++++++++++
blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServerHandler.java | 51 ++++++++++++
3 files changed, 178 insertions(+), 4 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 9fd26cd..9d89f2d 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
@@ -18,12 +18,10 @@
import org.springblade.core.cloud.feign.EnableBladeFeign;
import org.springblade.core.launch.BladeApplication;
-import org.springblade.core.launch.constant.AppConstant;
import org.springblade.jfpt.nettyServer.Server;
+import org.springblade.jfpt.nettyUdpServer.UdpServer;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.autoconfigure.SpringBootApplication;
-import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
-import org.springframework.cloud.client.SpringCloudApplication;
/**
* Desk启动器
@@ -41,7 +39,9 @@
@Override
public void run(String... args) throws Exception {
- Server server=new Server(8088);
+ UdpServer udpServer=new UdpServer(8088);
+ //Server server=new Server(8088);
}
+
}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServer.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServer.java
new file mode 100644
index 0000000..6083298
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServer.java
@@ -0,0 +1,123 @@
+package org.springblade.jfpt.nettyUdpServer;
+
+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 java.net.InetAddress;
+import java.net.InetSocketAddress;
+import java.util.concurrent.Executors;
+
+/**
+ * Created by wangshikai on 2016/7/18.
+ */
+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() {
+ //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 启动!");
+ // 等待服务器 socket 关闭 。
+ // 这不会发生,可以优雅地关闭服务器。
+ f.channel().closeFuture().sync();
+ } catch (InterruptedException e) {
+ e.printStackTrace();
+ } finally {
+ // 优雅退出 释放线程池资源
+ group.shutdownGracefully();
+ workerGroup.shutdownGracefully();
+ }
+ }
+
+}
diff --git a/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServerHandler.java b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServerHandler.java
new file mode 100644
index 0000000..6189787
--- /dev/null
+++ b/blade-service/blade-jfpts/src/main/java/org/springblade/jfpt/nettyUdpServer/UdpServerHandler.java
@@ -0,0 +1,51 @@
+package org.springblade.jfpt.nettyUdpServer;
+
+import io.netty.buffer.ByteBuf;
+import io.netty.buffer.Unpooled;
+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;
+
+/**
+ * Created by wangshikai on 2016/7/18.
+ */
+@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("【NOTE】>>>>>> 收到客户端的数据:"+body);
+
+ // 回复一条信息给客户端
+// 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();
+ }
+}
--
Gitblit v1.9.3