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
- Поддержка дескрипторов открытых сокетов (FD) ограничена только максимальным количеством файловых дескрипторов операционной системы, в то время как select поддерживает максимум 1024.
- select каждый раз сканирует все сокеты, а epoll сканирует только активные сокеты.
- Используйте 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();
}
}
Справочная статья:
- [Мифы о сплетнях и высоком параллелизме, посмотрите, как архитекторы Jingdong сняли это с алтаря
](Tickets.WeChat.QQ.com/Yes/LA В прошлом году 8cf SR…) 2. Анализ Java NIO 3. Подробное объяснение использования общих функций Berkeley API в протоколе TCP.4. «Руководство по программированию NIO и технологий сокетов»