Spec-Zone.ru › Kotlin 1.8

Выражение 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 с исключением. Мы можем использовать предложение onReceiveCatching, чтобы выполнить определенное действие при закрытии канала. Следующий пример также показывает, что select — это выражение, возвращающее результат выбранного предложения:

suspend fun selectAorB(a: ReceiveChannel<String>, b: ReceiveChannel<String>): String =
    select<String> {
        a.onReceiveCatching { it ->
            val value = it.getOrNull()
            if (value != null) {
                "a -> '$value'"
            } else {
                "Channel 'a' is closed"
            }
        }
        b.onReceiveCatching { it ->
            val value = it.getOrNull()
            if (value != null) {
                "b -> '$value'"
            } else {
                "Channel 'b' is closed"
            }
        }
    }

Давайте воспользуемся каналом 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.onReceiveCatching { it ->
            val value = it.getOrNull()
            if (value != null) {
                "a -> '$value'"
            } else {
                "Channel 'a' is closed"
            }
        }
        b.onReceiveCatching { it ->
            val value = it.getOrNull()
            if (value != null) {
                "b -> '$value'"
            } else {
                "Channel 'b' is closed"
            }
        }
    }
    
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 отправить данные тоже.

Второе наблюдение заключается в том, что onReceiveCatching выбирается немедленно, когда канал уже закрыт.

Выбор для отправки

Выражение 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

Переключение на канал отложенных значений

Давайте напишем функцию производителя канала, которая потребляет канал отложенных строковых значений, ожидает каждого полученного отложенного значения, но только до тех пор, пока не появится следующее отложенное значение или канал не будет закрыт. Этот пример объединяет предложения onReceiveCatching и 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.onReceiveCatching { update ->
                update.getOrNull()
            }
            current.onAwait { value ->
                send(value) // send value that current deferred has produced
                input.receiveCatching().getOrNull() // 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.onReceiveCatching { update ->
                update.getOrNull()
            }
            current.onAwait { value ->
                send(value) // send value that current deferred has produced
                input.receiveCatching().getOrNull() // 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
Последнее изменение: 10 января 2023
Общий изменяемый и конкурентный доступ к состоянию Отладка корутин с помощью IntelliJ IDEA — учебник

© 2010–2023 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/select-expression.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API