JAVA的AIO与NIO编程代码实现

本文的代码参考了Tomcat的NIO实现NioEndpoint以及AIO实现Nio2Endpoint代码逻辑。



上图是一个AIO模型。

下面我们通过代码判断其①②③线程是否相同

AIO服务端AsyncServer

public class AsyncServer {
    private ThreadPoolExecutor acceptorExecutor;  // 连接处理线程池
    private ThreadPoolExecutor callBackExecutor;  // 回调处理线程池
    private AsynchronousServerSocketChannel serverSocketChannel;

    public static void main(String[] args) throws Exception{
        AsyncServer asyncServer = new AsyncServer();
        asyncServer.listen();
        TimeUnit.SECONDS.sleep(500);  // 由于方法都是异步的,避免主程序执行完结束
    }

    private void listen() {
        AtomicInteger atomicInteger = new AtomicInteger(0);
        AtomicInteger atomic = new AtomicInteger(0);
        acceptorExecutor = new ThreadPoolExecutor(5, 10
                , 500, TimeUnit.SECONDS
                , new LinkedBlockingQueue<>(10),
                r -> {
                    Thread thread = new Thread(r);
                    thread.setName("acceptorExecutor" + atomicInteger.getAndIncrement());
                    return thread;
                });
        callBackExecutor = new ThreadPoolExecutor(5, 10
                , 500, TimeUnit.SECONDS
                , new LinkedBlockingQueue<>(10),
                r -> {
                    Thread thread = new Thread(r);
                    thread.setName("callBackExecutor" + atomic.getAndDecrement());
                    return thread;
                });
        try {
            AsynchronousChannelGroup channelGroup = AsynchronousChannelGroup.withThreadPool(callBackExecutor);
            serverSocketChannel = AsynchronousServerSocketChannel.open(channelGroup);   // 设置内核完成的处理回调线程池
            InetSocketAddress addr = new InetSocketAddress("0.0.0.0", 8088);
            serverSocketChannel.bind(addr);
            serverSocketChannel.accept(this, new AcceptorHandler());   // 该方法为异步方法,通过底层javadoc可以获悉
            System.out.println("accept end");   // 此命令的打印佐证了上面的方法为异步方法,不会像BIO或者NIO那样阻塞。
        } catch (IOException ex) {
            System.out.println(" listen failed");
        }
    }

    class AcceptorHandler implements CompletionHandler<AsynchronousSocketChannel, AsyncServer> {

        @Override
        public void completed(AsynchronousSocketChannel channel, AsyncServer attachment) {
            attachment.serverSocketChannel.accept(AsyncServer.this, this);   // 这个方法是异步方法,请求进来的时候需要处理请求的同时要再次调用来持续监听,如果不调用的话,则只能接收一个请求
            AsyncServer.this.handleWithExecutor(channel);
            System.out.println(Thread.currentThread().getName() + " submit to executor");
        }

        @Override
        public void failed(Throwable exc, AsyncServer attachment) {
            System.out.println(" accept failed");
        }
    }

    private void handleWithExecutor(AsynchronousSocketChannel socketChannel) {
        acceptorExecutor.execute(new SocketProcessor(socketChannel));
    }

    class SocketProcessor implements Runnable {
        private AsynchronousSocketChannel socketChannel;

        private CompletionHandler<Integer, ByteBuffer> readHandler =  new CompletionHandler<Integer, ByteBuffer>() {

            @Override
            public void completed(Integer result, ByteBuffer attachment) {
                attachment.flip();
                String input = new String(attachment.array()).trim();
                if (input.length() == 0) {
                    ByteBuffer buffer = ByteBuffer.allocate(1024);
                    socketChannel.read(buffer, buffer, readHandler);    // 客户端数据可能没用准备好,再次调用下
                    return;
                }
                System.out.println(Thread.currentThread().getName() + " 收到客户端消息: " + input);
                attachment.clear();
                attachment.put((input + " can do!").getBytes());
                attachment.flip();
                socketChannel.write(attachment);
            }

            @Override
            public void failed(Throwable exc, ByteBuffer attachment) {
                System.out.println("read failed");
            }
        };

        public SocketProcessor(AsynchronousSocketChannel socketChannel) {
            this.socketChannel = socketChannel;
        }

