Асинхронный поток
Функция приостановки асинхронно возвращает одно значение, но как мы можем вернуть несколько асинхронно вычисленных значений? Именно здесь появляются потоки Kotlin.
Представление нескольких значений
Несколько значений можно представить в 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
}
Однако, когда здесь вместо оператора zip используется оператор combine:
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, который был описан в разделе "Обработка последнего значения", существует соответствующий режим сглаживания "Последний", где сбор предыдущего потока отменяется, как только испускается новый поток. Он реализован оператором 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
Исключения потока
Коллекция потоков может завершиться с исключением, когда эмиттер или код внутри операторов бросают исключение. Есть несколько способов обработать эти исключения.
Обработка исключений коллектором
Коллектор может использовать блок try/catch Kotlin для обработки исключений:
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, а затем строку "Готово":
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 ждет завершения своей дочерней сопроцедуры и не позволяет основной функции возвращаться и завершать этот пример.
В реальных приложениях 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.
Хотя они отличаются, концептуально, поток является реактивным потоком, и его можно преобразовать в реактивный (совместимый со спецификацией и 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–2022 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/flow.html