Содержание
Выражение select (экспериментальное)
Выражение select позволяет одновременно ожидать выполнения нескольких приостанавливаемых функций и выбрать первую, которая станет доступной.
Выражения select являются экспериментальной функцией
kotlinx.coroutines. Их API ожидается, что будет развиваться в будущих обновлениях библиотекиkotlinx.coroutinesс потенциально несовместимыми изменениями.
Выбор из каналов
Представим два производителя строк: fizz и buzz. fizz производит строку "Fizz" каждые 300 мс:
fun CoroutineScope.fizz() = produce<String> {
while (true) { // sends "Fizz" every 300 ms
delay(300)
send("Fizz")
}
}
А buzz производит строку "Buzz!" каждые 500 мс:
fun CoroutineScope.buzz() = produce<String> {
while (true) { // sends "Buzz!" every 500 ms
delay(500)
send("Buzz!")
}
}
Используя приостанавливаемую функцию receive, мы можем получить значение либо из одного канала, либо из другого. Но выражение select позволяет получить значение из обоих каналов одновременно, используя его пункты onReceive:
suspend fun selectFizzBuzz(fizz: ReceiveChannel<String>, buzz: ReceiveChannel<String>) {
select<Unit> { // <Unit> means that this select expression does not produce any result
fizz.onReceive { value -> // this is the first select clause
println("fizz -> '$value'")
}
buzz.onReceive { value -> // this is the second select clause
println("buzz -> '$value'")
}
}
}
Давайте запустим его семь раз:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
import kotlinx.coroutines.selects.*
fun CoroutineScope.fizz() = produce<String> {
while (true) { // sends "Fizz" every 300 ms
delay(300)
send("Fizz")
}
}
fun CoroutineScope.buzz() = produce<String> {
while (true) { // sends "Buzz!" every 500 ms
delay(500)
send("Buzz!")
}
}
suspend fun selectFizzBuzz(fizz: ReceiveChannel<String>, buzz: ReceiveChannel<String>) {
select<Unit> { // <Unit> means that this select expression does not produce any result
fizz.onReceive { value -> // this is the first select clause
println("fizz -> '$value'")
}
buzz.onReceive { value -> // this is the second select clause
println("buzz -> '$value'")
}
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val fizz = fizz()
val buzz = buzz()
repeat(7) {
selectFizzBuzz(fizz, buzz)
}
coroutineContext.cancelChildren() // cancel fizz & buzz coroutines
//sampleEnd
}
Полный код можно найти здесь.
Результат выполнения этого кода:
fizz -> 'Fizz' buzz -> 'Buzz!' fizz -> 'Fizz' fizz -> 'Fizz' buzz -> 'Buzz!' fizz -> 'Fizz' buzz -> 'Buzz!'
Выбор при закрытии
Пункт onReceive в select завершается ошибкой, когда канал закрыт, вызывая соответствующее исключение в select. Мы можем использовать пункт onReceiveOrNull, чтобы выполнить определённое действие при закрытии канала. Следующий пример также демонстрирует, что select — это выражение, возвращающее результат выбранного пункта:
suspend fun selectAorB(a: ReceiveChannel<String>, b: ReceiveChannel<String>): String =
select<String> {
a.onReceiveOrNull { value ->
if (value == null)
"Channel 'a' is closed"
else
"a -> '$value'"
}
b.onReceiveOrNull { value ->
if (value == null)
"Channel 'b' is closed"
else
"b -> '$value'"
}
}
Обратите внимание, что onReceiveOrNull является функцией расширения, определённой только для каналов с не-null элементами, чтобы избежать случайного смешения закрытия канала и null значения.
Давайте используем его с каналом a, который производит строку "Hello" четыре раза, и каналом b, который производит строку "World" четыре раза:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
import kotlinx.coroutines.selects.*
suspend fun selectAorB(a: ReceiveChannel<String>, b: ReceiveChannel<String>): String =
select<String> {
a.onReceiveOrNull { value ->
if (value == null)
"Channel 'a' is closed"
else
"a -> '$value'"
}
b.onReceiveOrNull { value ->
if (value == null)
"Channel 'b' is closed"
else
"b -> '$value'"
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val a = produce<String> {
repeat(4) { send("Hello $it") }
}
val b = produce<String> {
repeat(4) { send("World $it") }
}
repeat(8) { // print first eight results
println(selectAorB(a, b))
}
coroutineContext.cancelChildren()
//sampleEnd
}
Полный код можно найти здесь.
Результат выполнения этого кода довольно интересный, поэтому мы подробно его проанализируем:
a -> 'Hello 0' a -> 'Hello 1' b -> 'World 0' a -> 'Hello 2' a -> 'Hello 3' b -> 'World 1' Channel 'a' is closed Channel 'a' is closed
Есть несколько наблюдений.
Во-первых, select имеет предпочтение к первому пункту. Когда несколько пунктов могут быть выбраны одновременно, выбирается первый из них. Здесь оба канала постоянно производят строки, поэтому канал a, являясь первым пунктом в select, побеждает. Однако, так как мы используем небуферизованный канал, a иногда приостанавливается при вызове send и даёт шанс на отправку значений в b.
Второе наблюдение состоит в том, что onReceiveOrNull выбирается сразу, когда канал уже закрыт.
Выбор для отправки
Выражение select имеет пункт onSend, который может быть полезен в сочетании с предвзятой природой выбора.
Давайте напишем пример производителя целых чисел, который отправляет свои значения в канал side , когда потребители в основном канале не успевают за ним:
fun CoroutineScope.produceNumbers(side: SendChannel<Int>) = produce<Int> {
for (num in 1..10) { // produce 10 numbers from 1 to 10
delay(100) // every 100 ms
select<Unit> {
onSend(num) {} // Send to the primary channel
side.onSend(num) {} // or to the side channel
}
}
}
Потребитель будет довольно медленным, тратя 250 мс на обработку каждого числа:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
import kotlinx.coroutines.selects.*
fun CoroutineScope.produceNumbers(side: SendChannel<Int>) = produce<Int> {
for (num in 1..10) { // produce 10 numbers from 1 to 10
delay(100) // every 100 ms
select<Unit> {
onSend(num) {} // Send to the primary channel
side.onSend(num) {} // or to the side channel
}
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val side = Channel<Int>() // allocate side channel
launch { // this is a very fast consumer for the side channel
side.consumeEach { println("Side channel has $it") }
}
produceNumbers(side).consumeEach {
println("Consuming $it")
delay(250) // let us digest the consumed number properly, do not hurry
}
println("Done consuming")
coroutineContext.cancelChildren()
//sampleEnd
}
Полный код можно найти здесь.
Давайте посмотрим, что произойдёт:
Consuming 1 Side channel has 2 Side channel has 3 Consuming 4 Side channel has 5 Side channel has 6 Consuming 7 Side channel has 8 Side channel has 9 Consuming 10 Done consuming
Выбор отложенных значений
Отложенные значения можно выбрать, используя пункт onAwait. Давайте начнём с асинхронной функции, которая возвращает отложенное строковое значение после случайной задержки:
fun CoroutineScope.asyncString(time: Int) = async {
delay(time.toLong())
"Waited for $time ms"
}
Давайте запустим дюжину из них со случайной задержкой.
fun CoroutineScope.asyncStringsList(): List<Deferred<String>> {
val random = Random(3)
return List(12) { asyncString(random.nextInt(1000)) }
}
Теперь главная функция ждёт завершения первого из них и подсчитывает количество отложенных значений, которые всё ещё активны. Обратите внимание, что мы использовали тот факт, что выражение select это Kotlin DSL, поэтому мы можем предоставить пункты для него с помощью произвольного кода. В этом случае мы проходим по списку отложенных значений, чтобы предоставить пункт onAwait для каждого отложенного значения.
import kotlinx.coroutines.*
import kotlinx.coroutines.selects.*
import java.util.*
fun CoroutineScope.asyncString(time: Int) = async {
delay(time.toLong())
"Waited for $time ms"
}
fun CoroutineScope.asyncStringsList(): List<Deferred<String>> {
val random = Random(3)
return List(12) { asyncString(random.nextInt(1000)) }
}
fun main() = runBlocking<Unit> {
//sampleStart
val list = asyncStringsList()
val result = select<String> {
list.withIndex().forEach { (index, deferred) ->
deferred.onAwait { answer ->
"Deferred $index produced answer '$answer'"
}
}
}
println(result)
val countActive = list.count { it.isActive }
println("$countActive coroutines are still active")
//sampleEnd
}
Полный код можно найти здесь.
Вывод:
Deferred 4 produced answer 'Waited for 128 ms' 11 coroutines are still active
Переключение по каналу отложенных значений
Давайте напишем функцию-производителя канала, которая потребляет канал отложенных строковых значений, ждёт получения каждого отложенного значения, но только до тех пор, пока не придёт следующее отложенное значение или канал не будет закрыт. Этот пример объединяет пункты onReceiveOrNull и onAwait в одном select.
fun CoroutineScope.switchMapDeferreds(input: ReceiveChannel<Deferred<String>>) = produce<String> {
var current = input.receive() // start with first received deferred value
while (isActive) { // loop while not cancelled/closed
val next = select<Deferred<String>?> { // return next deferred value from this select or null
input.onReceiveOrNull { update ->
update // replaces next value to wait
}
current.onAwait { value ->
send(value) // send value that current deferred has produced
input.receiveOrNull() // and use the next deferred from the input channel
}
}
if (next == null) {
println("Channel was closed")
break // out of loop
} else {
current = next
}
}
}
Для тестирования мы будем использовать простую асинхронную функцию, которая возвращает заданную строку через заданное время:
fun CoroutineScope.asyncString(str: String, time: Long) = async {
delay(time)
str
}
Главная функция просто запускает сопрограмму для вывода результатов switchMapDeferreds и отправляет некоторые тестовые данные ей:
import kotlinx.coroutines.*
import kotlinx.coroutines.channels.*
import kotlinx.coroutines.selects.*
fun CoroutineScope.switchMapDeferreds(input: ReceiveChannel<Deferred<String>>) = produce<String> {
var current = input.receive() // start with first received deferred value
while (isActive) { // loop while not cancelled/closed
val next = select<Deferred<String>?> { // return next deferred value from this select or null
input.onReceiveOrNull { update ->
update // replaces next value to wait
}
current.onAwait { value ->
send(value) // send value that current deferred has produced
input.receiveOrNull() // and use the next deferred from the input channel
}
}
if (next == null) {
println("Channel was closed")
break // out of loop
} else {
current = next
}
}
}
fun CoroutineScope.asyncString(str: String, time: Long) = async {
delay(time)
str
}
fun main() = runBlocking<Unit> {
//sampleStart
val chan = Channel<Deferred<String>>() // the channel for test
launch { // launch printing coroutine
for (s in switchMapDeferreds(chan))
println(s) // print each received string
}
chan.send(asyncString("BEGIN", 100))
delay(200) // enough time for "BEGIN" to be produced
chan.send(asyncString("Slow", 500))
delay(100) // not enough time to produce slow
chan.send(asyncString("Replace", 100))
delay(500) // give it time before the last one
chan.send(asyncString("END", 500))
delay(1000) // give it time to process
chan.close() // close the channel ...
delay(500) // and wait some time to let it finish
//sampleEnd
}
Полный код можно найти здесь.
Результат выполнения этого кода:
BEGIN Replace END Channel was closed
© 2010–2020 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/reference/coroutines/select-expression.html