Пожалуйста, помогите понять, как добиться вышеописанного.
- Селектор можно использовать для мониторинга выбираемых каналов Java NIO, когда канал[ы] готов к событию[ям], которые могут быть Чтение, Запись, Принятие , Connect (в соответствии с зарегистрированным интересом), Selector уведомит нас.
- Мы регистрируем наш Socket в Selector и указываем операции, которые нас интересуют (Чтение/Запись/Подключение/Принятие) в чтобы получить событие Ready с помощью API-интерфейса Register.
- Вызов Register возвращает нам ключ выбора, который представляет связь между сокетом и селектором.
- 'Ключ выбора ' содержит информацию о селекторе, канале, интересующем наборе (операциях, в которых заинтересован канал) и готовом наборе (наборе операций, которые готовы и могут быть выполнены без блокировки, подмножестве интересующего набора).
- "Селектор" поддерживает три набора ключей.
- Отмененный набор ключей:
a. Набор ключей выбора, регистрация которых в селекторе отменена.
b. Отмененный набор ключей не является потокобезопасным. - Зарегистрированный набор ключей:
a. Набор ключей выбора, зарегистрированных в селекторе.
b. Зарегистрированный набор ключей является потокобезопасным. - Выбранный набор ключей:
a. Набор ключей выбора, у которых есть хотя бы одна интересующая операция в состоянии готовности.
b. Выбранный набор ключей не является потокобезопасным.
- Мы блокируем наш поток (NETWORK_IO_IN) на артефакте селектора Вызов select(), который разблокируется, когда хотя бы один из зарегистрированных каналов будет готов выполнить интересующую операцию.
- Вызов select является потокобезопасным, он внутренне включает три шага:
- Сначала он перебирает отмененный набор ключей и освобождает ресурсы относительно канала/сокета, регистрация которых была отменена из селектора.
- На втором этапе он опрашивает ОС, все ли каналы готовы выполнить ввод-вывод, на основе их заинтересованного набора.
- На третьем этапе он заполняет выбранный ключ set.
- Примечание. Все эти шаги внутренне синхронизируются для селектора, затем для отмены и выбранного набора клавиш.
[*]Чтобы получить набор выбранных ключей, который содержит ключ выбора (содержащий информацию о канале), готовый к выполнению ввода-вывода, мы используем selectedKeys(), который возвращает нам ссылку из
выбранного набора ключей, который является закрытым элементом селектора.
- Мы выполняем итерацию по выбранному ключу установите, определите, какая интересующая операция готова (чтение/запись/подключение/принятие) и выполните этот ввод-вывод.
- Здесь 10 потоков заблокированы при вызове выбора артефакта селектора.
- Мы постоянно отправить данные на порт IP:Port, на котором канал/сокет прослушивается и зарегистрирован в селекторе.
- Мы получаем события Read всякий раз, когда канал/сокет готов читать данные из ОС буфер.
- Наблюдение:
- Все потоки блокируются вызов select и выходили один за другим.
- Поток, который выходит из вызова select, скажем, T1 получает выбранный набор ключей и работает над ним.
- Рассмотрение другого потока T2, который заблокирован при вызове выбора и находится на шаге, упомянутом в 7.c, также является работает с одним и тем же выбранным набором ключей.
- Затем мы получаем исключение одновременного изменения, поскольку оба потока пытаются одновременно выполнить изменение выбранного набор ключей.
общественное развлечение AsynIOArtefactLoop () {Код: Выделить всё
while (true) { println ("Thread " + Thread.currentThread().id + "blocked on select") // Threads will get blocked on the select calls until there is at least one channel /socket which is ready. selector.select () println ("Thread " + Thread.currentThread().id + "comes out of select") // Returns the reference of the selected key set. val selected_key_set = selector.selectedKeys() val iterator = selected_key_set.iterator() println ("Thread " + Thread.currentThread().id + “operating on the selected key set”) while (iterator.hasNext()) { // Extract the key (selection key object ) and increment the iterator. val key = iterator.next() // Remove the Seleciton Key object from the selected key set. iterator.remove() println (" Event receive corrsponding to descriptor id " + key.attachment() as Int) when { // Channel/Socket is ready for Receive. key.isReadable -> { … // Buffer in which to copy the receive data. var buffer = ByteBuffer.allocate(1024) // Get the channel from Selection key object and call the receive. (key.channel() as DatagramChannel).receive(buffer) … // Print the Receive data . println ("Message Received: ${String(receivedData)}") } key.isWritable -> println("Write Ready") } } }
- Здесь также блокируются 10 потоков при вызове select артефакта селектора.
- Мы постоянно отправляем данные к IP:порту, на котором канал/сокет прослушивается и зарегистрирован в селекторе.
- Мы получаем события чтения всякий раз, когда канал/сокет готов читать данные из буфера ОС.< /li>
Наблюдение:
- Все потоки блокируются при вызове select и выходят наружу один за другим.
- Поток, который выходит из вызова выбора, скажем, T1, получает блокировку selected_key_set, ссылка на который возвращается вызовом selectedKeys().
- Учитывая другой поток T2, который заблокирован при вызове выбора и находится на шаге, упомянутом в 7.c, если один из потоков уже заблокировал выбранный_key_set, он не сможет работать над ним, что синхронизируется. весь процесс итерации с вызовом select ().
Код: Выделить всё
while (true) {
println ("Thread " + Thread.currentThread().id + "blocked on select")
selector.select ()
println ("Thread " + Thread.currentThread().id + "comes out of select")
val selected_key_set = selector.selectedKeys()
// Here we have synchronize on the selected_key_set which ensures that thread which are blocked on the select call will not be able to use it if another thread is already operating over it.
synchronized (selected_key_set) {
println ("Thread " + Thread.currentThread().id + "blocked on selected key set")
val iterator = selected_key_set.iterator()
while (iterator.hasNext()) {
val key = iterator.next()
println (" Event receive corrsponding to descriptor id " + key.attachment() as Int)
iterator.remove()
when {
key.isReadable -> {
counter++;
var buffer = ByteBuffer.allocate(1024)
(key.channel() as DatagramChannel).receive(buffer)
buffer.flip()
val receivedData = ByteArray(buffer.remaining())
buffer.get(receivedData)
println ("Recv Event Number: "+ counter)
println ("Message Received: ${String(receivedData)}")
}
key.isWritable -> println("Write Ready")
}
}
}
println ("Thread " + Thread.currentThread().id + "comes out of selected key set")
}
Здесь мы выполняем синхронизацию выбранного набора ключей, когда мы выполняем обработку событий готовности (чтение/запись/принятие/подключение). ), что позволяет одновременно работать только одному потоку и блокирует другие потоки.
Если у меня многоядерный процессор, я не смогу воспользоваться преимуществом производительности, потому что только один моего потока будет выполнять обработку одновременно, а другие будут ждать этого.
Подробнее здесь: https://stackoverflow.com/questions/790 ... s-on-selec