Spec-Zone.ru › Kotlin 2

Выражение select (экспериментальная функция)

Выражение select позволяет одновременно ожидать несколько приостанавливающих функций и выбирать первую из них, которая становится доступной.

Выражения select — экспериментальная функция kotlinx.coroutines. Ожидается, что их API будет меняться в будущих обновлениях библиотеки kotlinx.coroutines, в том числе с потенциально несовместимыми изменениями.

Выбор из каналов

Предположим, у нас есть два производителя строк: fizz и buzz. fizz производит строку "Fizz" каждые 500 мс:

fun CoroutineScope.fizz() = produce<String> {
    while (true) { // sends "Fizz" every 500 ms
        delay(500)
        send("Fizz")
    }
}

А buzz производит строку "Buzz!" каждые 1000 мс:

fun CoroutineScope.buzz() = produce<String> {
    while (true) { // sends "Buzz!" every 1000 ms
        delay(1000)
        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 500 ms
        delay(500)
        send("Fizz")
    }
}

fun CoroutineScope.buzz() = produce<String> {
    while (true) { // sends "Buzz!" every 1000 ms
        delay(1000)
        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'
fizz -> 'Fizz'

Выбор при закрытии

Конструкция 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
27 сентября 2024 г.
Общее изменяемое состояние и параллелизмОтладка сопрограмм

© 2010–2026 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