Netty异步高性能高可用通信框架:
Netty 提供异步的、基于事件驱动的网络应用程序框架,用以快速开发高性能、高可靠性的网络 IO 程序,高性能、吞吐量更高:延迟更低;减少资源消耗;最小化不必要的内存复制,是目前最流行的 NIO 框架,
Netty在互联网领域、大数据分布式计算领域、游戏行业、通信行业等获得了广泛的应用,知名的 Elasticsearch 、Dubbo 框架内部都采用了Netty。
Netty 的线程模型是主要是基于主从 Reactor 多线程模型
所以这里得先介绍一下:主从 Reactor 多线程模型:
主从 Reactor 多线程
1.工作原理:
①Reactor主线程 MainReactor 对象通过select 监听连接事件, 收到事件后,通过Acceptor 处理连接事件
②当 Acceptor 处理连接事件后,MainReactor 将连接分配给SubReactor
③subReactor 将连接加入到连接队列进行监听,并创建handler进行各种事件处理
④当有新事件发生时, subreactor就会调用对应的handler处理
⑤handler 通过read 读取数据,分发给后面的worker 线程处理
⑥worker 线程池分配独立的worker 线程进行业务处理,并返回结果
⑦handler 收到响应的结果后,再通过send 将结果返回给client
⑧Reactor 主线程可以对应多个Reactor 子线程, 即MainRecator 可以关联多个SubReactor

.
.
了解Netty得工作原理:Netty线程模型
1.工作原理
Netty抽象出两组线程池 BossGroup 专门负责接收客户端的连接, WorkerGroup 专门负责网络的读写
BossGroup 和 WorkerGroup 类型都是 NioEventLoopGroup
NioEventLoopGroup 相当于一个事件循环组, 这个组中含有多个事件循环 ,每一个事件循环是 NioEventLoop
NioEventLoop 表示一个不断循环的执行处理任务的线程, 每个NioEventLoop 都有一个selector , 用于监听绑定在该通道上的socket的网络通讯
NioEventLoopGroup 可以有多个线程, 即可以含有多个NioEventLoop
每个Boss NioEventLoop 循环执行的步骤有3步
轮询accept 事件
处理accept 事件 , 与client端建立连接 , 生成NioScocketChannel , 并将其注册到某个worker NIOEventLoop 上的 selector 上
处理任务队列的任务 , 即 runAllTasks
每个 Worker NIOEventLoop 循环执行的步骤
轮询read, write 事件
处理i/o事件, 即read , write 事件,在对应NioScocketChannel 处理
处理任务队列的任务 , 即 runAllTasks
每个Worker NIOEventLoop 处理业务时,会使用pipeline(管道), pipeline 中包含了boss group上NioEventLoop注册到worker 的selector 的channel , 即通过pipeline 可以获取到对应通道, 管道中维护了很多的处理器

