liuyulin
发布于 2024-02-02 / 48 阅读
0
0

Netty

Netty

Netty是一个网络IO,可以理解为是对java原生IO的一个优化。Netty是一个异步基于事件驱动的网络应用程序框架,用于快速开发可维护的高性能协议服务器和客户端。Netty主要针对TCP协议下,面向client端的高并发应用,或者是P2P场景下的大量数据持续传输的应用。

Netty:

image-20230107113459578

IO模型

IO模型:用什么样的通道进行数据的发送和接收,很大程度上决定了程序通信性能。现在又BIO,NIO和AIO

BIO:同步并阻塞,服务器实现模式为一个连接对应一个线程,客户端有连接请求时,服务器就会有一个线程进行处理,如果连接不做任何事情,就会造成不必要的线程开销。

NIO:同步非阻塞,服务器实现模式为一个线程处理多个连接,即客户端发送的请求都会注册到一个多路复用器上,多路复用器轮询到有IO请求的时候,就会进行处理。

AIO:异步非阻塞,引入异步通道的概念,采用Proactor模式,简化程序编写,有效的程序才启动线程,它的特点是由操作系统完成后才通知服务端程序启动线程去处理,一般适用于连接数较多且连接时间长的应用,目前不成熟。

适用场景:

BIO:连接数目较少且固定的架构,对服务器资源要求较高,并发局限于应用中。

NIO:连接数目多且连接比较短的架构。

AIO:连接数目多且连接比较长的架构。

BIO

BIO就是传统的java io,同步阻塞的。传统的BIO可以通过线程池技术来改善,线程池不能减少连接的个数,只是让多个客户连接,就是并发。

BIO编程的简单流程:

  1. 服务器端启动一个ServerSocket

  2. 客户端启动一个Socket对服务器进行通信,默认情况下服务器端需要对每个客户建立一个线程,与之通信

  3. 客户端发出请求后,先咨询服务器是否有线程响应,如果没有则会等待,或者遭到拒绝

  4. 如果线程有响应,客户端线程会等待请求结束后才继续请求

简单示例:

package com.liuyulin10.bio;

import java.io.IOException;
import java.io.InputStream;
import java.net.ServerSocket;
import java.net.Socket;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

/**
 * BIOServer 适用线程池机制
 */
public class BIOServer {

    public static void main(String[] args) throws Exception {
        //创建一个线程池
        //如果有一个客户端连接 就创建一个线程于之通信(单独方法)
        ExecutorService newCachedThreadPool = Executors.newCachedThreadPool();
        //创建一个ServerSocket
        ServerSocket serverSocket = new ServerSocket(8778);
        System.out.println("服务器启动");

        while (true) {
            //监听 等待客户端连接
            final Socket socket = serverSocket.accept();
            System.out.println("连接到一个客户端");

            //创建一个线程于之通信
            newCachedThreadPool.execute(new Runnable() {
                @Override
                public void run() {
                    System.out.println("启动线程,ID:" + Thread.currentThread().getId());
                    //可以与客户端通信
                    handler(socket);
                }
            });
        }
    }

    //编写一个handler方法与客户端通信
    public static void handler(Socket socket) {
        byte[] bytes = new byte[1024];
        //通过socket获取输入流
        try {
            InputStream inputStream = socket.getInputStream();

            //循环读取客户端发来的数据
            while (true) {
                int read = inputStream.read(bytes);

                if (read != -1) {
                    //输出客户端发送的数据
                    System.out.println(new String(bytes, 0, read));
                } else {
                    //读取完毕了
                    break;
                }
            }


        } catch (IOException e) {
            e.printStackTrace();
        } finally {
            System.out.println("关闭客户端连接");
            try {
                socket.close();
            } catch (IOException e) {
                e.printStackTrace();
            }
        }
    }
}

NIO

NIO同步非阻塞,NIO有三大核心部分:Channel(通道),Buffer(缓冲区),Selector(选择器)。有了缓冲区,NIO就是面向缓冲区,或着面向块编程的,数据读取到一个它稍后处理的缓冲区,需要时可以前后移动,增加了灵活性,可以提供非阻塞的高伸缩性网络。

Buffer

缓冲区本质是一个可读写的内存块,可以理解为一个容器对象(含数组),该对象提供了一组方法,可以更轻松的适用内存块,channel读写数据必须经过buffer。

BasicBuffer:

package com.liuyulin10.nio;

import java.nio.IntBuffer;

/**
 * buffer的适用
 */
public class BasicBuffer {

    public static void main(String[] args) {

        //创建一个buffer,大小为5,即可以存5个int
        IntBuffer intBuffer = IntBuffer.allocate(5);

        //向buffer中存放数据
        for(int i = 0; i<intBuffer.capacity();i++){
            intBuffer.put(10+i);
        }

        //从buffer读取数据
        //buffer读写切换
        intBuffer.flip();

        while (intBuffer.hasRemaining()){
            //get中维护了一个索引,每取一次就后移
            System.out.println(intBuffer.get());
        }
        
    }

}
Buffer类及其子类

ByteBuffer, ShortBuffer, CharBuffer, IntBuffer, LongBuffer, DoubleBuffer, FloatBuffer.

在网络编程中最常用的是ByteBuffer,ReadOnlyBuffer是只读的Buffer,MappedByteBuffer可以让文件在文件在内存中直接修改,即操作系统不需要拷贝一次。

