ServerSocketChannel.open()创建一个服务器的socket,并且建立管道Channel。这个方式跟我们一个开始使用new ServerSocket(8080);没有太大区别。

public class NioTcpDemo01 {
    public static void main(String[] args) throws Exception{
        ServerSocketChannel ssc = ServerSocketChannel.open();
        ssc.bind(new InetSocketAddress(8080));
        SocketChannel accept = ssc.accept();
        System.out.println(accept);
    }
}

⬇接着是接收客户端发来的消息,跟前面没有用Channel是一个性质的。

public class NioTcpDemo01 {
    public static void main(String[] args) throws Exception{
        ByteBuffer buffer = ByteBuffer.allocate(1024);
        ServerSocketChannel ssc = ServerSocketChannel.open();
        ssc.bind(new InetSocketAddress(8080));
        SocketChannel accept = ssc.accept();
        System.out.println(accept);
        accept.read(buffer);
        System.out.println(new String(buffer.array()));
    }
}

⬇设置为非阻塞模式,充分利用系统的函数进行优化

public class NioTcpDemo01 {
	//设置为非阻塞
    public static void main(String[] args) throws Exception {
        ServerSocketChannel ssc = ServerSocketChannel.open();
        ssc.bind(new InetSocketAddress(8080));
        ssc.configureBlocking(false);
        while(true){
            SocketChannel accept = ssc.accept();//重点为非阻塞
            if(accept != null)
            	System.out.println(accept);
        }
    }
}

⬆设置为非阻塞之后accept()不再阻塞,如果没有连接进来,会返回null。通过主线程while循环进行不断轮询可以知道是否有连接进来。

将SocketChannel也设置为非阻塞。即当客户端连接进来时不阻塞read事件.

public static void main(String[] args) throws Exception {
        ByteBuffer buffer = ByteBuffer.allocate(1024);

        ServerSocketChannel ssc = ServerSocketChannel.open();
        ssc.bind(new InetSocketAddress(8080));
        ssc.configureBlocking(false);
        List<SocketChannel> socketChannels = new ArrayList<>();
        while(true){
            SocketChannel accept = ssc.accept();//重点为非阻塞
            if(accept != null){
                System.out.println(accept);
                accept.configureBlocking(false); //设置SocketChannel为非阻塞
                socketChannels.add(accept);
            }

            for (SocketChannel socketChannel : socketChannels){

                int read = socketChannel.read(buffer);//这里不阻塞了,如果读到数据,则返回数据长度,否则返回0
                if (read > 0){
                    buffer.flip();
                    byte[] data = new byte[read];
                    buffer.get(data);
                    System.out.println(new String(data));
                    buffer.clear();
                    socketChannel.write(ByteBuffer.wrap(data));
                }
                if (read == -1) { //read==-1就断开连接了
                    // 客户端断开连接
                    System.out.println("客户端断开连接: " + socketChannel);
                    socketChannel.close();
                    socketChannels.remove(socketChannel);
                    continue;
                }
            }
        }
    }

selector和channel建立连接

public class NioTcpDemo02 {
    public static void main(String[] args) throws Exception {
        ServerSocketChannel ssc = ServerSocketChannel.open();
        ssc.bind(new InetSocketAddress(8080));
        ssc.configureBlocking(false);//设置为非阻塞

        //创建多路复用
        Selector selector = Selector.open();
        //ssc注册到多路复用器,监听连接事件,selector管理和监听这个Channel
        ssc.register(selector, SelectionKey.OP_ACCEPT);
        while(true){
            int select = selector.select(); //阻塞,等待事件触发
            if(select > 0){
                Set<SelectionKey> selectionKeys = selector.selectedKeys(); //获取所有事件
                Iterator<SelectionKey> iterator = selectionKeys.iterator();
                while (iterator.hasNext()){
                    SelectionKey key = iterator.next();
                    if(key.isAcceptable()){//判断是否是连接事件
                        //连接事件的channel必定是ServerSocketChannel
                        ServerSocketChannel channel = (ServerSocketChannel) key.channel();
                        //接受连接
                        SocketChannel socketChannel = channel.accept();
                        System.out.println(socketChannel);
                        socketChannel.configureBlocking(false);
                        //socketChannel也交给Selector管理和监听
                        socketChannel.register(selector,SelectionKey.OP_READ);
                        System.out.println("连接成功");
                        iterator.remove(); // 时间处理完成移除事件
                    }
                    if(key.isReadable()){
                        SocketChannel socketChannel = (SocketChannel) key.channel();
                        byte[] buffer = new byte[1024];
                        int read = socketChannel.read(ByteBuffer.wrap(buffer));
                        if(read > 0){
                            System.out.println(new String(buffer));
                        }
                        socketChannel.write(ByteBuffer.wrap(buffer));
                        iterator.remove();
                    }
                }
            }
        }
    }
}

boss-worker

但实际上我们只用到了一个线程,一个selector,不能充分使用多核cpu的优势。现在我们创建多个selector。我们规定boss-selector只负责建立连接,worker-selector负责业务处理。每次连接创建一个线程,这个线程中创建worker-selector,这个线程中的selector负责处理其他事件。而且规定创建的worker数量。

连接事件->boss
读写事件->worker

假设worker最大数量=4个
当一个连接进来,我们创建一个worker,有4个连接进来,创建4个worker后,第5个连接进来时会复用原来的worker

多Reactor