Spec-Zone.ru › Kotlin 1.6

Асинхронный поток

Функция приостановления асинхронно возвращает единственное значение, но как мы можем вернуть несколько асинхронно вычисленных значений? Здесь на помощь приходят 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 называется flow.

  • Код внутри блока-строителя flow { ... } может приостанавливаться.

  • Функция simple больше не помечается модификатором suspend.

  • Значения выводятся из потока с помощью функции emit.

  • Значения собираются из потока с помощью функции collect.

Мы можем заменить delay на Thread.sleep в теле simple и увидеть, что главная нить в этом случае заблокирована.

Потоки — холодные потоки

Потоки — это холодные потоки, похожие на последовательности — код внутри 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.

  • Операторы для получения первого значения и для обеспечения того, что поток выводит единственное значение.

  • Сведение потока к значению с помощью reduce и fold.

Например:

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

Обратите внимание, что оператор flowOn использует тот же механизм буферизации, когда ему нужно изменить CoroutineDispatcher, но здесь мы явно запрашиваем буферизацию без изменения контекста выполнения.

Конфляция

Когда поток представляет частичные результаты операции или обновления статуса операции, может не потребоваться обрабатывать каждое значение, а только самые последние. В этом случае оператор 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 мс:

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

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

Обратите внимание, что flatMapMerge вызывает свой блок кода ({ requestFlow(it) } в этом примере) последовательно, но собирает результирующие потоки параллельно. Это эквивалентно выполнению последовательной map { requestFlow(it) }, а затем вызову flattenMerge на результате.

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

Обратите внимание, что flatMapLatest отменяет весь код в своём блоке ({ requestFlow(it) } в этом примере) при новом значении. В этом конкретном примере это не имеет значения, потому что вызов requestFlow сам по себе быстрый, не-подвешивающий и не может быть отменён. Однако это проявилось бы, если бы мы использовали подвешенные функции, такие как delay.

Исключения потока

Коллекция потока может завершиться с исключением, когда эмиттер или код внутри операторов выбрасывает исключение. Существует несколько способов обработки этих исключений.

Обработка исключений коллектором

Коллектор может использовать блок 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 и дружественные приостановкам способы работы с различными реактивными сущностями.

Последнее изменение: 04 апреля 2022
Контекст и диспетчеры корутин Каналы

© 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

Spec-Zone.ru

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