        @Override
        public void run() {
            ByteBuffer buffer = ByteBuffer.allocate(1024);
            socketChannel.read(buffer, buffer, readHandler);
            System.out.println(Thread.currentThread().getName() + " SocketProcessor run end");   // 这里打印submit to executor后的线程名
        }
    }
}

AIO客户端:

public class AsyncClient {
    private AsynchronousSocketChannel asc;
    private static ThreadPoolExecutor poolExecutor;

    public static void main(String[] args) throws Exception{
        AtomicInteger atomic = new AtomicInteger(0);
        poolExecutor = new ThreadPoolExecutor(5, 10
                , 10, TimeUnit.SECONDS
                , new LinkedBlockingQueue<>(1),
                r -> {
                    Thread thread = new Thread(r);
                    thread.setName("poolExecutor " + atomic.getAndIncrement());
                    return thread;
                });
        new Thread(() -> runOneClient("HONOR")).start();
        new Thread(() -> runOneClient("MEIZU")).start();
        new Thread(() -> runOneClient("HUAWEI")).start();
        new Thread(() -> runOneClient("VIVO")).start();
    }

    public AsyncClient() throws Exception{
        asc = AsynchronousSocketChannel.open(AsynchronousChannelGroup.withThreadPool(poolExecutor));  // 设置回调处理的线程池
    }

    private Future<Void> getSocket() {
        return asc.connect(new InetSocketAddress("127.0.0.1", 8088));
    }

    public static void runOneClient(String content) {
        try {
            AsyncClient asyncClient = new AsyncClient();
            Future<Void> socket = asyncClient.getSocket();
            while (!socket.isDone()) {
                TimeUnit.MILLISECONDS.sleep(100);
            }
            asyncClient.writeToServer(content);
            while (asyncClient.asc.isOpen()) {
                TimeUnit.MILLISECONDS.sleep(100);
            }
        } catch (Exception ex) {
            System.out.println( "something wrong occurs when runOneClient");
        }
    }

    private void writeToServer(String content) {
        ByteBuffer buffer = ByteBuffer.allocate(1024);
        asc.write(buffer, buffer, new CompletionHandler<Integer, ByteBuffer>() {
            @Override
            public void completed(Integer result, ByteBuffer attachment) {
                attachment.clear();
                attachment.put(content.getBytes());
                attachment.flip();
                asc.write(attachment);
                attachment.clear();
                readFromServer();
            }

            @Override
            public void failed(Throwable exc, ByteBuffer attachment) {
                System.out.println("writeToServer failed");
            }
        });
    }

    private void readFromServer() {
        ByteBuffer buffer = ByteBuffer.allocate(1024);
        asc.read(buffer, buffer, new CompletionHandler<Integer, ByteBuffer>() {
            @Override
            public void completed(Integer result, ByteBuffer attachment) {
                attachment.flip();
                String input = new String(attachment.array()).trim();
                attachment.clear();
                System.out.println(Thread.currentThread().getName() + " 收到服务端消息: " + input);
                try {
                    asc.close();
                } catch (IOException ex) {
                    System.out.println(" something wrong occurs when readFromServer");
                }
            }

            @Override
            public void failed(Throwable exc, ByteBuffer attachment) {
                System.out.println("readFromServer failed");
            }
        });
    }
}

NIO的服务端NioSocketServer代码:

public class NioSocketServer {
    private static ThreadPoolExecutor poolExecutor;

    private Selector selector;

    private Map<SelectionKey, SocketHandler> map = new HashMap<>();

    public static void main(String[] args) throws Exception{
        AtomicInteger atomic = new AtomicInteger(0);
        poolExecutor = new ThreadPoolExecutor(5, 10
                , 10, TimeUnit.SECONDS
                , new LinkedBlockingQueue<>(1),
                r -> {
                    Thread thread = new Thread(r);
                    thread.setName("poolExecutor " + atomic.getAndIncrement());
                    return thread;
                });
        NioSocketServer nioSocketServer = new NioSocketServer();
        nioSocketServer.listen();
    }