我个人简单理解netty:
在主reactoor中selector监听client连接,accept实现连接,连接完成生成channel对象,对象中标识了socket,将对象注册到从reactor得selector上,把任务放到队列中,轮询队列,判断是write还是read操作,然后交给相应得pipleline处理业务,pipeline中有channel,channel中有channelHandlerContext为节点得双向链接,每个节点中就是可以处理业务逻辑得处理器channelHandler。
.
.
springboot整合netty:
加入依赖:
<dependency>
<groupId>io.netty</groupId>
<artifactId>netty-all</artifactId>
<version>4.1.20.Final</version>
</dependency>
客户端
public class NettyClient {
public static void main(String[] args) {
// 定义一个循环事件组
NioEventLoopGroup group= new NioEventLoopGroup();
try {
// 创建客户端启动对象
Bootstrap bootstrap = new Bootstrap();
bootstrap.group(group)
.channel(NioSocketChannel.class)
.handler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel socketChannel) throws Exception {
socketChannel.pipeline().addLast(new NettyClientHandler());
}
});
System.err.println("client is ready...");
// 启动客户端连接服务端
ChannelFuture channelFuture = bootstrap.connect("127.0.0.1", 8082).sync();
// 设置通道关闭监听(当监听到通道关闭时,关闭client)
channelFuture.channel().closeFuture().sync();
} catch (InterruptedException e) {
e.printStackTrace();
} finally {
group.shutdownGracefully();
}
}
}
客户端处理器
public class NettyClientHandler extends ChannelInboundHandlerAdapter {
/**
* client端服务启动完成后续
* @param ctx
* @throws Exception
*/
@Override
public void channelActive(ChannelHandlerContext ctx) throws Exception {
System.err.println("client "+ctx);
// write + flush操作
// 将数据写入到缓存,并刷新 ,通常需要对发送的数据进行base64或其他编码
ctx.writeAndFlush(Unpooled.copiedBuffer("hello,netty server...", CharsetUtil.UTF_8));
}
/**
* 当通道有读事件时
* @param ctx
* @param msg
* @throws Exception
*/
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
ByteBuf byteBuf = (ByteBuf) msg;
System.err.println("服务器端回复消息:"+byteBuf.toString(CharsetUtil.UTF_8));
System.err.println("服务器端地址是:"+ctx.channel().remoteAddress());
}
/**
* 通道有异常时
* @param ctx
* @param cause
* @throws Exception
*/
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
//处理异常, 一般是需要关闭通道
cause.printStackTrace();
ctx.close();
}
}
.
.
服务端
public class NettyServer {
public static void main(String[] args) throws InterruptedException {
// 负责请求连接
NioEventLoopGroup bossGroup = new NioEventLoopGroup();
// 负责网络读写处理
NioEventLoopGroup workGroup = new NioEventLoopGroup();
// 服务器得启动对象,目的:为服务端启动配置一些服务参数
ServerBootstrap serverBootstrap = new ServerBootstrap();
// 链式编程配置参数
serverBootstrap.group(bossGroup,workGroup)
// 使用NioServerSocketChannel作为服务器的通道
.channel(NioServerSocketChannel.class)
// 设置线程等待的连接个数
.option(ChannelOption.SO_BACKLOG,128)
// 给PipeLine设置处理器
.childHandler(new ChannelInitializer<SocketChannel>() {
@Override
protected void initChannel(SocketChannel socketChannel) throws Exception {
// 添加处理事件
socketChannel.pipeline().addLast(new NettyServerHandler());
}
});
System.err.println("server is ready ...");
// 启动服务器,并绑定1个端口且同步生成一个ChannelFuture 对象
ChannelFuture channelFuture = serverBootstrap.bind(8082).sync();
//对关闭通道进行监听(netty异步模型) 当通道进行关闭时,才会触发这个关闭动作
channelFuture.channel().closeFuture().sync();
}
}
服务端处理器
public class NettyServerHandler extends ChannelInboundHandlerAdapter {
/**
* 读取数据
* @param ctx 上下文对象, 含pipeline管道 、通道channel,地址
* @param msg 客户端发送得数据 默认object
* @throws Exception
*/
@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
System.err.println("服务器读取线程"+Thread.currentThread().getName());
System.out.println("server ctx = "+ ctx);
System.out.println("查看channel 与 pipeline关系");
// 这里可以利用通道得exeut()或者shedule()线程去处理业务逻辑
Channel channel = ctx.channel();
// pipeline本质是一个双向'链表'的感觉
ChannelPipeline pipeline = ctx.pipeline();
// 两种线程继续处理业务逻辑
/* channel.eventLoop().execute(new Runnable() {
@Override
public void run() {
System.out.println("线程正在处理【审核】模块添加功能....");
}
});*/
channel.eventLoop().schedule(new Runnable() {
@Override
public void run() {
System.out.println("线程定时【10s】处理处理【审核】模块添加功能....");
}
},10, TimeUnit.SECONDS);
ByteBuf buf =(ByteBuf) msg;
System.out.println("客户端发送消息:"+buf.toString(CharsetUtil.UTF_8));
System.out.println("客户端地:"+channel.remoteAddress());
}
/**
* 读取数据完成后续操作
* @param ctx
* @throws Exception
*/
@Override
public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
System.err.println("client "+ctx);
// write + flush操作
// 将数据写入到缓存,并刷新 ,通常需要对发送的数据进行base64或其他编码
ctx.writeAndFlush(Unpooled.copiedBuffer("hello, netty client.....",CharsetUtil.UTF_8));
}
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
//处理异常, 一般是需要关闭通道
cause.printStackTrace();
ctx.close();
}
}