Посмотрите, как я снял NIO с алтаря

Java

1. Традиционный блокирующий ввод-вывод

Блокировка блокирующего ввода-вывода означает, что функции чтения и записи сокета заблокированы.

1.2 Модель программирования блокирующего ввода/вывода

public static void main(String[] args) throws IOException {
        ServerSocket serverSocket = new ServerSocket(8091);
        System.out.println("step1: bind 8091");
        while (true) {
            // 阻塞
            Socket socket = serverSocket.accept();
            System.out.println("step2: accept " + socket.getPort());
            new Thread(() -> {
                try (BufferedReader reader = new BufferedReader(new InputStreamReader(socket.getInputStream()));
                     PrintWriter out = new PrintWriter(socket.getOutputStream(), true)) {
                    String line;
                    // 阻塞
                    while ((line = reader.readLine()) != null) {
                        System.out.println(line);
                        out.println("Server recv:" + line);
                    }
                } catch (Exception e) {

                }
            }).start();
        }

    }

Поскольку функция принятия, функция чтения и функция записи сокета заблокированы синхронно, основной поток постоянно вызывает функцию принятия сокета, а состоянием опроса является установленное TCP-соединение.

Функция чтения прочитает подготовленные данные из буфера ядра и скопирует их в пользовательский процесс.Если данных в буфере ядра нет, поток будет приостановлен и соответствующие права использования процессора будут освобождены. Когда данные будут готовы в буфере ядра, процессор ответит на сигнал прерывания ввода-вывода и разбудит заблокированный поток для обработки данных.

Когда соединение обрабатывает ввод-вывод, система блокируется. Если это один поток, он должен зависнуть там, но ЦП освобождается, и включается многопоточность, чтобы ЦП мог обрабатывать больше вещей.

Блокирующая модель ввода/вывода

Недостатки блокировки ввода-вывода

Отсутствие масштабируемости и сильная зависимость от потоков. Потоки Java занимают от 512 000 до 1 М памяти. Слишком большое количество потоков приведет к переполнению памяти JVM. Большое количество переключений контекста потока серьезно снижает производительность ЦП. Активация большого количества потоков ввода-вывода может вызвать пилообразную нагрузку на систему.

2. Программирование НИО

Синхронная неблокирующая модель ввода/вывода

Для NIO при отсутствии данных в буфере ядра возвращается ошибка EWOULDBLOCK напрямую, вообще говоря, процесс может вызывать функцию чтения путем опроса, а при наличии данных в буфере данные копируются в пространство пользователя не прерывая нить..

такНеблокировка в синхронной неблокировке означает, что функции чтения и записи сокета не заблокированы, но пользовательскому процессу все равно необходимо опрашивать функции чтения и записи, поэтому он является синхронным. Но NIO дает нам возможность использовать ЦП, не запуская новый поток, то есть технологию мультиплексирования ввода-вывода.

2.1 Технология мультиплексирования ввода/вывода

В системе Linux вы можете использовать select/poll/epoll для мониторинга нескольких сокетов с помощью одного потока.Пока есть данные в кеше чтения одного сокета, метод немедленно возвращается, и тогда вы можете прочитать доступный для чтения сокет.Если Все буферы чтения сокета пусты, он заблокируется, то есть поток будет приостановлен.

Сначала в Linux использовался select, но select был медленным, и, наконец, был использован epoll.

2.1.1 Преимущества epoll

  1. Поддержка дескрипторов открытых сокетов (FD) ограничена только максимальным количеством файловых дескрипторов операционной системы, в то время как select поддерживает максимум 1024.
  2. select каждый раз сканирует все сокеты, а epoll сканирует только активные сокеты.
  3. Используйте mmap для ускорения копирования данных из пространства ядра в пространство пользователя.

2.2 Рабочий механизм НИО

НИО на самом делемодель, управляемая событиями, самое главное в NIO этоМультиплексор (селектор). В NIO предусмотрена возможность выбора готовых событий, нам достаточно поставитьКаналЗарегистрированный на Селекторе, Селектор будет непрерывно опрашивать зарегистрированный на нем Канал через метод select (собственно, операционная система через epoll).SelectionKey(Когда Channel зарегистрирован в Selector, он вернет связанный с ним SelectionKey) Вы можете получить готовую коллекцию Channel, иначе Selector заблокирует метод select.