Buffer中定义了所有缓冲区的四个属性,用来提供关于其所包含的数据元素的信息

image-20230107152038676

image-20230107152416614

Buffer的分散和聚集

NIO支持通过多个Buffer来完成读写操作,即Scattering(分散)和Gathering(聚集)。

分散指将数据写入到buffer时,可以采用buffer数组,依次写入;聚集指在读出到buffer中时可以采用数组。

public class ScatteringAndGatheringTest {
    public static void main(String[] args) throws IOException {

        //使用ServerSocket和SocketChannel网络
        ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
        InetSocketAddress inetSocketAddress = new InetSocketAddress(7771);

        //绑定端口到socket 并启动
        serverSocketChannel.socket().bind(inetSocketAddress);

        //创建Buffer数组
        ByteBuffer[] byteBuffers = new ByteBuffer[2];
        int buffer0size = 5;
        int buffer1size = 3;
        byteBuffers[0] = ByteBuffer.allocate(buffer0size);
        byteBuffers[1] = ByteBuffer.allocate(buffer1size);

        //等待客户端连接
        SocketChannel socketChannel = serverSocketChannel.accept();

        long bufferArraySize = buffer0size + buffer1size;

        //循环读取
        while (true) {
            long read = socketChannel.read(byteBuffers);
            if (read == -1) {
                break;
            }
            long psize = read;
            for (ByteBuffer buffer : byteBuffers) {
                //没有要打印的东西了
                if (psize <= 0) {
                    break;
                }
                //要打印的小于或等于buffer容量
                //打印这个buffer中的内容
                System.out.print(new String(buffer.array()).substring(0, buffer.position()));
                //更新要打印的数量
                psize -= buffer.position();
                buffer.clear();
            }
            //TODO:如果收到的数据大小是8的倍数byte 就没办法换行
            if (read < bufferArraySize) {
                System.out.print("\n");
            }
        }
        socketChannel.close();

    }
}

Channel

NIO通道类似于流,但是有以下区别

  • 通道可以同时进行读写,而流只能读或写

  • 通道可以异步读写数据

Channel是NIO中的一个接口,常用Channel有:FileChannel、DatagramChannel、ServerSocketChannel和SocketChannel,FileChannel用于文件的数据读写,DatagramChannel用于UDP的数据读写,ServerSocketChannel和SocketChannel用于TCP的数据读写

连接服务器的时候,server通过ServerSocketChannel给客户端创建一个SocketChannel,客户端通过SocketChannel于服务器通信

FIleChannel实现文件拷贝:

public class NIOFileCopy {
    public static void main(String[] args) throws Exception {
        //01源文件 02目标文件
        String fileName01 = "D:\\study\\untitled2\\file\\file01.txt";
        String fileName02 = "D:\\study\\untitled2\\file\\file02.txt";

        //buffer
        ByteBuffer byteBuffer = ByteBuffer.allocate(10);

        //输入输出channel
        FileInputStream fileInputStream = new FileInputStream(fileName01);
        FileChannel fileInputChannel = fileInputStream.getChannel();

        FileOutputStream fileOutputStream = new FileOutputStream(fileName02);
        FileChannel fileOutputChannel = fileOutputStream.getChannel();

        //buffer较小 需要循环读取 用于记录读取位置
        long position = 0;

        while (true) {
            //每次读取之前都要清空
            byteBuffer.clear();
            //读取
            int read = fileInputChannel.read(byteBuffer,position);
            //读取到末尾 跳出循环
            if (read == -1) {
                break;
            }
            //修改位置
            position += read;
            //buffer反转 转换为写
            byteBuffer.flip();
            //buffer写
            fileOutputChannel.write(byteBuffer);
            //buffer反转 转换为读 准备下次读取
            byteBuffer.flip();
        }
    }
}
public class NIOFileCopy2 {
    public static void main(String[] args) throws Exception {
        //01源文件 02目标文件
        String fileName01 = "D:\\study\\untitled2\\file\\java.jpg";
        String fileName02 = "D:\\study\\untitled2\\file\\java_cpy.jpg";

        //buffer
        ByteBuffer byteBuffer = ByteBuffer.allocate(10);

        //输入输出channel
        FileInputStream fileInputStream = new FileInputStream(fileName01);
        FileChannel fileInputChannel = fileInputStream.getChannel();

        FileOutputStream fileOutputStream = new FileOutputStream(fileName02);
        FileChannel fileOutputChannel = fileOutputStream.getChannel();

        //将输入通道中的内容拷贝到输出流通道 输出通道直接输出
        fileOutputChannel.transferFrom(fileInputChannel,0,fileInputChannel.size());

        fileInputChannel.close();
        fileOutputStream.close();
        fileInputStream.close();
        fileOutputChannel.close();
    }
}

Selector

java的NIO是非阻塞的方式,可以用一个线程处理多个客户端的连接,需要用到Selector选择器。

Selector能够检测多个注册的通道上是否有事件发生(多个Channel以事件的方式可以注册到同一个Selector),如果有事件发生,就获取这个事件进行处理,可以管理多个连接和请求。不是多线程,不需要在多个线程中频繁切换。当线程从客户端Socket通过通道进行读写的时候,若没有数据可用,该线程可以进行其他任务。

Netty的IO线程NioEventLoop聚合了Selector选择器(也叫多路复用器),可以同时并发处理成百上千个客户端连接。

