Асинхронный поток
Функция приостановления асинхронно возвращает одно значение, но как мы можем вернуть несколько асинхронно вычисленных значений? Здесь на помощь приходят 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 не выполняется, пока поток не собран. Это становится ясно в следующем примере:
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) } // emit 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, описанному в разделе "Обработка последнего значения", существует соответствующий режим сглаживания "Latest", где сбор предыдущего потока отменяется как только выдан новый поток. Он реализован оператором 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 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, поэтому, в то время как поток работает, этот scope runBlocking ожидает завершения своего дочернего сопроцесса и не позволяет функции main возвращаться и завершать этот пример.
В реальных приложениях 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, дизайн Flow может показаться очень знакомым.
Действительно, его дизайн был вдохновлен Reactive Streams и его различными реализациями. Но основная цель Flow — иметь как можно более простой дизайн, быть 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, интеграцию с реакторным Context и приостанавливаемыми способами работы с различными реактивными сущностями.
© 2010–2023 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/flow.html