Selector вызывает метод select, не поток выбирает готовый Channel через цикл for, а операционная система уведомляет поток JVM в виде события через epoll, какой канал имеет событие read-ready или write-ready. Таким образом, метод select больше похож на прослушиватель.

Основная цель мультиплексирования — использовать наименьшее количество потоков для управления большим количеством каналов, а внутри находится не только один поток. Количество создаваемых потоков определяется в соответствии с количеством каналов, и каждый раз, когда регистрируется 1023 канала, создается новый поток.

Ядром NIO являетсяМультиплексор и модель событий, Разобравшись с этими двумя моментами, можно собственно разобраться в основном принципе работы NIO. Оказалось, что выучить NIO очень сложно, при глубоком понимании TCP найти NIO несложно. При использовании NIO,Основная генерация заключается в регистрации канала и событий, которые необходимо отслеживать, в селекторе.

События, поддерживаемые различными типами каналов

Схематическая диаграмма модели событий NIO:

2.2.1 Пример кода

ServerReactor

@Slf4j
public class ServerReactor implements Runnable {
    private final Selector selector;
    private final ServerSocketChannel serverSocketChannel;
    private volatile boolean stop = false;

    public ServerReactor(int port, int backlog) throws IOException {
        selector = Selector.open();
        serverSocketChannel = ServerSocketChannel.open();
        ServerSocket serverSocket = serverSocketChannel.socket();
        serverSocket.bind(new InetSocketAddress(port), backlog);
        serverSocket.setReuseAddress(true);
        serverSocketChannel.configureBlocking(false);
        // 将channel注册到多路复用器上,并监听ACCEPT事件
        serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
    }

    public void setStop(boolean stop) {
        this.stop = stop;
    }

    @Override
    public void run() {
        try {
            // 无限的接收客户端连接
            while (!stop && !Thread.interrupted()) {
                int num = selector.select();
                Set<SelectionKey> selectionKeys = selector.selectedKeys();
                Iterator<SelectionKey> it = selectionKeys.iterator();
                while (it.hasNext()) {
                    SelectionKey key = it.next();
                    // 移除key,否则会导致事件重复消费
                    it.remove();
                    try {
                        handle(key);
                    } catch (Exception e) {
                        if (key != null) {
                            key.cancel();
                            if (key.channel() != null) {
                                key.channel().close();
                            }
                        }
                    }
                }
            }
        } catch (IOException e) {
            e.printStackTrace();
        }
        if (selector != null) {
            try {
                selector.close();
            } catch (IOException e) {
                e.printStackTrace();
            }

        }
    }

    private void handle(SelectionKey key) throws Exception {
        if (key.isValid()) {
            // 如果是ACCEPT事件,代表是一个新的连接请求
            if (key.isAcceptable()) {
                ServerSocketChannel serverSocketChannel = (ServerSocketChannel) key.channel();
                // 相当于三次握手后,从全连接队列中获取可用的连接
                // 必须使用accept方法消费ACCEPT事件,否则将导致多路复用器死循环
                SocketChannel socketChannel = serverSocketChannel.accept();
                // 设置为非阻塞模式,当没有可用的连接时直接返回null,而不是阻塞。
                socketChannel.configureBlocking(false);
                socketChannel.register(selector, SelectionKey.OP_READ);
            }

            if (key.isReadable()) {
                SocketChannel socketChannel = (SocketChannel) key.channel();
                ByteBuffer readBuffer = ByteBuffer.allocate(1024);
                int readBytes = socketChannel.read(readBuffer);
                if (readBytes > 0) {
                    readBuffer.flip();
                    byte[] bytes = new byte[readBuffer.remaining()];
                    readBuffer.get(bytes);
                    String content = new String(bytes);
                    System.out.println("recv client content: " + content);
                    ByteBuffer writeBuffer = ByteBuffer.allocate(1024);
                    writeBuffer.put(("服务端已收到: " + content).getBytes());
                    writeBuffer.flip();
                    socketChannel.write(writeBuffer);

                } else if (readBytes < 0) {
                    key.cancel();
                    socketChannel.close();
                }
            }

        }
    }
}