Selector相关方法:
  • select() 阻塞

  • wakeup() 唤醒

  • selectNow() 不阻塞 立马返回

Selector运行过程:
  1. 当客户端连接时,会通过ServerSocketChannel得到SocketChannel

  2. Selector进行监听select(),select()方法返回有事件发生的通道的个数

  3. 将得到的SocketChannel注册到Selector上,register(Selector sel, int ops),一个Selector上可以注册多个SocketChannel

  4. 注册后返回一个SelectionKey,会和该Selector关联(Set)

  5. 进一步得到各个SelectionKey(有事件发生)

  6. 再通过SelectionKey可以反向获取SocketChannel,方法channel()

  7. 可以通过得到的channel,完成业务处理

image-20230109194839765

服务端:

public class NIOServer {
    public static void main(String[] args) throws Exception {
        //创建ServerSocketChannel  类似ServerSocket
        ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
        //得到一个Selector对象
        Selector selector = Selector.open();
        //绑定一个端口 在服务器端监听
        serverSocketChannel.socket().bind(new InetSocketAddress(7771));
        //设置为非阻塞
        serverSocketChannel.configureBlocking(false);
        //把severSocketChannel注册到selector  关心事件为OP_ACCEPT连接事件
        serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
        //循环等待客户端连接
        while (true){
            //等待1s
            if (selector.select(1000) == 0){
                //没有任何事件发生 就返回
                System.out.println("服务器等待1s 无连接");
                continue;
            }
            //如果返回的不是0 就是有事件发生 就获取到SelectionKey集合
            //通过SelectionKey可以反向获取相关的通道
            //使用迭代器遍历通道
            Set<SelectionKey> selectionKeys = selector.selectedKeys();
            Iterator<SelectionKey> keyIterator = selectionKeys.iterator();
            while (keyIterator.hasNext()){
                //获取SelectionKey
                SelectionKey selectionKey = keyIterator.next();
                //根据key对应的通道发生的事件做响应的处理
                if(selectionKey.isAcceptable()){
                    //如果是有新的客户端连接 给该客户端生成SocketChannel
                    SocketChannel socketChannel = serverSocketChannel.accept();
                    //channel必须声明为非阻塞才能注册进selector
                    socketChannel.configureBlocking(false);
                    //将当前的SocketChannel注册到Selector 关注OP_READ读事件 同时给该channel关联一个buffer
                    socketChannel.register(selector,SelectionKey.OP_READ, ByteBuffer.allocate(1024));
                }
                if (selectionKey.isReadable()){
                    //发生事件为OP_READ事件 通过key反向获取channel
                    SocketChannel channel = (SocketChannel) selectionKey.channel();
                    //获取到该channel关联的buffer
                    ByteBuffer buffer = (ByteBuffer) selectionKey.attachment();
                    channel.read(buffer);
                    System.out.println("from client : " + new String(buffer.array()));
                }
                //手动从集合中删除key 防止重复操作
                keyIterator.remove();
            }
        }
    }
}

客户端:

public class NIOClient {
    public static void main(String[] args) throws Exception {
        //得到一个网络通道
        SocketChannel socketChannel = SocketChannel.open();
        //设置非阻塞模式
        socketChannel.configureBlocking(false);
        //提供服务器端的ip:port
        InetSocketAddress inetSocketAddress = new InetSocketAddress("127.0.0.1", 7771);
        //连接服务器
        if (!socketChannel.connect(inetSocketAddress)) {
            while (!socketChannel.finishConnect()) {
                System.out.println("Connecting....(doing other works)");
            }
        }
        //连接成果 发送数据
        String msg = "hello NIOServer";
        //产生一个与字节数组大小相同的buffer
        ByteBuffer buffer = ByteBuffer.wrap(msg.getBytes());
        for (int i = 0; i < 10; i++) {
            //发送数据 将buffer中的数据写入channel
            System.out.println("客户端写入消息 "+msg+" 到channel");
            socketChannel.write(buffer);
        }
        socketChannel.close();
    }
}

Selector、Channel和Buffer的关系

  • 每个channel都对应一个buffer

  • selector会对应一个线程,一个线程可以对应多个channel(连接)

  • 程序切换到哪个channel是由事件决定的

  • selector会根据不同的事件在各个通道上切换

  • buffer就是一个内存块,底层是有一个数据

  • 数据的读取和写入是通过buffer,这个和BIO是有本质不同的。BIO中要么是输入流,要么是输出流;NIO的buffer是可以读也可写的,但是需要 flip() 切换读写状态

  • channel是双向的,可以返回底层操作系统的情况,比如:linux底层操作系统通道就是双向的

零拷贝

零拷贝是网络编程的关键,很多性能优化都离不开零拷贝,java程序中,零拷贝有mmap(内存映射)和sendFile。零拷贝是从操作系统角度看的,是没有CPU拷贝。Linux2.4中避免了从内核缓冲区拷贝到SocketBuffer的操作,直接拷贝到协议栈,从而再减少了数据拷贝,其实有一次CPU拷贝(kernel buffer -> socket buffer),但是信息很少,消耗低,可以忽略。

mmap优化,mmap通过内存映射,将文件映射到内核缓冲区,同时用户可以共享内核空间的数据。这样,在进行网络传输时,就可以减少内核空间到用户空间的拷贝次数。mmap适合小数据量读写,需要4次上下文切换,3次数据拷贝。

