Асинхронный поток
Функция приостановления асинхронно возвращает единственное значение, но как мы можем вернуть несколько асинхронно вычисленных значений? Здесь на помощь приходят Kotlin Flows.
Представление нескольких значений
Несколько значений можно представить в Kotlin с помощью коллекций. Например, мы можем иметь функцию simple, которая возвращает список из трёх чисел, а затем распечатать их все с помощью forEach:
fun simple(): List<Int> = listOf(1, 2, 3)
fun main() {
simple().forEach { value -> println(value) }
}
Этот код выводит:
1 2 3
Последовательности
Если мы вычисляем числа с помощью некоторого ресурсоёмкого блокирующего кода (каждое вычисление занимает 100 мс), то мы можем представить числа с помощью последовательности:
fun simple(): Sequence<Int> = sequence { // sequence builder
for (i in 1..3) {
Thread.sleep(100) // pretend we are computing it
yield(i) // yield next value
}
}
fun main() {
simple().forEach { value -> println(value) }
}
Этот код выводит те же числа, но ожидает 100 мс перед печатью каждого из них.
Функции приостановки
Однако, это вычисление блокирует главную нить, которая выполняет код. Когда эти значения вычисляются асинхронным кодом, мы можем пометить функцию simple модификатором suspend, чтобы она могла выполнять свою работу без блокировки и возвращать результат в виде списка:
import kotlinx.coroutines.*
//sampleStart
suspend fun simple(): List<Int> {
delay(1000) // pretend we are doing something asynchronous here
return listOf(1, 2, 3)
}
fun main() = runBlocking<Unit> {
simple().forEach { value -> println(value) }
}
//sampleEnd
Этот код выводит числа после ожидания секунды.
Потоки
Использование типа результата List<Int> означает, что мы можем вернуть все значения только сразу. Чтобы представить поток значений, которые вычисляются асинхронно, мы можем использовать тип Flow<Int>, так же как мы бы использовали тип Sequence<Int> для синхронно вычисленных значений:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow { // flow builder
for (i in 1..3) {
delay(100) // pretend we are doing something useful here
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
// Launch a concurrent coroutine to check if the main thread is blocked
launch {
for (k in 1..3) {
println("I'm not blocked $k")
delay(100)
}
}
// Collect the flow
simple().collect { value -> println(value) }
}
//sampleEnd
Этот код ожидает 100 мс перед печатью каждого числа без блокировки основной нити. Это подтверждается печатью «Я не заблокирован» каждые 100 мс из отдельной корутины, выполняемой в основной нити:
I'm not blocked 1 1 I'm not blocked 2 2 I'm not blocked 3 3
Обратите внимание на следующие различия в коде с Flow из предыдущих примеров:
Потоки — холодные потоки
Потоки — это холодные потоки, похожие на последовательности — код внутри flow-строителя не выполняется, пока поток не собран. Это становится очевидным в следующем примере:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
println("Flow started")
for (i in 1..3) {
delay(100)
emit(i)
}
}
fun main() = runBlocking<Unit> {
println("Calling simple function...")
val flow = simple()
println("Calling collect...")
flow.collect { value -> println(value) }
println("Calling collect again...")
flow.collect { value -> println(value) }
}
//sampleEnd
Что выводит:
Calling simple function... Calling collect... Flow started 1 2 3 Calling collect again... Flow started 1 2 3
Это ключевая причина, по которой функция simple (которая возвращает поток) не помечена модификатором suspend. Сама по себе функция simple() возвращается быстро и не ждет ничего. Поток запускается каждый раз, когда он собирается, поэтому мы видим «Поток запущен», когда мы вызываем collect снова.
Основы отмены потока
Поток соответствует общему принципу кооперативной отмены корутин. Как обычно, сбор потока может быть отменён, когда поток приостановлен в функции приостановки с возможностью отмены (например, delay). Следующий пример показывает, как поток отменяется по таймауту при выполнении в блоке withTimeoutOrNull и прекращает выполнение своего кода:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
delay(100)
println("Emitting $i")
emit(i)
}
}
fun main() = runBlocking<Unit> {
withTimeoutOrNull(250) { // Timeout after 250ms
simple().collect { value -> println(value) }
}
println("Done")
}
//sampleEnd
Обратите внимание, как только два числа выводятся потоком в функции simple, что даёт следующий вывод:
Emitting 1 1 Emitting 2 2 Done
См. раздел Проверки отмены потока для получения более подробной информации.
Строители потоков
Строитель flow { ... } из предыдущих примеров является самым основным. Существуют другие строители для более удобной декларации потоков:
flowOf — строитель, который определяет поток, выводящий фиксированный набор значений.
Различные коллекции и последовательности могут быть преобразованы в потоки с помощью функций расширения
.asFlow().
Итак, пример, который выводит числа от 1 до 3 из потока, можно записать как:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
// Convert an integer range to a flow
(1..3).asFlow().collect { value -> println(value) }
//sampleEnd
}
Промежуточные операторы потоков
Потоки могут быть преобразованы с помощью операторов, точно так же, как это делается со списками и последовательностями. Промежуточные операторы применяются к входному потоку и возвращают выходной поток. Эти операторы являются ленивыми, как и потоки. Вызов такого оператора не является приостанавливаемой функцией. Он работает быстро, возвращая определение нового преобразованного потока.
Основные операторы имеют знакомые названия, такие как map и filter. Важное отличие от последовательностей заключается в том, что блоки кода внутри этих операторов могут вызывать приостанавливаемые функции.
Например, поток входящих запросов может быть преобразован в результаты с помощью оператора map, даже если выполнение запроса является длительной операцией, реализованной приостанавливаемой функцией:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
suspend fun performRequest(request: Int): String {
delay(1000) // imitate long-running asynchronous work
return "response $request"
}
fun main() = runBlocking<Unit> {
(1..3).asFlow() // a flow of requests
.map { request -> performRequest(request) }
.collect { response -> println(response) }
}
//sampleEnd
Он генерирует следующие три строки, каждая из которых появляется через секунду:
response 1 response 2 response 3
Оператор преобразования
Среди операторов преобразования потоков наиболее общим является оператор transform. Он может использоваться для имитации простых преобразований, таких как map и filter, а также для реализации более сложных преобразований. Используя оператор transform, мы можем выводить произвольные значения произвольное количество раз.
Например, используя transform, мы можем вывести строку перед выполнением длительного асинхронного запроса и последовать за ней ответом:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
suspend fun performRequest(request: Int): String {
delay(1000) // imitate long-running asynchronous work
return "response $request"
}
fun main() = runBlocking<Unit> {
//sampleStart
(1..3).asFlow() // a flow of requests
.transform { request ->
emit("Making request $request")
emit(performRequest(request))
}
.collect { response -> println(response) }
//sampleEnd
}
Выводом этого кода является:
Making request 1 response 1 Making request 2 response 2 Making request 3 response 3
Операторы ограничения размера
Промежуточные операторы ограничения размера, такие как take, отменяют выполнение потока, когда достигается соответствующее ограничение. Отмена в сопроцессах всегда выполняется путем выброса исключения, так что все функции управления ресурсами (например, try { ... } finally { ... } блоки) работают нормально в случае отмены:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun numbers(): Flow<Int> = flow {
try {
emit(1)
emit(2)
println("This line will not execute")
emit(3)
} finally {
println("Finally in numbers")
}
}
fun main() = runBlocking<Unit> {
numbers()
.take(2) // take only the first two
.collect { value -> println(value) }
}
//sampleEnd
Вывод этого кода ясно показывает, что выполнение flow { ... } тела в функции numbers() остановилось после вывода второго числа:
1 2 Finally in numbers
Конечные операторы потоков
Конечные операторы для потоков — это приостанавливаемые функции, которые запускают сбор потока. Оператор collect является самым основным, но существуют и другие конечные операторы, которые могут упростить задачу:
Преобразование в различные коллекции, такие как toList и toSet.
Операторы для получения первого значения и для обеспечения того, что поток выводит единственное значение.
Например:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
val sum = (1..5).asFlow()
.map { it * it } // squares of numbers from 1 to 5
.reduce { a, b -> a + b } // sum them (terminal operator)
println(sum)
//sampleEnd
}
Выводит одно число:
55
Потоки являются последовательными
Каждая отдельная коллекция потока выполняется последовательно, если не используются специальные операторы, работающие с несколькими потоками. Сбор выполняется непосредственно в сопроцессе, который вызывает конечный оператор. По умолчанию новые сопроцессы не запускаются. Каждое выводимое значение обрабатывается всеми промежуточными операторами сверху вниз, а затем передаётся в конечный оператор.
Рассмотрим следующий пример, который фильтрует чётные целые числа и преобразует их в строки:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
(1..5).asFlow()
.filter {
println("Filter $it")
it % 2 == 0
}
.map {
println("Map $it")
"string $it"
}.collect {
println("Collect $it")
}
//sampleEnd
}
Результат:
Filter 1 Filter 2 Map 2 Collect string 2 Filter 3 Filter 4 Map 4 Collect string 4 Filter 5
Контекст потока
Сбор потока всегда происходит в контексте вызывающей корутины. Например, если есть поток simple, то следующий код выполняется в контексте, указанном автором этого кода, независимо от реализации деталей потока simple:
withContext(context) {
simple().collect { value ->
println(value) // run in the specified context
}
}
Это свойство потока называется сохранением контекста.
Таким образом, по умолчанию код в билдере flow { ... } выполняется в контексте, предоставляемом коллектором соответствующего потока. Например, рассмотрим реализацию функции simple, которая выводит поток, в котором она вызывается, и генерирует три числа:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun log(msg: String) = println("[${Thread.currentThread().name}] $msg")
//sampleStart
fun simple(): Flow<Int> = flow {
log("Started simple flow")
for (i in 1..3) {
emit(i)
}
}
fun main() = runBlocking<Unit> {
simple().collect { value -> log("Collected $value") }
}
//sampleEnd
Выполнение этого кода даёт:
[main @coroutine#1] Started simple flow [main @coroutine#1] Collected 1 [main @coroutine#1] Collected 2 [main @coroutine#1] Collected 3
Поскольку simple().collect вызывается из основного потока, тело потока simple также вызывается в главном потоке. Это идеальное значение по умолчанию для быстро работающего или асинхронного кода, который не заботится о контексте выполнения и не блокирует вызывающего.
Неправильная генерация с withContext
Однако, длительный код, потребляющий ресурсы процессора, может потребовать выполнения в контексте Dispatchers.Default, а код для обновления пользовательского интерфейса может потребовать выполнения в контексте Dispatchers.Main. Обычно, withContext используется для изменения контекста в коде, использующем Kotlin coroutines, но код в билдере flow { ... } должен соблюдать свойство сохранения контекста и не должен разрешать генерацию из другого контекста.
Попробуйте запустить следующий код:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
// The WRONG way to change context for CPU-consuming code in flow builder
kotlinx.coroutines.withContext(Dispatchers.Default) {
for (i in 1..3) {
Thread.sleep(100) // pretend we are computing it in CPU-consuming way
emit(i) // emit next value
}
}
}
fun main() = runBlocking<Unit> {
simple().collect { value -> println(value) }
}
//sampleEnd
Этот код генерирует следующую ошибку:
Exception in thread "main" java.lang.IllegalStateException: Flow invariant is violated:
Flow was collected in [CoroutineId(1), "coroutine#1":BlockingCoroutine{Active}@5511c7f8, BlockingEventLoop@2eac3323],
but emission happened in [CoroutineId(1), "coroutine#1":DispatchedCoroutine{Active}@2dae0000, Dispatchers.Default].
Please refer to 'flow' documentation or use 'flowOn' instead
at ...
Оператор flowOn
Ошибка относится к функции flowOn, которая должна использоваться для изменения контекста генерации потока. Правильный способ изменения контекста потока показан в примере ниже, который также выводит имена соответствующих потоков, чтобы продемонстрировать, как всё работает:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun log(msg: String) = println("[${Thread.currentThread().name}] $msg")
//sampleStart
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
Thread.sleep(100) // pretend we are computing it in CPU-consuming way
log("Emitting $i")
emit(i) // emit next value
}
}.flowOn(Dispatchers.Default) // RIGHT way to change context for CPU-consuming code in flow builder
fun main() = runBlocking<Unit> {
simple().collect { value ->
log("Collected $value")
}
}
//sampleEnd
Обратите внимание, как flow { ... } работает в фоновом потоке, в то время как сбор происходит в главном потоке:
Ещё один момент, на который стоит обратить внимание, заключается в том, что оператор flowOn изменил по умолчанию последовательный характер потока. Теперь сбор происходит в одной корутине ("coroutine#1"), а генерация происходит в другой корутине ("coroutine#2"), которая работает в другом потоке одновременно с собирающей корутиной. Оператор flowOn создаёт другую корутину для входящего потока, когда ему необходимо изменить CoroutineDispatcher в своём контексте.
Буферизация
Запуск различных частей потока в различных корутинах может быть полезным с точки зрения общего времени, необходимого для сбора потока, особенно при участии длительных асинхронных операций. Например, рассмотрим случай, когда генерация потока simple медленная, занимая 100 мс для создания элемента; и коллектор также медленный, занимая 300 мс для обработки элемента. Посмотрим, сколько времени занимает сбор такого потока с тремя числами:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.system.*
//sampleStart
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
delay(100) // pretend we are asynchronously waiting 100 ms
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
val time = measureTimeMillis {
simple().collect { value ->
delay(300) // pretend we are processing it for 300 ms
println(value)
}
}
println("Collected in $time ms")
}
//sampleEnd
Это даёт результат примерно такой, при котором весь сбор занимает около 1200 мс (три числа, 400 мс для каждого):
1 2 3 Collected in 1220 ms
Мы можем использовать оператор buffer над потоком, чтобы запустить код генерации потока simple параллельно с кодом сбора, а не последовательно:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.system.*
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
delay(100) // pretend we are asynchronously waiting 100 ms
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val time = measureTimeMillis {
simple()
.buffer() // buffer emissions, don't wait
.collect { value ->
delay(300) // pretend we are processing it for 300 ms
println(value)
}
}
println("Collected in $time ms")
//sampleEnd
}
Это даёт те же числа, но быстрее, так как мы фактически создали конвейер обработки, ожидая только 100 мс для первого числа, а затем тратя только 300 мс на обработку каждого числа. Таким образом, это занимает около 1000 мс для выполнения:
1 2 3 Collected in 1071 ms
Конфляция
Когда поток представляет частичные результаты операции или обновления статуса операции, может не потребоваться обрабатывать каждое значение, а только самые последние. В этом случае оператор conflate может использоваться для пропуска промежуточных значений, когда коллектор слишком медленный для их обработки. Исходя из предыдущего примера:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.system.*
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
delay(100) // pretend we are asynchronously waiting 100 ms
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val time = measureTimeMillis {
simple()
.conflate() // conflate emissions, don't process each one
.collect { value ->
delay(300) // pretend we are processing it for 300 ms
println(value)
}
}
println("Collected in $time ms")
//sampleEnd
}
Мы видим, что в то время как первое число ещё обрабатывалось, второе и третье уже были произведены, поэтому второе было скомбинировано, и только самое последнее (третье) было доставлено коллектору:
1 3 Collected in 758 ms
Обработка последнего значения
Конфляция — один из способов ускорить обработку, когда и эмиттер, и коллектор медленные. Она делает это, отбрасывая выпущенные значения. Другой способ — отменить медленного коллектора и перезапустить его каждый раз, когда генерируется новое значение. Существует семейство xxxLatest операторов, которые выполняют ту же основную логику оператора xxx, но отменяют код в своём блоке при новом значении. Попробуем изменить conflate на collectLatest в предыдущем примере:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.system.*
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
delay(100) // pretend we are asynchronously waiting 100 ms
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
//sampleStart
val time = measureTimeMillis {
simple()
.collectLatest { value -> // cancel & restart on the latest value
println("Collecting $value")
delay(300) // pretend we are processing it for 300 ms
println("Done $value")
}
}
println("Collected in $time ms")
//sampleEnd
}
Так как тело collectLatest занимает 300 мс, а новые значения генерируются каждые 100 мс, мы видим, что блок выполняется для каждого значения, но завершается только для последнего значения:
Collecting 1 Collecting 2 Collecting 3 Done 3 Collected in 741 ms
Компоновка нескольких потоков
Существует множество способов компоновки нескольких потоков.
Zip
Так же, как и функция расширения Sequence.zip в стандартной библиотеке Kotlin, потоки имеют оператор zip, который объединяет соответствующие значения двух потоков:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
val nums = (1..3).asFlow() // numbers 1..3
val strs = flowOf("one", "two", "three") // strings
nums.zip(strs) { a, b -> "$a -> $b" } // compose a single string
.collect { println(it) } // collect and print
//sampleEnd
}
Этот пример выводит:
1 -> one 2 -> two 3 -> three
Combine
Когда поток представляет последнее значение переменной или операции (см. также соответствующий раздел о слиянии), может потребоваться выполнить вычисление, зависящее от последних значений соответствующих потоков, и пересчитывать его всякий раз, когда любой из потоков-источников испускает значение. Соответствующая группа операторов называется combine.
Например, если числа в предыдущем примере обновляются каждые 300 мс, а строки — каждые 400 мс, то использование оператора zip по-прежнему даст тот же результат, но результаты будут выводиться каждые 400 мс:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
val nums = (1..3).asFlow().onEach { delay(300) } // numbers 1..3 every 300 ms
val strs = flowOf("one", "two", "three").onEach { delay(400) } // strings every 400 ms
val startTime = System.currentTimeMillis() // remember the start time
nums.zip(strs) { a, b -> "$a -> $b" } // compose a single string with "zip"
.collect { value -> // collect and print
println("$value at ${System.currentTimeMillis() - startTime} ms from start")
}
//sampleEnd
}
Однако, при использовании оператора combine вместо оператора zip:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking<Unit> {
//sampleStart
val nums = (1..3).asFlow().onEach { delay(300) } // numbers 1..3 every 300 ms
val strs = flowOf("one", "two", "three").onEach { delay(400) } // strings every 400 ms
val startTime = System.currentTimeMillis() // remember the start time
nums.combine(strs) { a, b -> "$a -> $b" } // compose a single string with "combine"
.collect { value -> // collect and print
println("$value at ${System.currentTimeMillis() - startTime} ms from start")
}
//sampleEnd
}
Мы получаем совершенно другой вывод, где строка печатается при каждом испускании из потоков nums или strs:
1 -> one at 452 ms from start 2 -> one at 651 ms from start 2 -> two at 854 ms from start 3 -> two at 952 ms from start 3 -> three at 1256 ms from start
Сглаживание потоков
Потоки представляют асинхронно полученные последовательности значений, поэтому довольно легко попасть в ситуацию, где каждое значение инициирует запрос на другую последовательность значений. Например, у нас может быть следующая функция, которая возвращает поток из двух строк с интервалом в 500 мс:
fun requestFlow(i: Int): Flow<String> = flow {
emit("$i: First")
delay(500) // wait 500 ms
emit("$i: Second")
}
Теперь, если у нас есть поток из трех целых чисел и мы вызываем requestFlow для каждого из них так:
(1..3).asFlow().map { requestFlow(it) }
Тогда мы получаем поток из потоков (Flow<Flow<String>>), который необходимо сгладить в один поток для дальнейшей обработки. Коллекции и последовательности имеют операторы flatten и flatMap для этого. Однако, из-за асинхронной природы потоков они требуют разных режимов сглаживания, поэтому существует множество операторов сглаживания для потоков.
flatMapConcat
Режим конкатенации реализован операторами flatMapConcat и flattenConcat. Они являются наиболее прямыми аналогами соответствующих операторов для последовательностей. Они ожидают завершения внутреннего потока перед началом сбора следующего, как показано в следующем примере:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun requestFlow(i: Int): Flow<String> = flow {
emit("$i: First")
delay(500) // wait 500 ms
emit("$i: Second")
}
fun main() = runBlocking<Unit> {
//sampleStart
val startTime = System.currentTimeMillis() // remember the start time
(1..3).asFlow().onEach { delay(100) } // a number every 100 ms
.flatMapConcat { requestFlow(it) }
.collect { value -> // collect and print
println("$value at ${System.currentTimeMillis() - startTime} ms from start")
}
//sampleEnd
}
Последовательная природа flatMapConcat четко видна в выводе:
1: First at 121 ms from start 1: Second at 622 ms from start 2: First at 727 ms from start 2: Second at 1227 ms from start 3: First at 1328 ms from start 3: Second at 1829 ms from start
flatMapMerge
Другой режим сглаживания — одновременный сбор всех поступающих потоков и слияние их значений в один поток, так что значения испускаются как можно скорее. Он реализован операторами flatMapMerge и flattenMerge. Оба они принимают необязательный concurrency параметр, который ограничивает количество одновременных потоков, которые собираются одновременно (по умолчанию он равен DEFAULT_CONCURRENCY).
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun requestFlow(i: Int): Flow<String> = flow {
emit("$i: First")
delay(500) // wait 500 ms
emit("$i: Second")
}
fun main() = runBlocking<Unit> {
//sampleStart
val startTime = System.currentTimeMillis() // remember the start time
(1..3).asFlow().onEach { delay(100) } // a number every 100 ms
.flatMapMerge { requestFlow(it) }
.collect { value -> // collect and print
println("$value at ${System.currentTimeMillis() - startTime} ms from start")
}
//sampleEnd
}
Конкурентная природа flatMapMerge очевидна:
1: First at 136 ms from start 2: First at 231 ms from start 3: First at 333 ms from start 1: Second at 639 ms from start 2: Second at 732 ms from start 3: Second at 833 ms from start
flatMapLatest
Аналогично оператору collectLatest, который был показан в разделе "Обработка последнего значения", существует соответствующий режим сглаживания "Последний", где коллекция предыдущего потока отменяется, как только испускается новый поток. Он реализован оператором flatMapLatest.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun requestFlow(i: Int): Flow<String> = flow {
emit("$i: First")
delay(500) // wait 500 ms
emit("$i: Second")
}
fun main() = runBlocking<Unit> {
//sampleStart
val startTime = System.currentTimeMillis() // remember the start time
(1..3).asFlow().onEach { delay(100) } // a number every 100 ms
.flatMapLatest { requestFlow(it) }
.collect { value -> // collect and print
println("$value at ${System.currentTimeMillis() - startTime} ms from start")
}
//sampleEnd
}
Вывод здесь в этом примере хорошо демонстрирует, как работает flatMapLatest:
1: First at 142 ms from start 2: First at 322 ms from start 3: First at 425 ms from start 3: Second at 931 ms from start
Исключения потока
Коллекция потока может завершиться с исключением, когда эмиттер или код внутри операторов выбрасывает исключение. Существует несколько способов обработки этих исключений.
Обработка исключений коллектором
Коллектор может использовать блок Kotlin's try/catch для обработки исключений:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
println("Emitting $i")
emit(i) // emit next value
}
}
fun main() = runBlocking<Unit> {
try {
simple().collect { value ->
println(value)
check(value <= 1) { "Collected $value" }
}
} catch (e: Throwable) {
println("Caught $e")
}
}
//sampleEnd
Этот код успешно перехватывает исключение в терминальном операторе collect и, как мы видим, больше значений не излучаются после этого:
Emitting 1 1 Emitting 2 2 Caught java.lang.IllegalStateException: Collected 2
Всё перехватывается
Предыдущий пример фактически перехватывает любое исключение, возникающее в эмиттере или в промежуточных или терминальных операторах. Например, давайте изменим код так, чтобы излучаемые значения были преобразованными в строки, но соответствующий код генерирует исключение:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<String> =
flow {
for (i in 1..3) {
println("Emitting $i")
emit(i) // emit next value
}
}
.map { value ->
check(value <= 1) { "Crashed on $value" }
"string $value"
}
fun main() = runBlocking<Unit> {
try {
simple().collect { value -> println(value) }
} catch (e: Throwable) {
println("Caught $e")
}
}
//sampleEnd
Это исключение всё ещё перехватывается, и сбор данных останавливается:
Emitting 1 string 1 Emitting 2 Caught java.lang.IllegalStateException: Crashed on 2
Прозрачность исключений
Но как код эмиттера может инкапсулировать своё поведение обработки исключений?
Потоки должны быть прозрачными к исключениям, и нарушение прозрачности исключений — это попытка излучить значения в блоке flow { ... } изнутри блока try/catch. Это гарантирует, что коллектор, выбрасывающий исключение, всегда может его перехватить с помощью блока try/catch, как в предыдущем примере.
Эмиттер может использовать оператор catch, который сохраняет эту прозрачность исключений и позволяет инкапсулировать обработку исключений. Тело оператора catch может анализировать исключение и реагировать на него различными способами в зависимости от того, какое исключение было перехвачено:
Исключения могут быть повторно сгенерированы с помощью
throw.Исключения могут быть преобразованы в излучение значений с помощью emit из тела оператора catch.
Исключения могут быть проигнорированы, записаны в журнал или обработаны каким-либо другим кодом.
Например, давайте выведем текст при перехвате исключения:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun simple(): Flow<String> =
flow {
for (i in 1..3) {
println("Emitting $i")
emit(i) // emit next value
}
}
.map { value ->
check(value <= 1) { "Crashed on $value" }
"string $value"
}
fun main() = runBlocking<Unit> {
//sampleStart
simple()
.catch { e -> emit("Caught $e") } // emit on exception
.collect { value -> println(value) }
//sampleEnd
}
Вывод примера останется прежним, даже если у нас больше нет try/catch вокруг кода.
Прозрачный перехват
Оператор catch, соблюдая прозрачность исключений, перехватывает только исключения из потоков выше catch, но не ниже него. Если блок в collect { ... } (размещённый ниже catch) генерирует исключение, то оно выходит наружу:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
println("Emitting $i")
emit(i)
}
}
fun main() = runBlocking<Unit> {
simple()
.catch { e -> println("Caught $e") } // does not catch downstream exceptions
.collect { value ->
check(value <= 1) { "Collected $value" }
println(value)
}
}
//sampleEnd
Сообщение «Перехвачено…» не выводится, несмотря на наличие оператора catch:
Emitting 1 1 Emitting 2 Exception in thread "main" java.lang.IllegalStateException: Collected 2 at ...
Декларативное перехват исключений
Мы можем объединить декларативную природу оператора catch с желанием обрабатывать все исключения, переместив тело оператора collect в onEach и поместив его перед оператором catch. Вызов данного потока должен быть инициирован вызовом collect() без параметров:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun simple(): Flow<Int> = flow {
for (i in 1..3) {
println("Emitting $i")
emit(i)
}
}
fun main() = runBlocking<Unit> {
//sampleStart
simple()
.onEach { value ->
check(value <= 1) { "Collected $value" }
println(value)
}
.catch { e -> println("Caught $e") }
.collect()
//sampleEnd
}
Теперь мы видим, что сообщение «Перехвачено…» выводится, и мы можем перехватывать все исключения без явного использования блока try/catch:
Emitting 1 1 Emitting 2 Caught java.lang.IllegalStateException: Collected 2
Завершение потока
При завершении потока (в обычном или исключительном режиме) может потребоваться выполнить действие. Как вы уже могли заметить, это можно сделать двумя способами: императивным или декларативным.
Императивный блок finally
В дополнение к try/catch, коллектор также может использовать блок finally, чтобы выполнить действие при завершении collect.
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = (1..3).asFlow()
fun main() = runBlocking<Unit> {
try {
simple().collect { value -> println(value) }
} finally {
println("Done")
}
}
//sampleEnd
Этот код выводит три числа, сгенерированные потоком simple, а затем строку "Done":
1 2 3 Done
Декларативное управление
Для декларативного подхода поток имеет промежуточный оператор onCompletion, который вызывается, когда поток полностью собран.
Предыдущий пример можно переписать, используя оператор onCompletion, что даст тот же результат:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun simple(): Flow<Int> = (1..3).asFlow()
fun main() = runBlocking<Unit> {
//sampleStart
simple()
.onCompletion { println("Done") }
.collect { value -> println(value) }
//sampleEnd
}
Ключевым преимуществом оператора onCompletion является необязательный параметр лямбда-выражения Throwable, который можно использовать для определения, было ли завершение потока нормальным или исключительным. В следующем примере поток simple генерирует исключение после вывода числа 1:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = flow {
emit(1)
throw RuntimeException()
}
fun main() = runBlocking<Unit> {
simple()
.onCompletion { cause -> if (cause != null) println("Flow completed exceptionally") }
.catch { cause -> println("Caught exception") }
.collect { value -> println(value) }
}
//sampleEnd
Как ожидалось, он выведет:
1 Flow completed exceptionally Caught exception
Оператор onCompletion, в отличие от оператора catch, не обрабатывает исключение. Как видно из приведённого выше кода, исключение всё ещё передаётся вниз по потоку. Оно будет передано последующим onCompletion операторам и может быть обработано оператором catch.
Успешное завершение
Ещё одно отличие от оператора catch заключается в том, что onCompletion видит все исключения и получает исключение null только при успешном завершении потока (без отмены или сбоя).
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun simple(): Flow<Int> = (1..3).asFlow()
fun main() = runBlocking<Unit> {
simple()
.onCompletion { cause -> println("Flow completed with $cause") }
.collect { value ->
check(value <= 1) { "Collected $value" }
println(value)
}
}
//sampleEnd
Мы видим, что причина завершения не равна null, так как поток был прерван из-за исключения в нисходящем потоке:
1 Flow completed with java.lang.IllegalStateException: Collected 2 Exception in thread "main" java.lang.IllegalStateException: Collected 2
Императивный против декларативного подхода
Теперь мы знаем, как собирать поток и обрабатывать его завершение и исключения как императивно, так и декларативно. Естественный вопрос заключается в том, какой подход предпочтительнее и почему? Как библиотека, мы не выступаем за какой-либо определённый подход и считаем, что оба варианта допустимы и должны выбираться в зависимости от ваших предпочтений и стиля кодирования.
Запуск потока
Потоки легко используются для представления асинхронных событий, поступающих из какого-либо источника. В этом случае нам нужен аналог функции addEventListener, которая регистрирует фрагмент кода с реакцией на входящие события и продолжает дальнейшую работу. Оператор onEach может выполнять эту роль. Однако, onEach является промежуточным оператором. Нам также нужен терминальный оператор для сбора потока. В противном случае, просто вызов onEach не оказывает никакого влияния.
Если мы используем терминальный оператор collect после onEach, то код после него будет ожидать, пока поток не будет собран:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
// Imitate a flow of events
fun events(): Flow<Int> = (1..3).asFlow().onEach { delay(100) }
fun main() = runBlocking<Unit> {
events()
.onEach { event -> println("Event: $event") }
.collect() // <--- Collecting the flow waits
println("Done")
}
//sampleEnd
Как видно, он выводит:
Event: 1 Event: 2 Event: 3 Done
Здесь пригодится терминальный оператор launchIn. Заменив collect на launchIn, мы можем запустить сбор потока в отдельной корутине, чтобы выполнение дальнейшего кода немедленно продолжилось:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
// Imitate a flow of events
fun events(): Flow<Int> = (1..3).asFlow().onEach { delay(100) }
//sampleStart
fun main() = runBlocking<Unit> {
events()
.onEach { event -> println("Event: $event") }
.launchIn(this) // <--- Launching the flow in a separate coroutine
println("Done")
}
//sampleEnd
Он выводит:
Done Event: 1 Event: 2 Event: 3
Требуемый параметр для launchIn должен указывать CoroutineScope, в котором запускается корутина для сбора потока. В приведенном выше примере этот scope происходит из runBlocking билдера корутин, поэтому, пока поток работает, этот runBlocking scope ожидает завершения своей дочерней корутины и предотвращает возврат и завершение этого примера.
В реальных приложениях scope будет происходить из сущности с ограниченным жизненным циклом. Как только жизненный цикл этой сущности завершен, соответствующий scope отменяется, отменяя сбор соответствующего потока. Таким образом, пара onEach { ... }.launchIn(scope) работает как addEventListener. Однако, нет необходимости в соответствующей функции removeEventListener, так как отмена и структурированная конкурентность служат этой цели.
Обратите внимание, что launchIn также возвращает Job, который может быть использован для отмены соответствующей корутины сбора потока только без отмены всего scope или для соединения ее.
Проверки отмены потока
Для удобства, билдер flow выполняет дополнительные проверки ensureActive на отмену для каждого испускаемого значения. Это означает, что цикл с испусканием из flow { ... } является отменяемым:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun foo(): Flow<Int> = flow {
for (i in 1..5) {
println("Emitting $i")
emit(i)
}
}
fun main() = runBlocking<Unit> {
foo().collect { value ->
if (value == 3) cancel()
println(value)
}
}
//sampleEnd
Мы получим только числа до 3 и CancellationException после попытки испустить число 4:
Emitting 1
1
Emitting 2
2
Emitting 3
3
Emitting 4
Exception in thread "main" kotlinx.coroutines.JobCancellationException: BlockingCoroutine was cancelled; job="coroutine#1":BlockingCoroutine{Cancelled}@6d7b4f4c
Однако, большинство других операторов потока не выполняют дополнительные проверки на отмену самостоятельно по причинам производительности. Например, если вы используете расширение IntRange.asFlow для написания того же цикла с испусканием и не приостанавливаете его нигде, то проверок на отмену нет:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun main() = runBlocking<Unit> {
(1..5).asFlow().collect { value ->
if (value == 3) cancel()
println(value)
}
}
//sampleEnd
Все числа от 1 до 5 собираются, и отмена обнаруживается только перед возвратом из runBlocking:
1
2
3
4
5
Exception in thread "main" kotlinx.coroutines.JobCancellationException: BlockingCoroutine was cancelled; job="coroutine#1":BlockingCoroutine{Cancelled}@3327bd23
Делаем занятый поток отменяемым
В случае цикла с корутинами необходимо явно проверять отмену. Можно добавить .onEach { currentCoroutineContext().ensureActive() }, но существует готовый оператор cancellable, предназначенный для этого:
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
//sampleStart
fun main() = runBlocking<Unit> {
(1..5).asFlow().cancellable().collect { value ->
if (value == 3) cancel()
println(value)
}
}
//sampleEnd
С оператором cancellable собираются только числа от 1 до 3:
1
2
3
Exception in thread "main" kotlinx.coroutines.JobCancellationException: BlockingCoroutine was cancelled; job="coroutine#1":BlockingCoroutine{Cancelled}@5ec0a365
Потоки и Reactive Streams
Для тех, кто знаком с Reactive Streams или реактивными фреймворками, такими как RxJava и Project Reactor, дизайн потока может показаться очень знакомым.
Действительно, его дизайн был вдохновлен Reactive Streams и его различными реализациями. Но основной целью потока является максимально простой дизайн, дружелюбие к Kotlin и приостановкам, а также соблюдение структурированной конкурентности. Достижение этой цели было бы невозможно без пионеров реактивности и их огромной работы. Вы можете прочитать полную историю в статье Reactive Streams и Kotlin Flows.
Несмотря на различия, концептуально, Flow является реактивным потоком, и его можно преобразовать в реактивный (совместимый со спецификацией и TCK) Publisher и наоборот. Такие преобразователи предоставляются kotlinx.coroutines из коробки и можно найти в соответствующих реактивных модулях (kotlinx-coroutines-reactive для Reactive Streams, kotlinx-coroutines-reactor для Project Reactor и kotlinx-coroutines-rx2/kotlinx-coroutines-rx3 для RxJava2/RxJava3). Модули интеграции включают преобразования из и в Flow, интеграцию с Reactor's Context и дружественные приостановкам способы работы с различными реактивными сущностями.
© 2010–2022 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/flow.html