ClientReactor

public class ClientReactor implements Runnable {
    final String host;
    final int port;
    final SocketChannel socketChannel;
    final Selector selector;
    private volatile boolean stop = false;

    public ClientReactor(String host, int port) throws IOException {
        this.socketChannel = SocketChannel.open();
        this.socketChannel.configureBlocking(false);
        Socket socket = this.socketChannel.socket();
        socket.setTcpNoDelay(true);
        this.selector = Selector.open();
        this.host = host;
        this.port = port;

    }

    @Override
    public void run() {

        try {
            // 如果通道呈阻塞模式,则立即发起连接;
            // 如果呈非阻塞模式,则不是立即发起连接,而是在随后的某个时间才发起连接。

            // 如果连接是立即建立的,说明通道是阻塞模式,当连接成功时,则此方法返回true,连接失败出现异常。
            // 如果此通道处于阻塞模式,则此方法的调用将会阻塞,直到建立连接或发生I/O错误。

            // 如果连接不是立即建立的,说明通道是非阻塞模式,则此方法返回false,
            // 并且以后必须通过调用finishConnect()方法来验证连接是否完成
            // socketChannel.isConnectionPending()判断此通道是否正在进行连接
            if (socketChannel.connect(new InetSocketAddress(host, port))) {
                socketChannel.register(selector, SelectionKey.OP_READ);
                doWrite(socketChannel);
            } else {
                socketChannel.register(selector, SelectionKey.OP_CONNECT);

            }
            while (!stop && !Thread.interrupted()) {
                int num = selector.select();
                Set<SelectionKey> selectionKeys = selector.selectedKeys();
                Iterator<SelectionKey> it = selectionKeys.iterator();
                while (it.hasNext()) {
                    SelectionKey key = it.next();
                    // 移除key,否则会导致事件重复消费
                    it.remove();
                    try {
                        handle(key);
                    } catch (Exception e) {
                        if (key != null) {
                            key.cancel();
                            if (key.channel() != null) {
                                key.channel().close();
                            }
                        }
                    }
                }
            }
        } catch (IOException e) {
            e.printStackTrace();
        }

        if (selector != null) {
            try {
                selector.close();
            } catch (IOException e) {
                e.printStackTrace();
            }

        }


    }

    private void handle(SelectionKey key) throws IOException {

        if (key.isValid()) {

            SocketChannel socketChannel = (SocketChannel) key.channel();

            if (key.isConnectable()) {
                if (socketChannel.finishConnect()) {
                    socketChannel.register(selector, SelectionKey.OP_READ);
                    doWrite(socketChannel);
                }
            }

            if (key.isReadable()) {
                ByteBuffer readBuffer = ByteBuffer.allocate(1024);
                int readBytes = socketChannel.read(readBuffer);
                if (readBytes > 0) {
                    readBuffer.flip();
                    byte[] bytes = new byte[readBuffer.remaining()];
                    readBuffer.get(bytes);
                    System.out.println("recv server content: " + new String(bytes));
                } else if (readBytes < 0) {
                    key.cancel();
                    socketChannel.close();
                }
            }

        }
    }

    private void doWrite(SocketChannel socketChannel) {
        Scanner scanner = new Scanner(System.in);
        new Thread(() -> {
            while (scanner.hasNext()) {
                try {

                    ByteBuffer writeBuffer = ByteBuffer.allocate(1024);
                    writeBuffer.put(scanner.nextLine().getBytes());
                    writeBuffer.flip();
                    socketChannel.write(writeBuffer);
                } catch (Exception e) {

                }
            }
        }).start();
    }
}

Справочная статья:

  1. [Мифы о сплетнях и высоком параллелизме, посмотрите, как архитекторы Jingdong сняли это с алтаря

](Tickets.WeChat.QQ.com/Yes/LA В прошлом году 8cf SR…) 2. Анализ Java NIO 3. Подробное объяснение использования общих функций Berkeley API в протоколе TCP.4. «Руководство по программированию NIO и технологий сокетов»