sendFile函数,数据根本不经过用户态,直接从内核缓冲区进入到SocketBuffer,同时,由于和用户态完全无关,就减少了一次上下文切换。sendFile适合大文件传输,3次上下文切换,最少两次数据拷贝,sendFile可以使用DMA方式,减少CPU拷贝,mmap则不能(必须从内核拷贝到Socket缓冲区)。

零拷贝是从操作系统的角度上来说的,因为内核缓冲区之间,没有数据是重复的(只有bernel buffer有一份数据),就可以说是零拷贝了。零拷贝不仅带来更少的数据复制,还能带来其他的性能优势,例如更少的上下文切换,更少的CPU缓存伪共享以及无CPU校验和计算。

AIO

JDK7引入了Asynchronous IO(异步IO)。在进行IO编程中,常用到两种模式:Reactor(反应器)和Proactor。java NIO中使用的是Reactor,当有事件触发时,服务器得到通知,进行相应的处理。

AIO是异步不阻塞IO,AIO引入异步通道概念,采用Proactor模式,简化程序编写,有效的请求才启动线程,它的特点是先由操作系统完成之后才通知服务端程序启动线程去处理,一般使用于连接较多且连接时间较长的应用。

AIO还没有广泛应用。

为什么选用Netty

NIO缺点:

  • NIO类库和API繁杂,学习成本高,需要熟练掌握Selector、ServerSocketChannel、SocketChannel、ByteBuffer等;

  • 需要熟悉java多线程,因为NIO编程涉及Reactor模式,必须对多线程和网络编程非常熟悉,才能写出高质量的NIO程序;

  • 臭名昭著的epoll bug,导致Selector空轮询,最终导致CPU利用率100%。

Netty优点:

  • API简单,学习成本低;

  • 功能强大,内置了多种编码解码器,支持多种协议;

  • 性能高,对比其他主流NIO框架,Netty性能最优;

  • 社区活跃,发现BUG会及时修复,迭代版本周期短,不断加入新功能;

  • Dubbo、ES都采用Netty,质量得到验证;

线程模型

不同线程模型对程序的性能有很大影响,目前存在的线程模型有:传统阻塞IO模型,Reactor模式。

根据Reactor的数量和处理资源池线程的数量不同,有3种典型的实现:单Reactor单线程,单Reactor多线程,主从Reactor多线程。

Netty主要基于主从Reactor多线程模型做了一定的改进,主从Reactor多线程模型中有多个Reactor。

Reactor模型

Reactor基于IO复用模型:多个连接共用一个阻塞对象,应用程序只需要在一个阻塞对象等待,无需阻塞等待所有连接。当某个连接有新的数据可以处理时,操作系统通知应用程序,线程从阻塞状态返回,开始进行业务处理。

基于线程池复用线程资源:不必再为每个连接创建线程,将连接完成后的业务处理任务分配给线程进行处理,一个线程可以处理多个连接的业务。

image-20230110134638438

Reactor指,通过一个或多个输入请求同时传递给服务处理器的模式,是基于事件驱动的。服务器端程序处理传入的多个请求,并将他们同步分派到响应的处理线程,因此Reactor模式也叫Dispatch模式。Reactor模式使用了IO复用监听事件,收到事件后分发给某个线程(或进程),这就是网络服务高并发处理关键。

Reactor核心组成部分

Reactor:在一个单独的线程中运行,负责监听和分发事件,分发给适当的处理程序来对IO事件做出反应。

Handlers:处理程序IO事件要完成的实际事件,Reactor通过调度适当的处理程序来响应IO,处理程序执行非阻塞操作。

单Reactor单线程

一个Reactor一个Handler在同一个线程中,在高并发的时候Handler仍然容易阻塞。

单Reactor多线程

Reactor在主线程中运行,如果有事件要处理,就将业务处理转发给线程池中的线程处理。Reactor是单线程运行的,在高并发场景容易出现瓶颈。

主从Reactor多线程

一个Reactor对应多个Reactor子线程

Netty模型

Netty抽象出两种线程池BossGroup和WorkerGroup,BossGroup专门负责接收客户端的连接,WorkerGroup专门负责网络读写。

BossGroup和WorkerGroup类型都是NioEventLoopGroup,NioEventLoopGroup相当于一个事件循环组,这个组中含有多个事件循环,每个事件循环都叫NioEventLoop,NioEventLoop表示一个不断循环的执行处理任务的线程,每个NioEventLoop都有一个Selector,用于监听绑定在其上的Socket的网络通讯,NioEventLoopGroup可以有多个线程,即可以含有多个NioEventLoop。

每个Boss NioEventLoop循环执行的步骤有3步:轮询accept事件;处理accept事件,与client建立连接,生成NioSocketChannel,并将其注册到某个Worker NioEventLoop上的selector;处理任务队列的任务,即runAllTasks。

每个Worker NioEventLoop循环执行的步骤:轮询read或write事件;处理IO事件,在对应的NioSocketChannel处理;处理任务队列的其他任务,即runAllTasks。

每个Worker NioEventLoop处理业务时,会使用pipeline(管道),pipeline中包含了channel,即通过pipeline可以获取对应的通道,管道中维护了很多处理器

image-20230110142437424

架构图

image-20230106141240267

绿色核心部分包括零拷贝、API库、可扩展事件模型。

