netty实战(一)
2021-01-23 15:12
标签:ext art 通过 mybatis 目的 child reac autowired Matter 先分享一下自己的经历。 去年7月进入新公司没多久,部门领导就给我分配了一个任务:给公司的一个户外设备写一个采集数据程序,将数据入库,然后做一个web端。因为领导是做.NET的,当时在来之前有和领导沟通过,领导的意思是希望来一个会网络编程和多线程,部门急需一个可以来做采集程序的java,我当时有点心虚,是这样回复领导:自己也只是有1年多java的工作,不会网络编程,自己搭一个简单的web项目架构还是可以应付的,我只能尽自己的最大的可能做好这件事。 后来进入公司后,做采集程序是一头雾水,自己之前只是做web的,对web项目的架构以及业务还是比较熟悉,但是关于和硬件之间的通讯还真的是从来都没有接触过。It is the first step that is troublesome。问过一些朋友,也百度过看过一些技术博客,有用socket 连接,多线程后加while死循环来接受做采集,后来慢慢发现用netty这种NIO网络应用框架比较合适。 接受就是熟悉一下具体的业务,户外设备会主动连接服务器IP,定时向服务器的指定端口发送报文。所以当前的采集程序结合netty的情况就是需要写一个netty的server端,来监听服务器本机的该指定端口,接受户外设备发送过来的报文即可,而不需要再写一个netty的client端。 首先说明一下,本项目应用的框架:springboot+netty+rabbitMq+mybatis+lombok+logback,数据库选用的是MYSQL。 开发环境:idea2018+jdk1.8+mysql5.6.35+maven3.5.3 接下来就是搭建项目: 1.idea快速创建springboot项目,在pom.xml文件中配置依赖包 2.自定义netty的server端BootNettyServer类 3.自定义初始化类BootNettyInitializer类 继承ChannelInitializer,实现initChannel方法,初始化信道时可以配置心跳、编码器、解码器、以及业务处理ChannelHandler类。 4.自定义业务处理类BootNettyHandler 回到开始位置,自定义netty的server端BootNettyServer类该如何启动? 我的方案是在application启动类中来启动。 application类实现CommandLineRunner,来启动nettyserver服务 至此,采集程序完成。 此程序仅个人学习理解后完成,如有不足的地方,还请进一步多多给出指导,thank you !!! netty实战(一) 标签:ext art 通过 mybatis 目的 child reac autowired Matter 原文地址:https://www.cnblogs.com/lyzj/p/13278414.htmlpackage com.rtstjkx.jkxlistener.netty;
import io.netty.bootstrap.ServerBootstrap;
import io.netty.channel.AdaptiveRecvByteBufAllocator;
import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelOption;
import io.netty.channel.EventLoopGroup;
import io.netty.channel.nio.NioEventLoopGroup;
import io.netty.channel.socket.nio.NioServerSocketChannel;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
/**
* netty的server
*/
@Component
@Slf4j
public class BootNettyServer {
@Autowired
BootNettyInitializer bootNettyInitializer;
public void run(int port) throws Exception {
/**
* 配置服务端的NIO线程组
* NioEventLoopGroup 是用来处理I/O操作的Reactor线程组
* bossGroup:用来接收进来的连接,workerGroup:用来处理已经被接收的连接,进行socketChannel的网络读写,
* bossGroup接收到连接后就会把连接信息注册到workerGroup
* workerGroup的EventLoopGroup默认的线程数是CPU核数的二倍
*/
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();
try {
/**
* ServerBootstrap 是一个启动NIO服务的辅助启动类
*/
ServerBootstrap serverBootstrap = new ServerBootstrap();
/**
* 设置group,将bossGroup, workerGroup线程组传递到ServerBootstrap
*/
serverBootstrap = serverBootstrap.group(bossGroup, workerGroup);
/**
* ServerSocketChannel是以NIO的selector为基础进行实现的,用来接收新的连接,这里告诉Channel通过NioServerSocketChannel获取新的连接
*/
serverBootstrap = serverBootstrap.channel(NioServerSocketChannel.class);
/**
* option是设置 bossGroup,childOption是设置workerGroup
* netty 默认数据包传输大小为1024字节, 设置它可以自动调整下一次缓冲区建立时分配的空间大小,避免内存的浪费 最小 初始化 最大 (根据生产环境实际情况来定)
* 使用对象池,重用缓冲区
*/
serverBootstrap = serverBootstrap.option(ChannelOption.RCVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(64, 10496, 1048576));
serverBootstrap = serverBootstrap.childOption(ChannelOption.RCVBUF_ALLOCATOR, new AdaptiveRecvByteBufAllocator(64, 10496, 1048576));
System.out.println("正在监听8234端口中....");
/**
* 设置 I/O处理类,主要用于网络I/O事件,记录日志,编码、解码消息
*/
serverBootstrap = serverBootstrap.childHandler(bootNettyInitializer);
/**
* 绑定端口,同步等待成功
*/
ChannelFuture f = serverBootstrap.bind(port).sync();
/**
* 等待服务器监听端口关闭
*/
f.channel().closeFuture().sync();
} catch (InterruptedException e) {
} finally {
/**
* 退出,释放线程池资源
*/
bossGroup.shutdownGracefully();
workerGroup.shutdownGracefully();
}
}
}
package com.rtstjkx.jkxlistener.netty;
import com.rtstjkx.jkxlistener.netty.coder.MyDecoder;
import io.netty.channel.Channel;
import io.netty.channel.ChannelInitializer;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
/**
* 通道初始化
* @param
package com.rtstjkx.jkxlistener.netty;
import com.rtstjkx.jkxlistener.controller.OrderController;
import com.rtstjkx.jkxlistener.entity.Order;
import com.rtstjkx.jkxlistener.rabbitmq.RabbitProducer;
import com.rtstjkx.jkxlistener.service.serviceImpl.OrderServiceImpl;
import com.rtstjkx.jkxlistener.util.StringUtil;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.channel.*;
import io.netty.channel.socket.SocketChannel;
import io.netty.handler.timeout.IdleStateEvent;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.text.SimpleDateFormat;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.util.Date;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* I/O数据读写处理类
*/
@Component
@ChannelHandler.Sharable
@Slf4j
public class BootNettyHandler extends ChannelInboundHandlerAdapter {
// 将当前客户端连接 存入map 实现控制设备下发 参数
public static Map
//将报文信息记录到日志中
log.info(channel.remoteAddress().getHostString() + ": " + input);
System.out.println("服务端接受信息为: " + new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date()) + " 接收到消息:" + input);
//具体的处理业务逻辑,解析报文帧结构,将数据入库
......
}
/**
* 从客户端收到新的数据、读取完成时调用
*
* @param ctx
*/
@Override
public void channelReadComplete(ChannelHandlerContext ctx)
{
ctx.flush();
}
/**
* 当出现 Throwable 对象才会被调用,即当 Netty 由于 IO 错误或者处理器在处理事件时抛出的异常时
*
* @param ctx
* @param cause
*/
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause)
{
System.out.println("程序异常,断开客户端连接");
//获取客户端的请求地址 取到的值为客户端的 ip+端口号
InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress();
String clientIp = insocket.getAddress().getHostAddress();//设备IP地址(个人将设备的IP地址当作 map 的key,channel对象做value)
if(ctxMap.get(clientIp)!=null){//如果不为空就剔除
ctxMap.remove(clientIp, ctx.channel());
}else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key
}
cause.printStackTrace();
ctx.close();//抛出异常,断开与客户端的连接
log.info(clientIp+"连接断开");
}
/**
* 客户端与服务端第一次建立连接时 执行
*
* @param ctx
* @throws Exception
*/
@Override
public void channelActive(ChannelHandlerContext ctx)
{
InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress();
String clientIp = insocket.getAddress().getHostAddress();
if(ctxMap.get(clientIp)!=null){//如果不为空就不存
}else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key
ctxMap.put(clientIp, ctx.channel());
}
log.info(clientIp+"连接");
}
/**
* 客户端与服务端 断连时 执行
*
* @param ctx
* @throws Exception
*/
@Override
public void channelInactive(ChannelHandlerContext ctx) throws Exception
{
super.channelInactive(ctx);
InetSocketAddress insocket = (InetSocketAddress) ctx.channel().remoteAddress();
String clientIp = insocket.getAddress().getHostAddress();
if(ctxMap.get(clientIp)!=null){//如果不为空就剔除
ctxMap.remove(clientIp, ctx.channel());
}else{//否则就将当前的设备ip+端口存进map 当做下发设备的标识的key
}
ctx.close(); //断开连接时,必须关闭,否则造成资源浪费,并发量很大情况下可能造成宕机
log.info(clientIp+"断开");
}
/**
* 服务端当read超时, 会调用这个方法
*
* @param ctx
* @param evt
* @throws Exception
*/
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) throws Exception, IOException
{
if (evt instanceof IdleStateEvent) {
IdleStateEvent event = (IdleStateEvent) evt;
String eventType = null;
switch (event.state()) {
case READER_IDLE:
eventType = "读空闲";
break;
case WRITER_IDLE:
eventType = "些空闲";
break;
case ALL_IDLE:
eventType = "读写空闲";
break;
default:
}
String date = LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"));
log.info(date + " " + ctx.channel().remoteAddress() + " " + eventType);
}
}
}package com.rtstjkx.jkxlistener;
import com.rtstjkx.jkxlistener.netty.BootNettyServer;
import org.mybatis.spring.annotation.MapperScan;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.boot.CommandLineRunner;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
@MapperScan("com.rtstjkx.jkxlistener.repository")
@SpringBootApplication
public class KjxApplication implements CommandLineRunner {
@Value("${netty.port}")
private Integer port;
@Autowired
BootNettyServer bootNettyServer;
public static void main(String[] args) throws Exception {
SpringApplication.run(KjxApplication.class, args);
}
@Override
public void run(String... args) throws Exception {
/**
* 启动netty服务端服务
*/
bootNettyServer.run(port);
}
}
上一篇:JMeter生成HTML测试报告
下一篇:CTF-web8