предисловие
В этой статье в основном объясняется механизм ввода-вывода в Java и NIO, который обеспечивает высокий уровень параллелизма при сетевой связи.
Делится на две части:
Первый блок объясняет механизм ввода-вывода при многопоточности.
Второй блок объясняет, как оптимизировать растрату ресурсов ЦП в рамках механизма ввода-вывода (Новый ввод-вывод).
Эхо-сервер
Мне не нужно представлять механизм сокетов под один поток, если вы не понимаете, вы можете проверить информацию.
Что делать, если используются сокеты при таком количестве потоков?
Мы используем простейший эхо-сервер, чтобы помочь вам понять
Во-первых, давайте посмотрим на рабочий процесс сервера и клиента при многопоточности:
Как видите, несколько клиентов отправляют запросы на сервер одновременно.
Мера, принимаемая сервером, заключается в открытии нескольких потоков для соответствия соответствующему клиенту.
и каждый поток выполняет свой клиентский запрос самостоятельно
После того, как принцип закончен, давайте посмотрим, как он реализован
Вот я написал простой сервер
Для создания потоков используется технология пула потоков (я добавил комментарий для конкретной функции кода):
public class MyServer {
private static ExecutorService executorService = Executors.newCachedThreadPool(); //创建一个线程池
private static class HandleMsg implements Runnable{ //一旦有新的客户端请求,创建这个线程进行处理
Socket client; //创建一个客户端
public HandleMsg(Socket client){ //构造传参绑定
this.client = client;
}
@Override
public void run() {
BufferedReader bufferedReader = null; //创建字符缓存输入流
PrintWriter printWriter = null; //创建字符写入流
try {
bufferedReader = new BufferedReader(new InputStreamReader(client.getInputStream())); //获取客户端的输入流
printWriter = new PrintWriter(client.getOutputStream(),true); //获取客户端的输出流,true是随时刷新
String inputLine = null;
long a = System.currentTimeMillis();
while ((inputLine = bufferedReader.readLine())!=null){
printWriter.println(inputLine);
}
long b = System.currentTimeMillis();
System.out.println("此线程花费了:"+(b-a)+"秒!");
} catch (IOException e) {
e.printStackTrace();
}finally {
try {
bufferedReader.close();
printWriter.close();
client.close();
} catch (IOException e) {
e.printStackTrace();
}
}
}
}
public static void main(String[] args) throws IOException { //服务端的主线程是用来循环监听客户端请求
ServerSocket server = new ServerSocket(8686); //创建一个服务端且端口为8686
Socket client = null;
while (true){ //循环监听
client = server.accept(); //服务端监听到一个客户端请求
System.out.println(client.getRemoteSocketAddress()+"地址的客户端连接成功!");
executorService.submit(new HandleMsg(client)); //将该客户端请求通过线程池放入HandlMsg线程中进行处理
}
}
}
В приведенном выше коде мы используем класс для написания простого эхо-сервера и используем бесконечный цикл для открытия порта, прослушивающего основной поток.
простой клиент
С сервером мы можем получить к нему доступ и отправить некоторые строковые данные.Функция сервера состоит в том, чтобы вернуть эти строки и распечатать время, затраченное потоком.
Давайте напишем простой клиент для ответа серверу:
public class MyClient {
public static void main(String[] args) throws IOException {
Socket client = null;
PrintWriter printWriter = null;
BufferedReader bufferedReader = null;
try {
client = new Socket();
client.connect(new InetSocketAddress("localhost",8686));
printWriter = new PrintWriter(client.getOutputStream(),true);
printWriter.println("hello");
printWriter.flush();
bufferedReader = new BufferedReader(new InputStreamReader(client.getInputStream())); //读取服务器返回的信息并进行输出
System.out.println("来自服务器的信息是:"+bufferedReader.readLine());
} catch (IOException e) {
e.printStackTrace();
}finally {
printWriter.close();
bufferedReader.close();
client.close();
}
}
}
В коде мы используем поток символов для отправки строки приветствия в прошлом.Если код в порядке, сервер вернет данные приветствия и распечатает информацию журнала, которую мы установили.
отображение результатов эхо-сервера
Запускаем: 1. Открываем сервер и открываем цикл прослушивания:
2. Откройте клиент:
Вы можете видеть, что клиент распечатывает возвращенный результат
3. Просмотрите журнал сервера:
Очень хорошо, реализовано простое программирование многопоточного сокета
Но только представьте:
Если запрос клиента, добавьте Sleep в процессе записи IO на сервер,
Заставьте каждый запрос занимать 10 секунд в потоке сервера
Затем есть тонны клиентских запросов, каждый из которых занимает столько времени
Тогда параллельные возможности сервера будут значительно снижены.
Это не из-за того, сколько тяжелых задач у сервера, а просто потому, что служебный поток ожидает ввода-вывода (потому что прием, чтение, запись блокируются)
Очень неэкономично позволять высокоскоростному процессору ждать и его неэффективному сетевому вводу-выводу.
Что нам делать в это время?
NIO
Новый IO успешно решает вышеуказанные проблемы, как он их решает?
Наименьшая единица обработки клиентских запросов ввода-вывода — это поток.
И NIO использует блок, который на один уровень меньше, чем поток: канал (Channel)
Можно сказать, что только один поток в NIO может выполнять все операции приема, чтения, записи и другие операции.
Чтобы изучить NIO, вы должны сначала понять его три ядра.
Селектор, селектор
Буфер, буфер
Канал, канал
Блогер не талантлив, поэтому я нарисовала некрасивую картинку, чтобы всех впечатлить ^ ^
Дайте еще одну блок-схему работы NIO по TCP (так сложно рисовать линии...)
Всем достаточно понять, давайте пошагово
Buffer
Прежде всего, вам нужно знать, что такое Buffer
Взаимодействие данных в NIO больше не использует потоки, такие как механизмы ввода-вывода.
Вместо этого используйте буфер
Блогер считает, что проще всего понять картинку, поэтому...
Вы можете видеть положение буфера во всем рабочем процессе.
На практике конкретный код на приведенном выше рисунке выглядит следующим образом:
1. Сначала выделите место для буфера в байтах.
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
Создайте объект ByteBuffer и укажите размер памяти
2. Записать данные в буфер:
1).数据从Channel到Buffer:channel.read(byteBuffer);
2).数据从Client到Buffer:byteBuffer.put(...);
3. Чтение данных из буфера:
1).数据从Buffer到Channel:channel.write(byteBuffer);
2).数据从Buffer到Server:byteBuffer.get(...);
Selector
Селектор — это ядро NIO, это менеджер канала
Контролируйте, готов ли канал, выполнив метод блокировки select().
Как только данные доступны для чтения, возвращаемое значение этого метода представляет собой количество SelectionKeys.
Таким образом, сервер обычно выполняет метод select() в бесконечном цикле, пока канал не будет готов, а затем начинает работать.
Каждый канал привязывает событие к селектору, а затем генерирует объект SelectionKey.
нужно знать, это:
Когда канал привязан к селектору, канал должен быть в неблокирующем режиме.
И FileChannel не может переключиться в неблокирующий режим, потому что это не канал сокета, поэтому FileChannel не может привязывать события к Selector
В NIO есть четыре типа событий:
1.SelectionKey.OP_CONNECT: событие подключения
2.SelectionKey.OP_ACCEPT: получать события
3.SelectionKey.OP_READ: событие чтения
4.SelectionKey.OP_WRITE: событие записи
Channel
Есть четыре канала:
FileChannel: действует на файловый поток ввода-вывода.
DatagramChannel: действует по протоколу UDP.
SocketChannel: действует по протоколу TCP.
ServerSocketChannel: действует по протоколу TCP.
В этой статье объясняется NIO с помощью широко используемого протокола TCP.
Возьмем в качестве примера ServerSocketChannel:
Откройте канал ServerSocketChannel
ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
Закройте канал ServerSocketChannel:
serverSocketChannel.close();
Прослушивание цикла SocketChannel:
while(true){
SocketChannel socketChannel = serverSocketChannel.accept();
clientChannel.configureBlocking(false);
}
clientChannel.configureBlocking(false);Оператор должен установить этот канал как неблокирующий, то есть асинхронный.
Свободное управление блокировкой или неблокировкой — одна из характеристик NIO.
SelectionKey
SelectionKey — это основной компонент взаимодействия канала и селектора.Например, привяжите Selector к SocketChannel и зарегистрируйте его как событие подключения:
SocketChannel clientChannel = SocketChannel.open();
clientChannel.configureBlocking(false);
clientChannel.connect(new InetSocketAddress(port));
clientChannel.register(selector, SelectionKey.OP_CONNECT);
Ядро находится в методе register(), который возвращает объект SelectionKey для обнаружения события канала, которое может использовать следующие методы:
selectionKey.isAcceptable();
selectionKey.isConnectable();
selectionKey.isReadable();
selectionKey.isWritable();
Сервер использует эти методы для выполнения соответствующих операций при опросе.
Конечно, ключи, привязанные к Селектору через Канал, тоже можно получить по очереди.
Channel channel = selectionKey.channel();
Selector selector = selectionKey.selector();
Кстати, при регистрации событий на канале мы также можем привязать буфер:
clientChannel.register(key.selector(), SelectionKey.OP_READ,ByteBuffer.allocateDirect(1024));
Или привязать объект:
selectionKey.attach(Object);
Object anthorObj = selectionKey.attachment();
TCP-сервер NIO
Сказав так много, это все теория.Давайте взглянем на самый простой и самый основной код (не элегантно добавлять столько комментариев, но это легко понять каждому):
package cn.blog.test.NioTest;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.*;
import java.nio.charset.Charset;
import java.util.Iterator;
import java.util.Set;
public class MyNioServer {
private Selector selector; //创建一个选择器
private final static int port = 8686;
private final static int BUF_SIZE = 10240;
private void initServer() throws IOException {
//创建通道管理器对象selector
this.selector=Selector.open();
//创建一个通道对象channel
ServerSocketChannel channel = ServerSocketChannel.open();
channel.configureBlocking(false); //将通道设置为非阻塞
channel.socket().bind(new InetSocketAddress(port)); //将通道绑定在8686端口
//将上述的通道管理器和通道绑定,并为该通道注册OP_ACCEPT事件
//注册事件后,当该事件到达时,selector.select()会返回(一个key),如果该事件没到达selector.select()会一直阻塞
SelectionKey selectionKey = channel.register(selector,SelectionKey.OP_ACCEPT);
while (true){ //轮询
selector.select(); //这是一个阻塞方法,一直等待直到有数据可读,返回值是key的数量(可以有多个)
Set keys = selector.selectedKeys(); //如果channel有数据了,将生成的key访入keys集合中
Iterator iterator = keys.iterator(); //得到这个keys集合的迭代器
while (iterator.hasNext()){ //使用迭代器遍历集合
SelectionKey key = (SelectionKey) iterator.next(); //得到集合中的一个key实例
iterator.remove(); //拿到当前key实例之后记得在迭代器中将这个元素删除,非常重要,否则会出错
if (key.isAcceptable()){ //判断当前key所代表的channel是否在Acceptable状态,如果是就进行接收
doAccept(key);
}else if (key.isReadable()){
doRead(key);
}else if (key.isWritable() && key.isValid()){
doWrite(key);
}else if (key.isConnectable()){
System.out.println("连接成功!");
}
}
}
}
public void doAccept(SelectionKey key) throws IOException {
ServerSocketChannel serverChannel = (ServerSocketChannel) key.channel();
System.out.println("ServerSocketChannel正在循环监听");
SocketChannel clientChannel = serverChannel.accept();
clientChannel.configureBlocking(false);
clientChannel.register(key.selector(),SelectionKey.OP_READ);
}
public void doRead(SelectionKey key) throws IOException {
SocketChannel clientChannel = (SocketChannel) key.channel();
ByteBuffer byteBuffer = ByteBuffer.allocate(BUF_SIZE);
long bytesRead = clientChannel.read(byteBuffer);
while (bytesRead>0){
byteBuffer.flip();
byte[] data = byteBuffer.array();
String info = new String(data).trim();
System.out.println("从客户端发送过来的消息是:"+info);
byteBuffer.clear();
bytesRead = clientChannel.read(byteBuffer);
}
if (bytesRead==-1){
clientChannel.close();
}
}
public void doWrite(SelectionKey key) throws IOException {
ByteBuffer byteBuffer = ByteBuffer.allocate(BUF_SIZE);
byteBuffer.flip();
SocketChannel clientChannel = (SocketChannel) key.channel();
while (byteBuffer.hasRemaining()){
clientChannel.write(byteBuffer);
}
byteBuffer.compact();
}
public static void main(String[] args) throws IOException {
MyNioServer myNioServer = new MyNioServer();
myNioServer.initServer();
}
}
Я напечатал канал прослушивания, чтобы сообщить вам, когда запустился ServerSocketChannel.
Если вы сотрудничаете с отладкой клиента NIO, вы можете четко обнаружить, что перед входом в опрос select()
Хотя уже есть KEY события ACCEPT, select() не будет вызываться по умолчанию.
Вместо этого подождите, пока другие интересные события будут захвачены select(), прежде чем вызывать SelectionKey ACCEPT.
В это время ServerSocketChannel начинает выполнять мониторинг циклов.
То есть в селекторе всегда работает ServerSocketChannel.
иserverChannel.accept();Действительно асинхронный (channel.configureBlocking(false); в методе initServer)
Если соединение не принято, он вернет ноль
Если SocketChannel успешно подключен, этот SocketChannel регистрируется для события записи (READ).
и установить асинхронный
TCP-клиент для NIO
Если есть сервер, должен быть и клиент
На самом деле, если вы можете полностью понять сервер
Клиентский код аналогичен
package cn.blog.test.NioTest;
import java.io.IOException;
import java.net.InetSocketAddress;
import java.nio.ByteBuffer;
import java.nio.channels.SelectionKey;
import java.nio.channels.Selector;
import java.nio.channels.SocketChannel;
import java.util.Iterator;
public class MyNioClient {
private Selector selector; //创建一个选择器
private final static int port = 8686;
private final static int BUF_SIZE = 10240;
private static ByteBuffer byteBuffer = ByteBuffer.allocate(BUF_SIZE);
private void initClient() throws IOException {
this.selector = Selector.open();
SocketChannel clientChannel = SocketChannel.open();
clientChannel.configureBlocking(false);
clientChannel.connect(new InetSocketAddress(port));
clientChannel.register(selector, SelectionKey.OP_CONNECT);
while (true){
selector.select();
Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
while (iterator.hasNext()){
SelectionKey key = iterator.next();
iterator.remove();
if (key.isConnectable()){
doConnect(key);
}else if (key.isReadable()){
doRead(key);
}
}
}
}
public void doConnect(SelectionKey key) throws IOException {
SocketChannel clientChannel = (SocketChannel) key.channel();
if (clientChannel.isConnectionPending()){
clientChannel.finishConnect();
}
clientChannel.configureBlocking(false);
String info = "服务端你好!!";
byteBuffer.clear();
byteBuffer.put(info.getBytes("UTF-8"));
byteBuffer.flip();
clientChannel.write(byteBuffer);
//clientChannel.register(key.selector(),SelectionKey.OP_READ);
clientChannel.close();
}
public void doRead(SelectionKey key) throws IOException {
SocketChannel clientChannel = (SocketChannel) key.channel();
clientChannel.read(byteBuffer);
byte[] data = byteBuffer.array();
String msg = new String(data).trim();
System.out.println("服务端发送消息:"+msg);
clientChannel.close();
key.selector().close();
}
public static void main(String[] args) throws IOException {
MyNioClient myNioClient = new MyNioClient();
myNioClient.initClient();
}
}
выходной результат
Здесь я открываю сервер и два клиента:
Далее можно попробовать открыть тысячу клиентов одновременно, пока ваш ЦП достаточно мощный, сервер не сможет снизить производительность из-за блокировки
Выше приведено подробное объяснение основ Java NIO.
Спасибо за прочтение и подписку~