橙色Protocol Support协议支持,包括Http协议、webSocket、SSL、谷歌Protobu协议、zlib/gzip压缩与解压缩、Large File Transfer大文件传输等。

红色部分Transport Service传输服务,包括Socket、Datagram、Http Tunnel等。

以上可以看出Netty的功能、协议、传输方式都比较全,比较强大

Hello Word

引入依赖

<!-- https://mvnrepository.com/artifact/io.netty/netty-all -->
<dependency>
    <groupId>io.netty</groupId>
    <artifactId>netty-all</artifactId>
    <version>4.1.42.Final</version>
</dependency>

NettyClient

public class NettyClient {
    public static void main(String[] args) throws Exception {
        //客户端只需要一个事件循环组
        EventLoopGroup group = new NioEventLoopGroup();

        try {
            //创建客户端启动对象
            //客户端使用的不是ServerBootStrap 而是 BootStra
            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.out.println("Client is ready.");

            //启动客户端去连接服务器
            //关于ChannelFuture要分析 涉及到netty的异步模型
            ChannelFuture channelFuture = bootstrap.connect("127.0.0.1", 7770).sync();

            //给关闭通道进行监听
            channelFuture.channel().closeFuture().sync();
        }finally {
            group.shutdownGracefully();
        }

    }
}

NettyClientHandler

public class NettyClientHandler extends ChannelInboundHandlerAdapter {

    /**
     * 通道就绪 就会触发该方法
     * @param ctx
     * @throws Exception
     */
    @Override
    public void channelActive(ChannelHandlerContext ctx) throws Exception {
//        System.out.println(ctx);
        ctx.writeAndFlush(Unpooled.copiedBuffer("Hello Server, this is Client.", CharsetUtil.UTF_8));
    }

    /**
     *
     * @param ctx 当通道有读取事件时 会触发
     * @param msg
     * @throws Exception
     */
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        ByteBuf buf = (ByteBuf) msg;
        System.out.println("Server msg:"+buf.toString(CharsetUtil.UTF_8));
        System.out.println("Server address:" + ctx.channel().remoteAddress());
    }

    /**
     * 发生异常
     * @param ctx
     * @param cause
     * @throws Exception
     */
    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        cause.printStackTrace();
        ctx.close();
    }
}

NettyServer

public class NettyServer {
    public static void main(String[] args) throws Exception{
        //创建BossGroup和WorkerGroup
        //创建了两个线程组 bossGroup只是处理连接请求 workerGroup完成真正的和客户端的业务处理
        //这两个都是无限循环
        //bossGroup和workerGroup 含有的子线程(NioEventLoop)个数默认是CPU核数*4
        EventLoopGroup bossGroup = new NioEventLoopGroup();
        EventLoopGroup workerGroup = new NioEventLoopGroup();

        try {
            //创建服务器端的启动对象 配置启动参数
            ServerBootstrap bootstrap = new ServerBootstrap();
            bootstrap.group(bossGroup, workerGroup)//设置两个线程组
                    .channel(NioServerSocketChannel.class)//使用NioServerSocketChannel作为服务器的通道实现
                    .option(ChannelOption.SO_BACKLOG, 128)//设置线程队列等待连接的个数
                    .childOption(ChannelOption.SO_KEEPALIVE, true)//设置保持活动连接状态
                    .childHandler(new ChannelInitializer<SocketChannel>() {//使用匿名对象方式创建一个通道测试对象
                        //给pipeline设置处理器
                        @Override
                        protected void initChannel(SocketChannel socketChannel) throws Exception {
                            socketChannel.pipeline().addLast(new NettyServerHandler());
                        }
                    });//给workerGroup的EventLoop对应的管道设置处理器

            System.out.println("Server is ready.");

            //绑定一个端口并且同步,生成了一个ChannelFuture对象
            //启动服务器并绑定端口
            ChannelFuture channelFuture = bootstrap.bind(7770).sync();

            //关闭对通道进行监听
            channelFuture.channel().closeFuture().sync();
        } finally {
            //优雅地关闭
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }

    }
}

NettyServerHandler

/**
 * 我们自定义一个Handler 需要集成Netty规定好的某个HandlerAdapter
 * 这时我们自定义的Handler才能称之为一个Handler
 */
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.out.println(ctx);
        //将msg转为ByteBuffer 不是NIO的ByteBuffer
        ByteBuf buf = (ByteBuf) msg;
        System.out.println("client msg:" + buf.toString(CharsetUtil.UTF_8));
        System.out.println("client address:"+ctx.channel().remoteAddress());
    }

    /**
     * 数据读取完毕
     * @param ctx
     * @throws Exception
     */
    @Override
    public void channelReadComplete(ChannelHandlerContext ctx) throws Exception {
        //write()加flush() 将数据写入到缓冲 并刷新
        //一般来讲 要将数据编码后再发送
        ctx.writeAndFlush(Unpooled.copiedBuffer("Hello Client, this is Server.", CharsetUtil.UTF_8));
    }

    /**
     * 处理异常 关闭通道
     * @param ctx
     * @param cause
     * @throws Exception
     */
    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        cause.printStackTrace();
        ctx.close();
    }
}

任务队列

任务队列的三种典型场景:

  1. 用户程序自定义的普通任务 提交到taskQueue

  2. 用户自定义定时任务 提交到scheduleTaskQueue

  3. 非当前Reactor线程调用Channel的各种方法