    private void listen() {
        try {
            selector = Selector.open();
            ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
            serverSocketChannel.bind(new InetSocketAddress("0.0.0.0", 8088));
            serverSocketChannel.configureBlocking(true);   // 设置监听阻塞,避免没有请求到来时线程空转
            poolExecutor.execute(() -> {
                while (true) {
                    try {
                        SocketChannel socketChannel = serverSocketChannel.accept();
                        System.out.println(Thread.currentThread().getName() + " accept ");
                        if (socketChannel != null) {
                            socketChannel.configureBlocking(false);   // 设置连接为非阻塞的,只有非阻塞的才可以注册到selector上去,这个很好理解,因为selector需要轮询监听,没用数据时需要立即返回
                            socketChannel.register(selector, SelectionKey.OP_READ | SelectionKey.OP_WRITE, null);
                        }
                    } catch (Exception ex) {
                        System.out.println("something wrong occurs when accept");
                    }
                }
            });
            selectorHandle();
        } catch (Exception ex) {
            System.out.println("something wrong occurs when listen");
        }
    }

    private void selectorHandle() {
        try {
            while (true) {
                int selectNum = selector.select(50);//  select方法是阻塞的,会占用锁,导致无法Register,这里调用此方法,具体的可以点进去看底层代码实现。
                if (selectNum == 0) {
                    continue;
                }
                Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
                while (iterator.hasNext()) {
                    SelectionKey key = iterator.next();
                    synchronized (key) {
                        SocketHandler socketHandler = map.get(key);
                        if (socketHandler == null) {
                            socketHandler = new SocketHandler(key);
                            map.put(key, socketHandler);
                        }
                        iterator.remove();
                        poolExecutor.execute(socketHandler);
                    }
                }
            }
        } catch (Exception ex) {
            System.out.println("something wrong occurs when selectorHandle");
        }
    }

    private void cancelKey(SelectionKey selectionKey, SocketChannel channel) {
        try {
            selectionKey.channel();
            channel.close();
        } catch (Exception ex) {
            System.out.println(" cancelKey occurs exception");
        }
    }

    class SocketHandler implements Runnable{
        private SelectionKey selectionKey;
        private SocketChannel channel;

        public SocketHandler(SelectionKey selectionKey) {
            this.selectionKey = selectionKey;
            channel = (SocketChannel) selectionKey.channel();
        }

        @Override
        public void run() {
            synchronized (selectionKey) {
                try {
                    if (!(selectionKey.isValid() && selectionKey.isReadable() && channel.isOpen())) {
                        return;
                    }
                    ByteBuffer buffer = ByteBuffer.allocate(1024);
                    channel.read(buffer);
                    String input = new String(buffer.array()).trim();
                    if (input.length() == 0) {
                        return;
                    }
                    System.out.println(Thread.currentThread().getName() + "收到客户端的消息 : " + input);
                    buffer.clear();
                    buffer.put((input + " can do!").getBytes());
                    buffer.flip();
                    channel.write(buffer);
                    cancelKey(selectionKey, channel);
                } catch (Exception ex) {
                    System.out.println("something wrong occurs when SocketHandler");
                    cancelKey(selectionKey, channel);
                }
            }
        }
    }
}

普通客户端socketClient代码

public class SocketClient {
    public static void main(String[] args) {
        runOneClient("HONOR");
        runOneClient("HUAWEI");
        runOneClient("MEIZU");
        runOneClient("XIAOMI");
    }

    private static void runOneClient(String content) {
        new Thread( () -> {
            try {
                Socket socket = new Socket();
                InetSocketAddress addr = new InetSocketAddress("127.0.0.1", 8088);
                socket.connect(addr);
                InputStream inputStream = socket.getInputStream();
                OutputStream outputStream = socket.getOutputStream();
                outputStream.write(content.getBytes());
                byte[] bytes = new byte[1024];
                inputStream.read(bytes);
                System.out.println("收到服务端消息 :" + new String(bytes));
            } catch (IOException ex) {
                System.out.println("something wrong occurs");
            }
        }).start();
    }
}
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,193评论 6 498
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,306评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 162,130评论 0 353
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,110评论 1 292
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,118评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,085评论 1 295
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,007评论 3 417
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,844评论 0 273
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,283评论 1 310
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,508评论 2 332
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,667评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,395评论 5 343
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,985评论 3 325
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,630评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,797评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,653评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,553评论 2 352

推荐阅读更多精彩内容