Spec-Zone.ru › Kotlin 1.4

Содержание

  • Асинхронный поток
    • Представление нескольких значений
      • Последовательности
      • Подвешивающие функции
      • Потоки
    • Потоки — холодные
    • Основы отмены потока
    • Строители потоков
    • Промежуточные операторы потоков
      • Оператор преобразования
      • Операторы ограничения размера
    • Конечные операторы потоков
    • Потоки являются последовательными
    • Контекст потока
      • Неправильная отправка с withContext
      • Оператор flowOn
    • Буферизация
      • Конфляция
      • Обработка последнего значения
    • Компоновка нескольких потоков
      • Zip
      • Combine
    • Разворачивание потоков
      • flatMapConcat
      • flatMapMerge
      • flatMapLatest
    • Исключения потока
      • Collector try и catch
      • Всё перехватывается
    • Прозрачность исключений
      • Прозрачный перехват
      • Декларативный перехват
    • Завершение потока
      • Императивный блок finally
      • Деларативная обработка
      • Успешное завершение
    • Императивное против декларативного
    • Запуск потока
    • Проверки отмены потока
      • Делаем занятый поток отменяемым
    • Поток и реактивные потоки

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

Подвешивающая функция асинхронно возвращает единственное значение, но как мы можем вернуть несколько асинхронно вычисленных значений? Здесь на помощь приходят 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. Обычно для изменения контекста в коде с использованием Kotlin coroutines используется withContext, но код в flow { ... } билдере должен соблюдать свойство сохранения контекста и не допускается emit из другого контекста.

Попробуйте запустить следующий код:

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.

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

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

Обработчик try-catch

Обработчик может использовать блок 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

Всё перехватывается

Предыдущий пример фактически перехватывает любое исключение, возникающее в излучателе или в любых промежуточных или терминальных операторах. Например, изменим код так, чтобы излучаемые значения преобразовывались в строки с помощью map, но соответствующий код генерирует исключение:

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, соблюдая прозрачность исключений, перехватывает только исключения из потока данных, идущего от операторов выше 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

Полный код можно найти здесь.

Сообщение «Caught …» не выводится, несмотря на наличие оператора 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
}            

Полный код можно найти здесь.

Теперь мы видим, что сообщение «Caught …» выводится, и таким образом мы можем перехватить все исключения без явного использования блока 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

Полный код можно найти здесь.

Мы видим, что причина завершения не нулевая, потому что поток данных был прерван из-за исключения в последующих операторах:

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

Однако, большинство других операторов flow не выполняют дополнительных проверок на отмену по причинам производительности. Например, если вы используете расширение 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

Делаем цикл flow с задержкой отменяемым

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

Flow и 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–2020 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/reference/coroutines/flow.html

Spec-Zone.ru

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