普通任务

@Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        ctx.channel().eventLoop().execute(new Runnable() {
            @Override
            public void run() {
                try {
                    Thread.sleep(10 * 1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        });
        ctx.writeAndFlush(Unpooled.copiedBuffer("Hello Server, this is Client. 10s after", CharsetUtil.UTF_8));
    }

定时任务

 /**
     * 向scheduleTaskQueue提交定时任务
     * @param ctx
     * @param msg
     * @throws Exception
     */
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        ctx.channel().eventLoop().schedule(new Runnable() {
            @Override
            public void run() {
                ctx.writeAndFlush(Unpooled.copiedBuffer("This is a scheduled task. The scheduled time is five seconds", CharsetUtil.UTF_8));
            }
        }, 5, TimeUnit.SECONDS);
    }

异步模型

异步过程调用发出后,调用者不能立刻得到结果,在实际处理这个调用的组件完成后,通过状态、通知和回调来通知调用者。Netty中的IO是异步的,包括Bind,Write,Connect等操作会简单的返回一个ChannelFuture。调用者并不能立刻获取结果,而是通过Future-Listener机制,用户可以方便的主动获取或通过机制获取IO操作的结果。Netty的异步模型就是建立在future和callback的之上的,callback就是回调。Future的核心思想是:假设一个方法fun,计算过程可能非常耗时,等待fun返回显然是不合适的。那么可以在调用fun的时候,立马返回一个Future,后续可以通过Future去监控方法fun的处理过程,即Future-Listener机制。

Future说明:

  • 表示异步的执行结果,可以通过它提供的方法来检测执行是否完成,比如检索计算等等

  • ChannelFuture是一个接口,集成了Future,我们可以添加监听器,当监听的时间发生时,就会通知监听器。

Future-Linster机制:

当Future对象刚刚创建时,处于非完成状态,调用者可以通过返回的ChannelFuture来获取操作执行的状态,注册监听函数来执行完成后的操作。常见的操作如下:

  • 通过isDone方式来判断当前操作是否完成

  • 通过isSuccess方法来判断已完成的当前操作是否成功

  • 通过getCause方法来获取已完成的当前操作失败的原因

  • 通过isCancelled方法来判断已完成的当前操作是否被取消

  • 通过addListener方法来注册监听器,当操作已完成(isDone方法返回完成),将会通知指定的监听器;如果Future对象已完成,则通知指定的监听器

HTTP服务

HttpServer:

public class HttpServer {
    public static void main(String[] args) throws Exception {
        //创建bossGroup和workerGroup
        EventLoopGroup bossGroup = new NioEventLoopGroup();
        EventLoopGroup workerGroup = new NioEventLoopGroup();

        try {
            //创建启动服务器
            ServerBootstrap bootstrap = new ServerBootstrap();
            bootstrap.group(bossGroup, workerGroup)//设置两个线程组
                    .channel(NioServerSocketChannel.class)//使用NioServerSocketChannel作为服务器的通道实现
                    .childHandler(new HttpServerInitializer());//给workerGroup的EventLoop对应的管道设置处理器
            //绑定端口
            ChannelFuture channelFuture = bootstrap.bind(7771).sync();
            channelFuture.channel().closeFuture().sync();
        }
        finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }
    }
}

HttpServerHandler:

/**
 * SimpleChannelInboundHandler 继承了 ChannelInboundHandlerAdapter
 * HttpObject 客户端和服务器端相互通讯的数据被封装成 HttpObject
 */
public class HttpServerHandler extends SimpleChannelInboundHandler<HttpObject> {

    //读取客户端数据
    @Override
    protected void channelRead0(ChannelHandlerContext ctx, HttpObject httpObject) throws Exception {
        //判断httpObject是不是一个httpRequest请求
        if (httpObject instanceof HttpRequest) {
            System.out.println("client address:" + ctx.channel().remoteAddress());

            //回复信息给浏览器
            ByteBuf content = Unpooled.copiedBuffer("hello, this is server!", CharsetUtil.UTF_8);

//            //过滤请求 过滤图标请求
//            HttpRequest httpRequest = (HttpRequest) httpObject;
//            //获取URI
//            URI uri = new URI(httpRequest.uri());
//            if ("/favicon.ico".equals(uri.getPath())) {
//                System.out.println("client request favicon.ico, do nothing");
//                return;
//            }

            //构造一个http响应
            FullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, content);
            response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/plain");
            response.headers().set(HttpHeaderNames.CONTENT_LENGTH, content.readableBytes());

            //将构建好的response返回
            ctx.writeAndFlush(response);
        }

    }
}

HttpServerInitializer:

public class HttpServerInitializer extends ChannelInitializer<SocketChannel> {
    @Override
    protected void initChannel(SocketChannel socketChannel) throws Exception {

        //向管道加入处理器
        //得到管道
        ChannelPipeline pipeline = socketChannel.pipeline();
        //加入一个netty提供的httpServerCodec
        //HttpServerCodec是Netty提供的处理http的编码解码器
        pipeline.addLast(new HttpServerCodec());
        //增加一个自定义的Handler
        pipeline.addLast(new HttpServerHandler());

    }
}

Netty核心模块

Bootstrap&ServerBootstrap

