карта разума
Обучение похоже на греблю вверх по течению
1 обзор НИО
1.1 Определения
java.nioполное имяjava non-blocking IO,Относится кJDK1.4 и вышеНовый API (Новый IO), представленный в версии, предоставляется для всех примитивных типов (кроме логических типов).Кэшированные контейнеры данных, использование которого обеспечиваетнеблокирующийВысококалетабельная сеть (от энциклопедии Baidu).
1.2 Зачем использовать NIO
Как упоминалось в описании выше, NIO предоставляется только в версиях выше JDK1.4, так что же он использовался раньше? Ответ прост, этоBIO(блокирующий ввод-вывод), который является нашим обычно используемым потоком ввода-вывода.
Излишне говорить о проблеме BIO, потому что при использовании BIO основной поток будет переходить в состояние блокировки, что сильно влияет на производительность программы.Неспособность полностью использовать ресурсы машины. Но тогда некоторые люди будут задавать вопросы, тогда яИспользовать многопоточностьВы не можете сделать это?
Однако в случае высокого параллелизма будет создано много потоков, потоки будут занимать память, а переключение между потоками также приведет к потере ресурсов.
И НИОТолько если соединение/канал действительно имеет события чтения и записикогда (управляемый событиями),прочти и напиши, что значительно снижает нагрузку на систему. Вам не нужно создавать поток для каждого соединения, и вам не нужно поддерживать несколько потоков.
Избегает переключения контекста между несколькими потоками, что приводит к пустой трате ресурсов.
2 Три ядра NIO
| Ядро НИО | Соответствующий класс или интерфейс | заявление | эффект |
|---|---|---|---|
| буфер | java.nio.Buffer | Файловый ввод-вывод/сетевой ввод-вывод | Хранение данных |
| ряд | java.nio.channels.Channel | Файловый ввод-вывод/сетевой ввод-вывод | транспорт |
| Селектор | java.nio.channels.Selector | сетевой ввод-вывод | контроллер |
2.1 Буфер (буфер)
2.1.1 Что такое буфер
Давайте сначала посмотрим на следующую диаграмму классов, мы можем видетьBufferЕсть семь типов.
Bufferэто блок памяти. существуетNIO, все данные используютсяBufferОбработка, есть два режима чтения и записи. Таким образом, здесь отражена разница между NIO и традиционным IO. Традиционный ввод-вывод дляStreamпоток,NIOвместо этого ориентированный на буфер (Buffer).
2.1.2 Обычно используемые типы ByteBuffer
Как правило, типы, которые мы обычно используем,ByteBuffer, преобразовать данные в байты для обработки. по существуbyte[]множество.
public abstract class ByteBuffer extends Buffer implements Comparable<ByteBuffer>{
//存储数据的数组
final byte[] hb;
//构造器方法
ByteBuffer(int mark, int pos, int lim, int cap, byte[] hb, int offset) {
super(mark, pos, lim, cap);
//初始化数组
this.hb = hb;
this.offset = offset;
}
}
2.1.3 Как создать буфер
В основном он делится на два типа: буфер блока памяти JVM в куче и буфер блока памяти вне кучи.
Способ создания блока памяти в куче (непрямой буфер):
//创建堆内内存块HeapByteBuffer
ByteBuffer byteBuffer1 = ByteBuffer.allocate(1024);
String msg = "java技术爱好者";
//包装一个byte[]数组获得一个Buffer,实际类型是HeapByteBuffer
ByteBuffer byteBuffer2 = ByteBuffer.wrap(msg.getBytes());
Методы создания блоков памяти вне кучи (прямые буферы):
//创建堆外内存块DirectByteBuffer
ByteBuffer byteBuffer3 = ByteBuffer.allocateDirect(1024);
2.1.3.1 Разница между HeapByteBuffer и DirectByteBuffer
Собственно, это видно из названия класса,HeapByteBufferСозданный буфер байтов находится в куче JVM, то есть в массиве байтов, поддерживаемом внутри JVM. а такжеDirectByteBufferдаПрямая манипуляция собственным кодом операционной системысозданныймассив буферов памяти.
DirectByteBufferсценарии использования:
-
java программа с локальным диском, передача данных через сокет
-
Можно использовать большие файловые объекты. Не ограничен размером динамической памяти.
-
Его не нужно создавать часто, жизненный цикл длинный, и его можно использовать повторно.
HeapByteBufferсценарии использования:
В дополнение к вышеперечисленным сценариям рекомендуется использовать и другие ситуации.HeapByteBuffer, не достигает определенного порядка, фактически используяDirectByteBufferне отражает преимущество.
2.1.3.2 Начальный опыт работы с Buffer
Далее используйтеByteBufferСделайте небольшой пример, чтобы ознакомиться с ним:
public static void main(String[] args) throws Exception {
String msg = "java技术爱好者,起飞!";
//创建一个固定大小的buffer(返回的是HeapByteBuffer)
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
byte[] bytes = msg.getBytes();
//写入数据到Buffer中
byteBuffer.put(bytes);
//切换成读模式,关键一步
byteBuffer.flip();
//创建一个临时数组,用于存储获取到的数据
byte[] tempByte = new byte[bytes.length];
int i = 0;
//如果还有数据,就循环。循环判断条件
while (byteBuffer.hasRemaining()) {
//获取byteBuffer中的数据
byte b = byteBuffer.get();
//放到临时数组中
tempByte[i] = b;
i++;
}
//打印结果
System.out.println(new String(tempByte));//java技术爱好者,起飞!
}
Этоflip()Метод важен. Значит переключиться в режим чтения. упомянутый вышеБуферы двунаправленные,Вы можете как записывать данные в буфер, так и читать данные из буфера.. Но не одновременно, нужно переключаться. Так в чем же суть этого режима переключения?
2.1.4 Три важных параметра
//位置,默认是从第一个开始
private int position = 0;
//限制,不能读取或者写入的位置索引
private int limit;
//容量,缓冲区所包含的元素的数量
private int capacity;
Затем мы используем приведенный выше пример для анализа кода предложение за предложением:
String msg = "java技术爱好者,起飞!";
//创建一个固定大小的buffer(返回的是HeapByteBuffer)
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
При создании буфера значение параметра выглядит так:
при выполнении кbyteBuffer.put(bytes),когдаput()Сколько вводится данных, насколько увеличится позиция, и изменятся параметры:
Следующий ключевой шагbyteBuffer.flip(), произойдут следующие изменения:
flip()Исходный код метода выглядит следующим образом:
public final Buffer flip() {
limit = position;
position = 0;
mark = -1;
return this;
}
Зачем так назначать? Потому что ниже есть условное суждение цикла:
byteBuffer.hasRemaining();
public final boolean hasRemaining() {
//判断position的索引是否小于limit。
//所以可以看出limit的作用就是记录写入数据的位置,那么当读取数据时,就知道读到哪个位置
return position < limit;
}
Далее вwhileв циклеget()Прочтите данные после прочтения.
наконец, когдаpositionравныйlimitКогда условие оценки цикла не установлено, он выходит из цикла, и чтение завершается.
Так что видно, что на самом делеcapacityВеличина емкости постоянна и фактически контролируетсяpositionа такжеlimitЗначение для управления чтением и записью данных.
2.2 Канал (Канал)
Во-первых, давайте посмотрим, какие подклассы есть у Channel:
Обычно используются четыре канала:
FileChannel, чтение и запись данных в файлы. SocketChannel, чтение и запись данных в сети через TCP. ServerSocketChannel, прослушивающий входящие соединения TCP, например веб-сервер. SocketChannel создается для каждого нового входящего соединения. DatagramChannel, чтение и запись данных в сети через UDP.
Сам канал не хранит данные, он отвечает только за транспортировку данных. должен иBufferиспользовать вместе.
2.2.1 Как получить канал
2.2.1.1 FileChannel
Способ получения FileChannel описан ниже на примере копирования файла:
Сначала подготовьте «1.txt» в корневом каталоге проекта, а затем напишите основной метод:
public static void main(String[] args) throws Exception {
//获取文件输入流
File file = new File("1.txt");
FileInputStream inputStream = new FileInputStream(file);
//从文件输入流获取通道
FileChannel inputStreamChannel = inputStream.getChannel();
//获取文件输出流
FileOutputStream outputStream = new FileOutputStream(new File("2.txt"));
//从文件输出流获取通道
FileChannel outputStreamChannel = outputStream.getChannel();
//创建一个byteBuffer,小文件所以就直接一次读取,不分多次循环了
ByteBuffer byteBuffer = ByteBuffer.allocate((int)file.length());
//把输入流通道的数据读取到缓冲区
inputStreamChannel.read(byteBuffer);
//切换成读模式
byteBuffer.flip();
//把数据从缓冲区写入到输出流通道
outputStreamChannel.write(byteBuffer);
//关闭通道
outputStream.close();
inputStream.close();
outputStreamChannel.close();
inputStreamChannel.close();
}
После выполнения получаем файл "2.txt". выполнение удалось.
Вышеприведенный пример, который можно представить схематически, выглядит так:
2.2.1.2 SocketChannel
Далее мы учимся получатьSocketChannelПуть.
Тем не менее, давайте быстро начнем с примера:
public static void main(String[] args) throws Exception {
//获取ServerSocketChannel
ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
InetSocketAddress address = new InetSocketAddress("127.0.0.1", 6666);
//绑定地址,端口号
serverSocketChannel.bind(address);
//创建一个缓冲区
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
while (true) {
//获取SocketChannel
SocketChannel socketChannel = serverSocketChannel.accept();
while (socketChannel.read(byteBuffer) != -1){
//打印结果
System.out.println(new String(byteBuffer.array()));
//清空缓冲区
byteBuffer.clear();
}
}
}
Затем запустите метод main(), мы можем передатьtelnetКоманда для проверки соединения:
Из приведенного выше примера мы можем узнать, что поServerSocketChannel.open()Метод может получить канал сервера, затем привязать номер порта адреса, а затемaccept()метод полученияSocketChannelКанал, то есть канал подключения клиента.
Наконец, с использованиемBufferВы можете читать и писать.
Это простой пример, на самом деле приведенный выше пример является блокирующим. Чтобы быть неблокирующим, вам также нужно использовать селекторыSelector.
2.3 Селектор
Selectorперевести наСелектор, некоторые также переводятся какмультиплексор, фактически говоря об одном и том же.
Селекторы используются только для сетевого ввода-вывода, а не для файлового ввода-вывода.
Можно сказать, что селектор является основным компонентом NIO, который может отслеживать состояние канала для достижения асинхронного неблокирующего ввода-вывода. Другими словами, он управляется событиями. достичь этогоУправление несколькими каналами с помощью одного потокацель.
2.3.1 Основной API
| Имя метода API | эффект |
|---|---|
| Selector.open() | Откройте селектор. |
| select() | Выбирает группу клавиш, соответствующий канал которой готов к операции ввода/вывода. |
| selectedKeys() | Возвращает выбранный набор ключей для этого селектора. |
Приведенный выше API будет использоваться в следующих примерах, поэтому сначала у меня сложилось впечатление.
3 Быстрый старт НИО
3.1 Файловый ввод-вывод
3.1.1 Передача данных между каналами
Здесь мы в основном вводим способ передачи данных между двумя каналами:
transferTo(): передача данных исходного канала в целевой канал.
public static void main(String[] args) throws Exception {
//获取文件输入流
File file = new File("1.txt");
FileInputStream inputStream = new FileInputStream(file);
//从文件输入流获取通道
FileChannel inputStreamChannel = inputStream.getChannel();
//获取文件输出流
FileOutputStream outputStream = new FileOutputStream(new File("2.txt"));
//从文件输出流获取通道
FileChannel outputStreamChannel = outputStream.getChannel();
//创建一个byteBuffer,小文件所以就直接一次读取,不分多次循环了
ByteBuffer byteBuffer = ByteBuffer.allocate((int) file.length());
//把输入流通道的数据读取到输出流的通道
inputStreamChannel.transferTo(0, byteBuffer.limit(), outputStreamChannel);
//关闭通道
outputStream.close();
inputStream.close();
outputStreamChannel.close();
inputStreamChannel.close();
}
transferFrom(): передача данных из исходного канала в целевой канал.
public static void main(String[] args) throws Exception {
//获取文件输入流
File file = new File("1.txt");
FileInputStream inputStream = new FileInputStream(file);
//从文件输入流获取通道
FileChannel inputStreamChannel = inputStream.getChannel();
//获取文件输出流
FileOutputStream outputStream = new FileOutputStream(new File("2.txt"));
//从文件输出流获取通道
FileChannel outputStreamChannel = outputStream.getChannel();
//创建一个byteBuffer,小文件所以就直接一次读取,不分多次循环了
ByteBuffer byteBuffer = ByteBuffer.allocate((int) file.length());
//把输入流通道的数据读取到输出流的通道
outputStreamChannel.transferFrom(inputStreamChannel,0,byteBuffer.limit());
//关闭通道
outputStream.close();
inputStream.close();
outputStreamChannel.close();
inputStreamChannel.close();
}
3.1.2 Рассеянное чтение и агрегированная запись
Давайте сначала посмотрим на исходный код FileChannel:
public abstract class FileChannel extends AbstractInterruptibleChannel
implements SeekableByteChannel, GatheringByteChannel, ScatteringByteChannel {
}
Из исходного кода видно, что реализованы интерфейсы GatheringByteChannel и ScatteringByteChannel. То есть операции, поддерживающие чтение с разбросом и агрегированную запись. Как его использовать, смотрите в следующем примере:
Мы пишем основной метод для копирования файла 1.txt.Содержимое файла:
abcdefghijklmnopqrstuvwxyz//26个字母
код показывает, как показано ниже:
public static void main(String[] args) throws Exception {
//获取文件输入流
File file = new File("1.txt");
FileInputStream inputStream = new FileInputStream(file);
//从文件输入流获取通道
FileChannel inputStreamChannel = inputStream.getChannel();
//获取文件输出流
FileOutputStream outputStream = new FileOutputStream(new File("2.txt"));
//从文件输出流获取通道
FileChannel outputStreamChannel = outputStream.getChannel();
//创建三个缓冲区,分别都是5
ByteBuffer byteBuffer1 = ByteBuffer.allocate(5);
ByteBuffer byteBuffer2 = ByteBuffer.allocate(5);
ByteBuffer byteBuffer3 = ByteBuffer.allocate(5);
//创建一个缓冲区数组
ByteBuffer[] buffers = new ByteBuffer[]{byteBuffer1, byteBuffer2, byteBuffer3};
//循环写入到buffers缓冲区数组中,分散读取
long read;
long sumLength = 0;
while ((read = inputStreamChannel.read(buffers)) != -1) {
sumLength += read;
Arrays.stream(buffers)
.map(buffer -> "posstion=" + buffer.position() + ",limit=" + buffer.limit())
.forEach(System.out::println);
//切换模式
Arrays.stream(buffers).forEach(Buffer::flip);
//聚合写入到文件输出通道
outputStreamChannel.write(buffers);
//清空缓冲区
Arrays.stream(buffers).forEach(Buffer::clear);
}
System.out.println("总长度:" + sumLength);
//关闭通道
outputStream.close();
inputStream.close();
outputStreamChannel.close();
inputStreamChannel.close();
}
распечатать результат:
posstion=5,limit=5
posstion=5,limit=5
posstion=5,limit=5
posstion=5,limit=5
posstion=5,limit=5
posstion=1,limit=5
总长度:26
Вы можете видеть, что он зацикливается дважды. В первом цикле все три буфера считывают 5 байтов, всего 15 байтов, то есть заполнены. Осталось 11 байт, поэтому во втором цикле первым двум буферам выделяется 5 байт, а последнему буферу выделяется 1 байт, только что закончил чтение. Всего 26 байт.
Это процесс разбросанного чтения и агрегированной записи.
Сценарий - это то, что вы можете использоватьИспользуйте массив буферов для автоматического выделения размера буфера по мере необходимости. Может уменьшить потребление памяти. Также можно использовать сетевой ввод-вывод, поэтому я не буду приводить здесь демонстрационный пример.
3.1.3 Косвенные/прямые буферы
Как создается косвенный буфер:
static ByteBuffer allocate(int capacity)
Как создается прямой буфер:
static ByteBuffer allocateDirect(int capacity)
Схематическая диаграмма разницы между непрямыми/прямыми буферами:
Из диаграммы видно, что самое большое отличие состоит в том, что прямым буферам не нужно файлировать, а затем копировать содержимое физической памяти. Это значительно повышает производительность. Фактически, во введении Buffer мы соприкоснулись с этой концепцией. Прямая буферная куча памяти снаружи, эффективность локального файлового ввода-вывода будет немного выше.
Далее сравним эффективность на примере видеофайла размером 136 МБ:
public static void main(String[] args) throws Exception {
long starTime = System.currentTimeMillis();
//获取文件输入流
File file = new File("D:\\小电影.mp4");//文件大小136 MB
FileInputStream inputStream = new FileInputStream(file);
//从文件输入流获取通道
FileChannel inputStreamChannel = inputStream.getChannel();
//获取文件输出流
FileOutputStream outputStream = new FileOutputStream(new File("D:\\test.mp4"));
//从文件输出流获取通道
FileChannel outputStreamChannel = outputStream.getChannel();
//创建一个直接缓冲区
ByteBuffer byteBuffer = ByteBuffer.allocateDirect(5 * 1024 * 1024);
//创建一个非直接缓冲区
//ByteBuffer byteBuffer = ByteBuffer.allocate(5 * 1024 * 1024);
//写入到缓冲区
while (inputStreamChannel.read(byteBuffer) != -1) {
//切换读模式
byteBuffer.flip();
outputStreamChannel.write(byteBuffer);
byteBuffer.clear();
}
//关闭通道
outputStream.close();
inputStream.close();
outputStreamChannel.close();
inputStreamChannel.close();
long endTime = System.currentTimeMillis();
System.out.println("消耗时间:" + (endTime - starTime) + "毫秒");
}
результат:
Прошедшее время для прямого буфера: 283 мс
Прошедшее время для непрямых буферов: 487 мс
3.2 Сетевой ввод-вывод
На самом деле, основным целям NIO является сеть IO. До Nio, это было полезно только для Java, чтобы использовать сетевое программирование.Socket. а такжеSocketЭто блокировка, очевидно, не подходящая для сценариев с высоким параллелизмом. Таким образом, появление NIO должно решить эту проблему.
Основная идея заключается в том, чтобы прописать канал Channel в Селекторе, и использовать Селектор для мониторинга состояния события в Канале, чтобы не было необходимости блокировать ожидание подключения клиента, от активного ожидания подключения клиент должен управляться событиями. Никаких событий для прослушивания нет, сервер может заниматься своими делами.
3.2.1 Небольшой пример с использованием Selector
Далее, куй железо, пока горячо, сделаем пример сервера, принимающего клиентские сообщения:
Код первого сервера:
public class NIOServer {
public static void main(String[] args) throws Exception {
//打开一个ServerSocketChannel
ServerSocketChannel serverSocketChannel = ServerSocketChannel.open();
InetSocketAddress address = new InetSocketAddress("127.0.0.1", 6666);
//绑定地址
serverSocketChannel.bind(address);
//设置为非阻塞
serverSocketChannel.configureBlocking(false);
//打开一个选择器
Selector selector = Selector.open();
//serverSocketChannel注册到选择器中,监听连接事件
serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
//循环等待客户端的连接
while (true) {
//等待3秒,(返回0相当于没有事件)如果没有事件,则跳过
if (selector.select(3000) == 0) {
System.out.println("服务器等待3秒,没有连接");
continue;
}
//如果有事件selector.select(3000)>0的情况,获取事件
Set<SelectionKey> selectionKeys = selector.selectedKeys();
//获取迭代器遍历
Iterator<SelectionKey> it = selectionKeys.iterator();
while (it.hasNext()) {
//获取到事件
SelectionKey selectionKey = it.next();
//判断如果是连接事件
if (selectionKey.isAcceptable()) {
//服务器与客户端建立连接,获取socketChannel
SocketChannel socketChannel = serverSocketChannel.accept();
//设置成非阻塞
socketChannel.configureBlocking(false);
//把socketChannel注册到selector中,监听读事件,并绑定一个缓冲区
socketChannel.register(selector, SelectionKey.OP_READ, ByteBuffer.allocate(1024));
}
//如果是读事件
if (selectionKey.isReadable()) {
//获取通道
SocketChannel socketChannel = (SocketChannel) selectionKey.channel();
//获取关联的ByteBuffer
ByteBuffer buffer = (ByteBuffer) selectionKey.attachment();
//打印从客户端获取到的数据
socketChannel.read(buffer);
System.out.println("from 客户端:" + new String(buffer.array()));
}
//从事件集合中删除已处理的事件,防止重复处理
it.remove();
}
}
}
}
Код клиента:
public class NIOClient {
public static void main(String[] args) throws Exception {
SocketChannel socketChannel = SocketChannel.open();
InetSocketAddress address = new InetSocketAddress("127.0.0.1", 6666);
socketChannel.configureBlocking(false);
//连接服务器
boolean connect = socketChannel.connect(address);
//判断是否连接成功
if(!connect){
//等待连接的过程中
while (!socketChannel.finishConnect()){
System.out.println("连接服务器需要时间,期间可以做其他事情...");
}
}
String msg = "hello java技术爱好者!";
ByteBuffer byteBuffer = ByteBuffer.wrap(msg.getBytes());
//把byteBuffer数据写入到通道中
socketChannel.write(byteBuffer);
//让程序卡在这个位置,不关闭连接
System.in.read();
}
}
Затем запустите сервер, а затем запустите клиент, мы можем увидеть, как консоль выводит следующую информацию:
服务器等待3秒,没有连接
服务器等待3秒,没有连接
from 客户端:hello java技术爱好者!
服务器等待3秒,没有连接
服务器等待3秒,没有连接
На этом примере мы рисуем следующие точки знаний.
3.2.2 SelectionKey
существуетSelectionKeyВ классе есть четыре константы для представления четырех типов событий, посмотрите исходный код:
public abstract class SelectionKey {
//读事件
public static final int OP_READ = 1 << 0; //2^0=1
//写事件
public static final int OP_WRITE = 1 << 2; // 2^2=4
//连接操作,Client端支持的一种操作
public static final int OP_CONNECT = 1 << 3; // 2^3=8
//连接可接受操作,仅ServerSocketChannel支持
public static final int OP_ACCEPT = 1 << 4; // 2^4=16
}
Прикрепленный объект (необязательно), объект можно прикрепить при регистрации канала с помощью селектора.
public final SelectionKey register(Selector sel, int ops, Object att)
отselectionKeyЧтобы получить объект вложения, вы можете использоватьattachment()метод
public final Object attachment() {
return attachment;
}
4 Использование NIO для реализации многопользовательских чатов
Далее давайте возьмем практический пример, используя NIO для реализации многопользовательской спортивной версии чата.
Код сервера:
public class GroupChatServer {
private Selector selector;
private ServerSocketChannel serverSocketChannel;
public static final int PORT = 6667;
//构造器初始化成员变量
public GroupChatServer() {
try {
//打开一个选择器
this.selector = Selector.open();
//打开serverSocketChannel
this.serverSocketChannel = ServerSocketChannel.open();
//绑定地址,端口号
this.serverSocketChannel.bind(new InetSocketAddress("127.0.0.1", PORT));
//设置为非阻塞
serverSocketChannel.configureBlocking(false);
//把通道注册到选择器中
serverSocketChannel.register(selector, SelectionKey.OP_ACCEPT);
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 监听,并且接受客户端消息,转发到其他客户端
*/
public void listen() {
try {
while (true) {
//获取监听的事件总数
int count = selector.select(2000);
if (count > 0) {
Set<SelectionKey> selectionKeys = selector.selectedKeys();
//获取SelectionKey集合
Iterator<SelectionKey> it = selectionKeys.iterator();
while (it.hasNext()) {
SelectionKey key = it.next();
//如果是获取连接事件
if (key.isAcceptable()) {
SocketChannel socketChannel = serverSocketChannel.accept();
//设置为非阻塞
socketChannel.configureBlocking(false);
//注册到选择器中
socketChannel.register(selector, SelectionKey.OP_READ);
System.out.println(socketChannel.getRemoteAddress() + "上线了~");
}
//如果是读就绪事件
if (key.isReadable()) {
//读取消息,并且转发到其他客户端
readData(key);
}
it.remove();
}
} else {
System.out.println("等待...");
}
}
} catch (Exception e) {
e.printStackTrace();
}
}
//获取客户端发送过来的消息
private void readData(SelectionKey selectionKey) {
SocketChannel socketChannel = null;
try {
//从selectionKey中获取channel
socketChannel = (SocketChannel) selectionKey.channel();
//创建一个缓冲区
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
//把通道的数据写入到缓冲区
int count = socketChannel.read(byteBuffer);
//判断返回的count是否大于0,大于0表示读取到了数据
if (count > 0) {
//把缓冲区的byte[]转成字符串
String msg = new String(byteBuffer.array());
//输出该消息到控制台
System.out.println("from 客户端:" + msg);
//转发到其他客户端
notifyAllClient(msg, socketChannel);
}
} catch (Exception e) {
try {
//打印离线的通知
System.out.println(socketChannel.getRemoteAddress() + "离线了...");
//取消注册
selectionKey.cancel();
//关闭流
socketChannel.close();
} catch (IOException e1) {
e1.printStackTrace();
}
}
}
/**
* 转发消息到其他客户端
* msg 消息
* noNotifyChannel 不需要通知的Channel
*/
private void notifyAllClient(String msg, SocketChannel noNotifyChannel) throws Exception {
System.out.println("服务器转发消息~");
for (SelectionKey selectionKey : selector.keys()) {
Channel channel = selectionKey.channel();
//channel的类型实际类型是SocketChannel,并且排除不需要通知的通道
if (channel instanceof SocketChannel && channel != noNotifyChannel) {
//强转成SocketChannel类型
SocketChannel socketChannel = (SocketChannel) channel;
//通过消息,包裹获取一个缓冲区
ByteBuffer byteBuffer = ByteBuffer.wrap(msg.getBytes());
socketChannel.write(byteBuffer);
}
}
}
public static void main(String[] args) throws Exception {
GroupChatServer chatServer = new GroupChatServer();
//启动服务器,监听
chatServer.listen();
}
}
Код клиента:
public class GroupChatClinet {
private Selector selector;
private SocketChannel socketChannel;
private String userName;
public GroupChatClinet() {
try {
//打开选择器
this.selector = Selector.open();
//连接服务器
socketChannel = SocketChannel.open(new InetSocketAddress("127.0.0.1", GroupChatServer.PORT));
//设置为非阻塞
socketChannel.configureBlocking(false);
//注册到选择器中
socketChannel.register(selector, SelectionKey.OP_READ);
//获取用户名
userName = socketChannel.getLocalAddress().toString().substring(1);
System.out.println(userName + " is ok~");
} catch (Exception e) {
e.printStackTrace();
}
}
//发送消息到服务端
private void sendMsg(String msg) {
msg = userName + "说:" + msg;
try {
socketChannel.write(ByteBuffer.wrap(msg.getBytes()));
} catch (Exception e) {
e.printStackTrace();
}
}
//读取服务端发送过来的消息
private void readMsg() {
try {
int count = selector.select();
if (count > 0) {
Iterator<SelectionKey> iterator = selector.selectedKeys().iterator();
while (iterator.hasNext()) {
SelectionKey selectionKey = iterator.next();
//判断是读就绪事件
if (selectionKey.isReadable()) {
SocketChannel channel = (SocketChannel) selectionKey.channel();
//创建一个缓冲区
ByteBuffer byteBuffer = ByteBuffer.allocate(1024);
//从服务器的通道中读取数据到缓冲区
channel.read(byteBuffer);
//缓冲区的数据,转成字符串,并打印
System.out.println(new String(byteBuffer.array()));
}
iterator.remove();
}
}
} catch (Exception e) {
e.printStackTrace();
}
}
public static void main(String[] args) throws Exception {
GroupChatClinet chatClinet = new GroupChatClinet();
//启动线程,读取服务器转发过来的消息
new Thread(() -> {
while (true) {
chatClinet.readMsg();
try {
Thread.sleep(3000);
} catch (Exception e) {
e.printStackTrace();
}
}
}).start();
//主线程发送消息到服务器
Scanner scanner = new Scanner(System.in);
while (scanner.hasNextLine()) {
String msg = scanner.nextLine();
chatClinet.sendMsg(msg);
}
}
}
Сначала запустите основной метод сервера, а затем запустите основные методы двух клиентов:
Затем используйте двух клиентов, чтобы начать общение~
Вышеупомянутое - использовать NIOРеализовать чат на несколько человекНапример, учащиеся могут посмотреть мой пример и выполнить его самостоятельно. Для понимания этих концепций требуется больше кода.
напиши в конце
Творить не просто, если вы найдете это полезнымпоставить лайкБар.
Если вы хотите впервые увидеть мои обновленные статьи, вы можете выполнить поиск в общедоступной учетной записи на WeChat "java技术爱好者",Не хочу быть соленой рыбой, я программист, стремящийся запомниться всем. Увидимся в следующий раз! ! !
Возможности ограничены, если есть какие-то ошибки или неуместности, пожалуйста, критикуйте и исправьте их, учитесь и общайтесь вместе!