• Netty异步高性能高可用通信框架


    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>
    
    • 1
    • 2
    • 3
    • 4
    • 5

    客户端

    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();
            }
    
        }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33

    客户端处理器

    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();
        } 
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42

    .
    .

    服务端

    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();
        }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35

    服务端处理器

    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();
        }
    }
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
    • 41
    • 42
    • 43
    • 44
    • 45
    • 46
    • 47
    • 48
    • 49
    • 50
    • 51
    • 52
    • 53
    • 54
    • 55
    • 56
    • 57
    • 58
    • 59
  • 相关阅读:
    Jetson系列设置Python脚本开机自启
    GZ038 物联网应用开发赛题第7套
    leetcode第311场周赛
    【狂神说Java】Mybatis学习笔记(下)
    联邦学习综述二
    python趣味编程-5分钟实现一个蛇梯游戏(含源码、步骤讲解)
    低代码之光!轻量级 GUI 的设计与实现
    Advanced .Net Debugging 9:平台互用性
    Python基础教程之五:Python中的数据类型
    万物皆可长按:SwiftUI 5.0(iOS 17)极简原生实现任意视图长按惯性加速功能
  • 原文地址:https://blog.csdn.net/qq_45399396/article/details/126766393