Netty程序通常是由一个Bootstrap开始,主要作用是配置整个Netty程序,串联各个组件。Netty中Bootstrap类是客户端启动引导类,ServerBootstrap是服务端启动引导类。

group(),方法用于设置线程组,服务端设置一个boss一个worker,客户端设置一个,用EventLoopGroup

channel(),用来设置一个服务器端的通道实现

option(),用于给ServerChannel添加配置

childOption(),用来给接收到的通道添加配置

handler(),用来设置handler

childHandler(),用来设置业务处理类,服务端使用handler()是在bossGroup生效,使用childHandler()是在workerGroup生效

bind(),用来设置端口号

connect(),用来连接服务端

Future&ChannelFuture

Netty中的IO操作都是异步的,不能立刻得知消息是否被正确的处理,但是可以过一会等它执行完成或者直接注册一个监听,具体实现就是通过Future和ChannelFuture,他们可以注册一个监听,当执行成功或者失败时会监听会自动触发注册的监听事件。

Channel channel(),返回当前正在进行的IO操作的通道

ChannelFuture sync(),等待异步操作执行完毕

Channel

Netty的网络通信组件,用于执行网络的IO操作,通过Channel可以获取当前网络连接的通道状态,可以获取网络连接的配置参数。Channel提供异步的网络IO操作,异步意味着任何IO调用都将立即返回,并且不保证在调用结束时所请求的IO操作已完成。调用立即返回一个ChannelFuture实例,通过注册监听器到ChannelFuture上,可以IO操作成功、失败或取消时回调通知调用方。Channel支持关联IO操作域对应的处理程序,不同协议、不同阻塞类型的连接都有不同的Channel类型与之对应,常用的Channel类型:

  • NioSocketChannel,异步的客户端TCP Socket连接

  • NioServerSocketChannel,异步的服务器端TCP Socket连接

  • NioDatagramChannel,异步的UDP连接

  • NioSctpChannel,异步的客户端Sctp连接

  • NioSctpServerChannel,异步的Sctp服务器端连接,这些通道涵盖了UDP和TCP网络IO以及文件IO

Selector

Netty基于Selector对象实现IO多路复用,通过Selector一个线程可以监听多个连接的Channel事件。当向一个Selector中注册Channel后,Selector内部的机制就可以自动不断地查询这些注册的Channel是否有已就绪的IO事件,这样程序就可以很简单地使用一个线程高效地管理多个Channel

ChannelHandler

ChannelHandler是一个接口,处理IO事件或者拦截IO事件,并将其转发到其ChannelPipeline(业务处理链)中的下一个处理程序。ChannelHandler本身并没有提供很多方法,因为这个接口有许多方法需要实现,方便使用期间,可以继承它的子类。

ChannelInboundHandler用于处理入站IO事件,ChannelOutboundHandler用于处理出战IO操作。

ChannelInboundHandlerAdapter用于处理入栈IO事件,ChannelOutboundHandlerAdapter用于处理出战IO事件,ChannelDuplexHandler用于处理入站和出战事件。

ChannelPipeline提供了ChannelHandler链的容器,如果事件的运动方向是从客户端到服务端的,那么我们称这些事件为出战,即客户端发送给服务端的数据会通过pipeline中的一系列ChannelOutboundHandler,并被这些Handler处理,反之称为入站。

channelActive(),通道就绪事件,channelInactive(),通道非活动状态

channelRead(),通道读取数据事件,channelReadComplete(),通道读取数据完毕

handlerAdded(),handler注册事件,handlerRemoved(),handler移除事件

Pipeline&ChannelPipeline

ChannelPipeline是一个Handler的集合,它负责处理和拦截 inbound 或者 outbound 的事件和操作,相当于一个贯穿Netty的链。ChannelPipeline是保存ChannelHandler的List,用于处理或拦截Channel的入站事件和出战操作。ChannelPipeline实现了一种高级形式的拦截过滤器模式,使用户可以完全控制事件的处理方式,以及Channel中各个的ChannelHandler如何交互。

一个Channel包含了一个ChannelPipeline,而ChannelPipeline中又维护了一个由ChannelHandlerContext组成的双向链表,并且每个ChannelHandlerContext中又关联着一个ChannelHandler。

入站事件和出战事件在一个双向链表中,入站事件会从链表head往后传递到最后一个入站的handler,出战事件会从链表tail往前传递到最前一个出战的handler,两种类型的handler互不干扰。

addFirst(),把一个业务处理类(handler)添加到链中的第一个位置,addLast(),把一个业务处理类(handler)加到链中的最后一个位置。

ChannelHandlerContext

保存Channel相关的所有上下文信息,同时关联一个ChannelHandler对象,即ChannelHandlerContext中包含一个具体的事件处理器ChannelHandler,同时ChannelHandlerContext中也绑定了对应的pipeline和Channel的信息,方便对ChannelHandler进行说明。

close(),关闭通道,flush(),刷新,writeAndFlush(),将数据写到ChannelPipeline中当前ChannelHandler的下一个ChannelHandler开始处理(出站)。

ChannelOption

Netty在创建Channel实例后,一般都需要设置ChannelOption参数。

ChannelOption.SO_BACKLOG:对应TCP/IP协议的listen函数中的backlog参数,用来初始化服务器可连接队列大小。服务端处理客户端连接请求时顺序处理的,所以同一时间只能处理一个客户端连接。多个客户端来的时候,服务端将不能处理的客户端连接请求放在队列中等待处理,backlog参数指定了队列的大小。

ChannelOption.SO_KEEPALIVE:一直保持连接活动状态

EventLoopGroup&NioEventLoopGroup

EventLoopGroup是一组EventLoop的抽象,Netty为了更好的利用多核CPU资源,一般会有多个EventLoop同时工作,每个EventLoop维护着一个Selector实例。

EventLoopGroup提供next接口,可以从组里面按照一定规则获取其中一个EventLoop来处理任务。在Netty服务器端编程中,我们一般都需要提供两个EventLoopGroup,boss和worker。

shutdownGracefully(),断开连接,关闭线程

Unpooled类

Netty提供一个专门用来操作缓冲区(即Netty的数据容器)的工具类

copiedBuffer(),通过给定的数据和字符编码返回一个ByteBuf对象,类似NIO中的ByteBuffer,但有区别

Netty心跳

WebSocket长连接

WebSocket建立在TCP协议之上,没有同源限制,客户端可以于任意服务器通信。

WSServer:

public class WSServer {
    public static void main(String[] args) throws Exception {
        EventLoopGroup bossGroup = new NioEventLoopGroup();
        EventLoopGroup workerGroup = new NioEventLoopGroup();
        try {
            ServerBootstrap serverBootstrap = new ServerBootstrap();

            serverBootstrap.group(bossGroup,workerGroup);
            serverBootstrap.channel(NioServerSocketChannel.class);
            serverBootstrap.childHandler(new ChannelInitializer<SocketChannel>() {
                @Override
                protected void initChannel(SocketChannel socketChannel) throws Exception {
                    ChannelPipeline pipeline = socketChannel.pipeline();

                    //因为是基于http协议,使用http编解码器
                    pipeline.addLast(new HttpServerCodec());
                    //是以块方式写,添加ChunkedWriteHandler处理器
                    pipeline.addLast(new ChunkedWriteHandler());

                    //因为http数据在传输过程中是分段的,HttpObjectAggregator可以将多个段聚合起来
                    pipeline.addLast(new HttpObjectAggregator(8192));

                    //websocket数据是以帧的形式传递的 用来指定uri
                    //WebSocketServerProtocolHandler将http协议升级为ws协议 保持长连接
                    pipeline.addLast(new WebSocketServerProtocolHandler("/ws"));

                    pipeline.addLast(new TextWebSocketFrameHandler());
                }
            });

            ChannelFuture channelFuture = serverBootstrap.bind(7772).sync();
            channelFuture.channel().closeFuture().sync();
        }finally {
            bossGroup.shutdownGracefully();
            workerGroup.shutdownGracefully();
        }
    }
}

TextWebSocketFrameHandler:

//TextWebSocketFrame表示一个文本帧(Frame)
public class TextWebSocketFrameHandler extends SimpleChannelInboundHandler<TextWebSocketFrame> {

    @Override
    protected void channelRead0(ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception {
        System.out.println("服务器收到消息:" + msg.text());
        //回复消息
        ctx.channel().writeAndFlush(new TextWebSocketFrame("服务器时间" + new Date() + " " + msg.text()));
    }

    //当web客户端连接后 触发方法
    @Override
    public void handlerAdded(ChannelHandlerContext ctx) throws Exception {
        //id 表示唯一的值 LongText 是唯一的 ShortText 不是唯一的
        System.out.println("handlerAdded " + ctx.channel().id().asLongText());
        System.out.println("handlerAdded " + ctx.channel().id().asShortText());
    }

    @Override
    public void handlerRemoved(ChannelHandlerContext ctx) throws Exception {
        System.out.println("handlerRemoved "+ctx.channel().id().asLongText());
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        System.out.println("Exception:" + cause.getMessage());
        ctx.close();
    }
}

ws.html:

<!DOCTYPE html>
<html lang="en">
<head>
    <meta charset="UTF-8">
    <title>Title</title>
</head>
<body>
<script>
    var socket;
    //判断当前浏览器是否支持webSocket
    if (window.WebSocket) {
        //go on
        socket = new WebSocket("ws://localhost:7772/ws");
        //相当于channelRead0,ev 收到服务器端回送
        socket.onmessage = function (ev) {
            var rt = document.getElementById("responseText");
            rt.value = rt.value + "\n" + ev.data;
        }
        //相当于连接开启
        socket.onopen = function (ev) {
            var rt = document.getElementById("responseText");
            rt.value = "连接开启了。。";
        }
        //连接关闭
        socket.onclose = function (ev) {
            var rt = document.getElementById("responseText");
            rt.value = rt.value + "\n" + "连接关闭了。。";
        }
    } else {
        alert("当前浏览器不支持webSocket");
    }
​
    //发送消息到服务器
    function send(message) {
        if (!window.socket) {
            return;
        }
        if (socket.readyState == WebSocket.OPEN) {
            //发送消息
            socket.send(message);
        } else {
            alert("连接没有开启");
        }
    }
​
</script>
​
<form onsubmit="return false">
    <textarea name="message" style="height: 300px; width: 300px"></textarea>
    <input type="button" value="发送消息" onclick="send(this.form.message.value)">
    <textarea id="responseText" style="height: 300px; width: 300px"></textarea>
    <input type="button" value="清空内容" onclick="document.getElementById('responseText').value=''">
</form>
</body>
</html>


评论