Spec-Zone.ru › Kotlin 2

Потоки

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

С помощью потоков можно создавать конвейеры потоков, которые постепенно загружают данные, реагируют на потоки событий и моделируют API в стиле подписки.

Конвейер потока — это последовательность операций, в которой участвуют следующие роли:

  • Производитель: создает значения.

  • Промежуточные операторы (необязательно): получают значения из потока, применяют к ним операцию и возвращают другой поток.

  • Сборщик: получает значения из потока.

Вот простой пример, показывающий, как взаимодействуют эти роли в конвейере:

import kotlinx.coroutines.flow.*

//sampleStart
suspend fun main() {
    // The emitter produces values
    flowOf(0x4B, 0x6F, 0x74, 0x6C, 0x69, 0x6E)
        // The intermediate operator consumes values,
        // applies an operation, and returns another flow
        .map { value -> value.toChar() }
        // The collector consumes the transformed values
        .collect { updatedValue ->
            println("Say '$updatedValue'!")
        }
}
//sampleEnd

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

Parts of a flow: emitter, intermediate operator (optional), collector. Values move from upstream to downstream.

В Kotlin доступны следующие типы потоков:

  • Холодные потоки начинают создавать значения при сборе. Каждый сборщик запускает новое независимое выполнение потока.

  • Горячие потоки создают значения независимо от сборщиков и передают всем сборщикам один и тот же поток значений.

Для тестирования потоков Kotlin можно использовать библиотеку Turbine. Она упрощает сбор и проверку значений потока в модульных тестах, в том числе проверку завершения и ошибок.

Холодные потоки

Как и последовательности, холодные потоки являются ленивыми.

Код блока построителя холодного потока не выполняется, пока сборщик не начнет сбор. Каждый новый сборщик запускает новое выполнение потока.

Создание холодного потока

Чтобы создать холодный поток, используйте функцию-построитель flow(). Внутри ее блока используйте функцию emit(), чтобы передавать значения сборщикам:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
fun main() {
    // Creates a flow
    val pageFlow = flow {
        for (page in 1..3) {
            println("Loading page $page...")

            // Emits each page as it is loaded
            emit("Page $page")
        }
    }
    println("Creating a cold flow doesn't run it!")
}
//sampleEnd

В этом примере функция-построитель flow() возвращает Flow<T>, но не начинает выполнять свой блок. Холодный поток похож на рецепт: он определяет, как создавать значения, но начинает создавать их только тогда, когда вы начинаете его сбор.

Холодные потоки также можно создавать с помощью следующих функций:

  • flowOf(): создает поток из переданных значений.

  • .asFlow(): преобразует существующую итерируемую последовательность, например диапазон, в поток.

Вот пример:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

fun main() {
    // Creates a flow from provided values
    val predefinedPageFlow = flowOf("Page 1", "Page 2", "Page 3")
    // Creates a flow from a range
    val generatedPageFlow = (1..3).asFlow()
}

Сбор холодного потока

Чтобы собрать холодный поток, используйте функцию collect(), которая запускает создание значений в потоке с начала. Если передать лямбда-выражение в collect(), оно будет получать каждое созданное значение:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        val pageFlow = flow {
            for (page in 1..3) {
                println("Loading page $page...")
                emit("Page $page")
            }
        }
        // Collects the flow with a lambda that receives each emitted page
        pageFlow.collect { page ->
            println("Processing $page...")
            delay(100.milliseconds)
            println("Done processing $page.")
        }
    }
}
//sampleEnd

Каждый вызов collect() запускает весь холодный поток с самого начала. Если несколько сборщиков собирают один и тот же холодный поток, каждый сборщик запускает собственный сбор:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
suspend fun main() {
    val pageFlow = flow {
        // Reads the name of the current coroutine
        val coroutineName = currentCoroutineContext()[CoroutineName]?.name

        println("Starting emissions in $coroutineName")
        for (page in 1..3) {
            println("Loading page $page in $coroutineName")
            emit("Page $page")
        }
        println("Done emitting in $coroutineName")
    }

    withContext(Dispatchers.Default) {
        // Launches a collector that processes each page slowly
        launch(CoroutineName("a slow coroutine")) {
            pageFlow.collect {
                println("Processing $it slowly")
                delay(100.milliseconds)
                println("Done processing $it slowly")
            }
        }

        // Launches a collector that processes each page quickly
        launch(CoroutineName("a fast coroutine")) {
            pageFlow.collect {
                println("Processing $it quickly")
                delay(10.milliseconds)
                println("Done processing $it quickly")
            }
        }
    }
}
//sampleEnd

В этом примере CoroutineName задает имя каждой корутине. Для отладки можно использовать CoroutineName. Здесь это помогает показать, какой сборщик выполняет каждый сбор.

Промежуточные операторы потоков

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

Библиотека kotlinx.coroutines предлагает широкий набор промежуточных операторов потоков для преобразования и обработки потоков. Если встроенные операторы не обеспечивают нужное поведение, можно также определить собственные операторы.

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

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
// A simplified custom implementation of the default .map() operator
fun <T, R> Flow<T>.myMap(transform: suspend (value: T) -> R): Flow<R> = flow {
    // Collects values from the upstream flow
    this@myMap.collect { value ->
        // Transforms each collected value and emits the result
        emit(transform(value))
    }
}

suspend fun main() {
    // Creates a flow, applies the custom map operator, and collects the transformed values
    flowOf(1, 2, 3).myMap { 2 * it }.collect {
        println("Collecting $it")
    }
}
//sampleEnd

Вызов приостанавливаемых функций внутри построителя потока

В отличие от последовательностей, внутри функции-построителя flow() можно вызывать приостанавливаемые функции:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
suspend fun loadPage(): Int {
    delay(100)
    return 3
}

suspend fun main() {
    flow {
        emit(loadPage())
    }.collect {
        println(it)
        // 3
    }
}
//sampleEnd

Однако функция-построитель flow() должна создавать значения в том же контексте корутины, в котором она выполняется. Нельзя запускать другую корутину, которая вызывает emit() в своем блоке, а также нельзя менять контекст корутины с помощью withContext():

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
suspend fun main() {
    // This fails with an exception!
    flow {
        // Changes the coroutine context with withContext()
        withContext(Dispatchers.IO) {
            emit('a')
        }
    }.collect { 
        println("This never prints")
    }
}
//sampleEnd

Это ограничение относится к функции-построителю flow().

Если поток с начала должен выполняться в другом контексте корутины, можно изменить контекст с помощью оператора .flowOn().

В качестве альтернативы можно использовать channelFlow(), чтобы создавать значения из нескольких корутин.

Изменение контекста корутины холодного потока с помощью .flowOn()

По умолчанию холодный поток выполняется в том же контексте корутины, что и сборщик.

Чтобы поток выполнялся в другом контексте корутины, используйте оператор .flowOn(). Этот оператор сохраняет контекст. Он меняет только контекст корутины потока с начала, оставляя поток, направленный к концу, в контексте вызывающего кода.

Вот пример холодного потока, который создает значения в одном контексте корутины, а собирает их в другом:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default + CoroutineName("downstream")) {
        flow {
            val coroutineName = currentCoroutineContext()[CoroutineName]?.name

            // Emits in the coroutine context applied with .flowOn()
            println("Emitting '1' in $coroutineName")
            // Emitting '1' in upstream
            emit(1)

        // Changes the coroutine context of the upstream flow
        }.flowOn(Dispatchers.IO + CoroutineName("upstream"))
            .collect {
            val coroutineName = currentCoroutineContext()[CoroutineName]?.name

            // Collects in the caller's coroutine context
            println("Collecting '$it' in $coroutineName")
            // Collecting '1' in downstream
        }
    }
}
//sampleEnd

Обработка исключений в потоках

И производители, и сборщики могут выбрасывать исключения.

Если не обработать исключение во время сбора потока, оно распространяется от сборщика вверх по потоку и выбрасывается вызывающему коду функции collect().

Такие исключения можно обработать, заключив функцию collect() в блок try-catch, например:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

class MyFlowException(message: String) : Exception(message)

//sampleStart
suspend fun main() {
    val myFlow = flow {
        try {
            // The emit() function calls the lambda passed to collect()
            emit('a')
        } catch (e: MyFlowException) {
            println("Collector threw $e")

            // Rethrows the downstream exception
            throw e
        }
    }
    // Wraps flow collection in try-catch
    try {
        myFlow.collect {
            // Throws an exception from the collect() lambda
            throw MyFlowException("Can't process '$it'!")
        }
    } catch (e: MyFlowException) {
        println("Flow collection failed with $e")
        // Rethrows the exception to the caller
        throw e
    }
}
//sampleEnd

В этом примере сборщик выбрасывает исключение при получении значения из функции emit(). Функция-построитель flow() перехватывает это исключение, возникшее ниже по потоку.

Если внутри функции-построителя потока вы перехватываете исключение, выброшенное сборщиком, выбросьте его повторно. Это сохраняет прозрачность исключений и позволяет вызывающему коду collect() обработать исключение.

Обработка исключений в потоке с начала с помощью оператора .catch()

Чтобы обрабатывать исключения до того, как они достигнут сборщика, используйте оператор .catch().

Оператор .catch() можно использовать для обработки исключений из потока с начала, например, применяя функцию emit() для передачи резервного значения дальше по потоку:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

//sampleStart
suspend fun main() {
    flow {
        emit("a")
        emit("b")

        // Throws an exception from the upstream flow
        throw UnsupportedOperationException(
            "I am tired of listing letters"
        )
    }.catch { upstreamException ->
        println("Upstream completed with $upstreamException!")

        // Emits a fallback value downstream
        emit("Upstream terminated with an exception!")
    }.collect {
        println("Got '$it'")
    }
}
//sampleEnd

В этом примере поток с начала создает значения, а затем выбрасывает исключение. Оператор .catch() обрабатывает исключение и создает "Upstream terminated with an exception!" в качестве резервного значения.

Если в штатном режиме работы потока ожидается возникновение некоторых исключений, обрабатывайте восстанавливаемые исключения в .catch(), а неожиданные исключения выбрасывайте повторно.

Вот пример потока, который загружает данные и сообщает о ходе загрузки:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
sealed interface LoadingState {
    sealed interface Terminal: LoadingState
    object Started: LoadingState
    data class Percentage(val percents: Int): LoadingState
    object Failed: Terminal
    object Done: Terminal
}

fun loadBlob(url: String) = flow {
    emit(LoadingState.Started)

    val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)

    repeat(10) { step ->
        if (Random.nextDouble() < failureChancePerStep)
            throw IOException("Failed to load!")
        emit(LoadingState.Percentage((step + 1) * 10))
        delay(10.milliseconds)
    }
    emit(LoadingState.Done)
}.catch { e ->
    println("Loading data failed with $e")
    if (e is IOException) {
        // Handles an expected exception
        emit(LoadingState.Failed)
    } else {
        // Rethrows unexpected exceptions, so the collect() fails with them
        throw e
    }
}

suspend fun main() {
    loadBlob("https://example.com/").collect {
        println("Got '$it'")
    }
}
//sampleEnd

В этом примере, если загрузка завершается ожидаемым исключением, оператор .catch() использует функцию emit() для создания резервного состояния. Неожиданные исключения повторно выбрасываются в операторе .catch(). Это позволяет вызывающему коду функции collect() получить исключения, которые поток не обрабатывает.

Оператор .catch() не обрабатывает исключения, выброшенные сборщиком. Если лямбда-выражение, переданное в collect(), выбрасывает исключение, обработайте его с помощью блока try-catch вокруг функции collect():

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

suspend fun main() {
    val myFlow = flow {
        for (char in listOf('a', 'o', '5', 'c')) {
            try {
                emit(char)
            } catch (e: IllegalArgumentException) {
                println("Collector doesn't support character '$char': $e")

                // Rethrows the downstream exception
                throw e
            }
        }
    }.catch { e ->
        // Doesn't run because the exception happens downstream
        println("Upstream threw an exception: $e")
    }

    try {
        myFlow.collect {
            require(!it.isDigit()) { "Digits are not allowed!" }
        }
    } catch (e: IllegalArgumentException) {
        // Handles the exception from the collect() lambda
        println("Flow collection failed with $e")
    }
}

Поскольку лямбда-выражение collect() выполняется после .catch(), с помощью .catch() нельзя обрабатывать выброшенные им исключения. Чтобы обрабатывать исключения из кода, который выполняется для каждого созданного значения с помощью .catch(), поместите этот код в .onEach() перед .catch().

Оператор .onEach() выполняет свою лямбда-функцию перед передачей каждого значения дальше по потоку. Если .catch() обрабатывает исключение из .onEach(), поток завершается и следующее значение не создается:

Вот пример:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

suspend fun main() {
    flowOf('a', 'o', '5', 'c')
        // Runs before each value is emitted downstream
        .onEach {
            require(!it.isDigit()) { "Digits are not allowed!" }
            println("Got '$it'")
        }
        .catch { e ->
            println("Caught an exception: $e")
        }
        .collect()
}

В этом примере оператор .onEach() выполняется перед .catch(), поэтому оператор .catch() обрабатывает исключение, когда проверка require() завершается неудачно для '5'.

Перезапуск потока с начала после исключения

Некоторые операции могут временно завершаться неудачно, например сетевой запрос может потерять соединение. В таких случаях можно использовать оператор .retry(), чтобы перезапустить поток с начала после исключения.

Оператор .retry() получает исключение и перезапускает сбор, если его лямбда-выражение возвращает true, но не больше указанного числа повторных попыток. Например, .retry(3) повторяет сбор потока с начала до трех раз после первой неудачной попытки.

Если лямбда-выражение возвращает false, .retry() прекращает повторные попытки и повторно выбрасывает исключение.

Для более гибкого управления логикой повторных попыток используйте оператор .retryWhen(). Как и .retry(), он получает исключение, но также получает номер текущей попытки и может создавать значения перед повторной попыткой.

Вот пример, в котором загрузка повторяется до трех раз после возникновения IOException:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds

//sampleStart
sealed interface LoadingState {
    sealed interface Terminal: LoadingState
    object Started: LoadingState
    data class Percentage(val percents: Int): LoadingState
    object Failed: Terminal
    object Done: Terminal
}

fun loadBlob(url: String) = flow {
    emit(LoadingState.Started)

    val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)

    repeat(10) { step ->
        if (Random.nextDouble() < failureChancePerStep)
            throw IOException("Failed to load!")
        emit(LoadingState.Percentage((step + 1) * 10))
        delay(10.milliseconds)
    }
    emit(LoadingState.Done)
}.retry(3) { e ->
    if (e is IOException) {
        // This is an expected error
        // Waits for one second before retrying
        delay(1.seconds)
        true
    } else {
        // Stops retrying and rethrows unexpected exceptions
        false
    }
}

suspend fun main() {
    loadBlob("https://example.org/").collect {
        println("Got $it")
    }
}
//sampleEnd

Отмена потока

Отмена потока прекращает сбор, когда результат больше не нужен, например, если истекло время ожидания запроса.

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

Чтобы отменить сбор потока, вызовите функцию cancel() для Job собирающей корутины:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
val myFlow = flow {
    var i = 0
    try {
        while (true) {
            println("Emitting $i")
            emit(i)
            println("Emitted $i")
            ++i
            delay(10.milliseconds)
        }
    } catch (e: Throwable) {
        println("Upstream finished with $e")
        throw e
    }
}

suspend fun main() {
    coroutineScope {
        val job = launch {
            try {
                myFlow.collect {
                    println("Processing $it")
                    delay(5.milliseconds)
                }
            } catch (e: Throwable) {
                println("Collection finished with $e")
                throw e
            }
        }
        delay(100.milliseconds)

        // Cancels the coroutine that collects the flow
        job.cancel()
    }
}
//sampleEnd

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

Оператор .take() использует это поведение, чтобы прекращать сбор после фиксированного числа значений. Например, .take(3) собирает только первые три значения из потока с начала, а затем отменяет его.

Вот пример упрощенной версии оператора .take():

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
// Defines a simplified version of the default .take() operator
fun <T> Flow<T>.myTake(count: Int): Flow<T> = flow {
    require(count > 0)
    val cancellationException = CancellationException()
    var elementsRemaining = count
    try {
        this@myTake.collect {
            emit(it)
            --elementsRemaining
            if (elementsRemaining == 0) {
                // Cancels the upstream flow after the requested number of values
                throw cancellationException
            }
        }
    } catch (e: Throwable) {
        if (e === cancellationException) {
            // Handles the CancellationException used to cancel the upstream flow
            // Completes the flow after the set number of values in .myTake()
        } else {
            // Rethrows unexpected exceptions
            throw e
        }
    }
}

suspend fun main() {
    (0..1000).asFlow().myTake(3).collect {
        println("Got $it")
    }
}
//sampleEnd

В этом примере функция .myTake() создает значения из потока с начала, пока не будут созданы все запрошенные значения. Затем она выбрасывает исключение CancellationException, чтобы отменить поток с начала.

Создание значений одновременно с помощью channelFlow()

Функция-построитель flow() проста и эффективна для потоков, создающих значения из одной корутины. Если нужно одновременно создавать значения для одного потока из нескольких корутин, используйте функцию-построитель channelFlow(). Ее можно применять для параллельной работы, которая постепенно сообщает результаты, например для загрузки данных из нескольких источников.

Функция-построитель channelFlow() создает холодный поток, использующий канал для передачи значений из нескольких корутин. Внутри построителя используйте функцию send() вместо функции emit() для создания значений.

Вот пример, в котором channelFlow() одновременно собирает два потока и повторно передает их значения, используя упрощенную версию оператора .merge():

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
// Defines a simplified version of the default .merge() operator
fun <T> Flow<T>.myMerge(other: Flow<T>): Flow<T> = channelFlow {
    // CoroutineScope and SendChannel are available as receivers here
    // Launches a coroutine that collects the receiver flow
    launch {
        // Collects the receiver flow
        this@myMerge.collect {
            send(it)
        }
    }
    launch {
        // Launches a coroutine that collects the other flow
        other.collect {
            // Calls SendChannel.send
            send(it)
        }
    }
}

suspend fun main() {
    val flow1 = (0..3).asFlow().onEach { delay(20.milliseconds) }
    val flow2 = (6..9).asFlow().onEach { delay(50.milliseconds) }
    flow1.myMerge(flow2).collect { println(it) }
}
//sampleEnd

Функция-построитель channelFlow() использует буферизованный канал, позволяя производителям отправлять значения с опережением сборщика, пока буфер не заполнится. По умолчанию буфер может содержать до 64 значений. Когда буфер заполнен, производители приостанавливаются, пока в буфере не освободится место.

Размер буфера можно изменить с помощью оператора .buffer(). Например, .buffer(12) позволяет производителям отправлять до 12 значений с опережением сборщика, а .buffer(0) удаляет буфер, поэтому каждое значение отправляется только тогда, когда сборщик может его принять.

Вот пример:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
suspend fun main() {
    val oneHundredNumbers = channelFlow {
        repeat(100) {
            println("Sending $it")
            send(it)
        }
    }

    // Uses the default buffer capacity
    oneHundredNumbers.collect {
        println("Processing $it")
        delay(10.milliseconds)
    }
  
    // Removes the buffer so sending and processing interleave from the start
    oneHundredNumbers.buffer(0).collect {
        println("Processing $it")
        delay(10.milliseconds)
    }
}
//sampleEnd

В этом примере для потока oneHundredNumbers используется размер буфера по умолчанию, а поток oneHundredNumbers.buffer(0) не имеет буфера.

При размере буфера по умолчанию производитель быстро отправляет значения, пока буфер не заполнится. После этого send() приостанавливается, пока в буфере не освободится место, поэтому сообщения Sending и Processing начинают чередоваться.

При использовании .buffer(0) каждый вызов send() ожидает, пока сборщик сможет принять значение, поэтому Sending и Processing чередуются с самого начала.

Горячие потоки

Горячие потоки — это общие потоки, которые выдают значения независимо от сборщиков. Они продолжают выдавать значения, даже когда ни один сборщик не активен, а несколько сборщиков могут собирать одни и те же данные из уже активного потока, вместо того чтобы запускать новое выполнение.

Сборщик горячего потока называется подписчиком.

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

В Kotlin предусмотрено два типа горячих потоков:

  • SharedFlow передаёт значения нескольким подписчикам. Используйте его, когда нужно передавать события, происходящие с течением времени, например сообщения или уведомления.

  • StateFlow — это специализированный SharedFlow, который всегда хранит последнее значение состояния. Используйте его, когда нужно представить состояние, изменяющееся с течением времени, например состояние пользовательского интерфейса.

Создание SharedFlow

SharedFlow — это горячий поток, который передаёт подписчикам выданные значения, появляющиеся с течением времени.

Создать SharedFlow можно с помощью функции MutableSharedFlow().

MutableSharedFlow предоставляет функции для выдачи значений. Если предоставить его напрямую, код вне класса сможет выдавать значения в поток.

Чтобы этого избежать, сохраните изменяемый поток в приватном резервном свойстве, а наружу предоставьте доступный только для чтения SharedFlow с помощью функции .asSharedFlow(). Чтобы выдавать значения подписчикам, используйте функцию emit() для MutableSharedFlow:

data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

class Chatroom {
    // Stores the SharedFlow in a private backing property
    private val _messages = MutableSharedFlow<Message>()

    // Exposes a read-only SharedFlow to subscribers
    val messages: SharedFlow<Message>
        get() = _messages.asSharedFlow()

    suspend fun sendMessageToEveryone(message: Message) {
        // Emits the message to subscribers
        _messages.emit(message)
    }
}

Как и для холодных потоков, для сбора значений из SharedFlow можно использовать функцию collect().

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

Чтобы задать число предыдущих выдач, которые получает новый подписчик, используйте параметр replay в MutableSharedFlow():

// Sets the number of already emitted messages new subscribers receive on subscription
const val MESSAGES_TO_REMEMBER = 10

class Chatroom {
    private val _messages = MutableSharedFlow<Message>(

        // Replays the set amount of last emitted messages to new subscribers
        replay = MESSAGES_TO_REMEMBER
    )

    val messages: SharedFlow<Message>
        get() = _messages.asSharedFlow()

    suspend fun sendMessageToEveryone(message: Message) {
        // Emits the message to subscribers of the messages flow 
        _messages.emit(message)
    }
}

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

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

Рассмотрим пример, в котором SharedFlow используется для моделирования комнаты чата: она отправляет каждое новое сообщение активным подписчикам и повторно передаёт недавние сообщения подписчикам, присоединившимся позднее:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.*
import java.io.IOException
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.*

data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

// Sets the number of already emitted messages that new subscribers receive on subscription
const val MESSAGES_TO_REMEMBER = 10

class Chatroom {
    // Stores the SharedFlow in a private backing property
    private val _messages = MutableSharedFlow<Message>(

        // Replays the set amount of last emitted messages to new subscribers
        replay = MESSAGES_TO_REMEMBER
    )

    // Exposes a read-only SharedFlow to subscribers
    val messages: SharedFlow<Message>
        get() = _messages.asSharedFlow()


    // Emits the message to subscribers
    suspend fun sendMessageToEveryone(message: Message) {
        _messages.emit(message)
    }
}

suspend fun main() {
    val nUsers = 3
    val chatroom = Chatroom()
    withContext(Dispatchers.Default) {
        // Starts a message reader for each user
        val messageReaders = List(nUsers) { userId ->
            // Starts collection before messages are emitted
            launch(start = CoroutineStart.UNDISPATCHED) {
                chatroom.messages.collect { message ->
                    println("User $userId received $message")
                }
            }
        }
        // Sends a greeting from each user
        repeat(nUsers) { userId ->
            chatroom.sendMessageToEveryone(
                Message(
                    userId,
                    Clock.System.now(),
                    "Hello from $userId!"
                )
            )
        }
        // Delays to make sure people have enough time to chat
        delay(100.milliseconds)
        // Cancels readers because SharedFlow collection doesn't finish by itself
        messageReaders.forEach { it.cancel() }
    }

}

В этом примере CoroutineStart.UNDISPATCHED немедленно запускает каждую собирающую корутину.

Это гарантирует, что каждая корутина дойдёт до collect(), подпишется на messages и приостановится до того, как sendMessageToEveryone() начнёт выдавать сообщения. Без этого собирающая корутина может запуститься позже и пропустить ранние выдачи, если кэш повторного воспроизведения слишком мал.

Использование явных резервных полей для предоставления горячих потоков

Используйте явные резервные поля, чтобы предоставить доступный только для чтения SharedFlow, сохраняя изменяемое резервное поле внутри класса.

Явные резервные поля задают тип реализации в объявлении field. Внутри класса компилятор выполняет интеллектуальное приведение свойства к типу резервного поля, поэтому можно вызывать функцию emit() без отдельного приватного резервного свойства.

Явные резервные поля не создают обёртку только для чтения, которую предоставляет .asSharedFlow(). Используйте этот шаблон, только если вас не беспокоит приведение предоставленного потока к более конкретному типу.

Пример:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Clock
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.ExperimentalTime
import kotlin.time.Instant
        
data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

const val MESSAGES_TO_REMEMBER = 10

//sampleStart
class Chatroom {
    // Exposes a read-only SharedFlow with a mutable backing field
    val messages: SharedFlow<Message>
        field = MutableSharedFlow<Message>(
            replay = MESSAGES_TO_REMEMBER
        )

    suspend fun sendMessageToEveryone(message: Message) {
        // Emits through the mutable backing field inside Chatroom
        messages.emit(message)
    }
}
//sampleEnd

suspend fun main() {
    val nUsers = 3
    val chatroom = Chatroom()

    withContext(Dispatchers.Default) {
        val messageReaders = List(nUsers) { userId ->
            launch(start = CoroutineStart.UNDISPATCHED) {
                chatroom.messages.collect { message ->
                    println("User $userId received $message")
                }
            }
        }

        repeat(nUsers) { userId ->
            chatroom.sendMessageToEveryone(
                Message(
                    senderId = userId,
                    time = Clock.System.now(),
                    text = "Hello from $userId!"
                )
            )
        }

        delay(100.milliseconds)
        messageReaders.forEach { it.cancel() }
    }
}

Создание StateFlow

StateFlow — это горячий поток, который хранит одно значение состояния и выдаёт обновления, когда оно заменяется новым. Новые подписчики получают текущее значение сразу после начала сбора, а затем получают новые значения при каждом обновлении состояния.

Используйте StateFlow для представления состояния, изменяющегося с течением времени, например хода загрузки, состояния пользовательского интерфейса или состояния объекта.

Чтобы создать StateFlow, используйте функцию MutableStateFlow() с начальным значением:

// Creates a MutableStateFlow with LoadingState.Started as the initial value
val result = MutableStateFlow<LoadingState>(LoadingState.Started)

Чтобы задать текущее состояние, используйте свойство value:

fun loadBlob(url: String): StateFlow<LoadingState> {
    val result = MutableStateFlow<LoadingState>(LoadingState.Started)

    DownloadManager.startLoading(
        url,
        onPercentageLoaded = { percentage ->
            // Replaces the current state with the latest progress
            result.value = LoadingState.Percentage(percentage)
        },
        onCompletion = {
            // Replaces the current state with the completion state
            result.value = LoadingState.Done
        },
        onFailure = {
            // Replaces the current state with the failure state
            result.value = LoadingState.Failed
        }
    )
}

Присваивание value является потокобезопасным и заменяет текущее состояние, однако обновление value на основе его предыдущего значения не является атомарным. Если новое состояние зависит от предыдущего, используйте вместо этого .update().

Как и MutableSharedFlow, MutableStateFlow предоставляет API для выдачи обновлений. Если предоставить его напрямую, любой получивший его код сможет обновить состояние, приведя его к MutableStateFlow.

Чтобы этого избежать, предоставьте изменяемый поток в виде доступного только для чтения StateFlow с помощью функции .asStateFlow():

fun loadBlob(url: String): StateFlow<LoadingState> {
    val result = MutableStateFlow<LoadingState>(LoadingState.Started)

    DownloadManager.startLoading(
        url,
        onPercentageLoaded = { percentage ->
            result.value = LoadingState.Percentage(percentage)
        },
        onCompletion = {
            result.value = LoadingState.Done
        },
        onFailure = {
            result.value = LoadingState.Failed
        }
    )

    // Exposes the loading state as a read-only StateFlow
    return result.asStateFlow()
}

Вот пример использования StateFlow для передачи хода загрузки из API на основе обратных вызовов:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random

//sampleStart
sealed interface LoadingState {
    sealed interface Terminal: LoadingState
    object Started: LoadingState
    data class Percentage(val percents: Int): LoadingState
    object Failed: Terminal
    object Done: Terminal
}

fun loadBlob(url: String): StateFlow<LoadingState> {
    // Creates a mutable StateFlow with the initial loading state
    val result = MutableStateFlow<LoadingState>(LoadingState.Started)
    DownloadManager.startLoading(
        url,
        onPercentageLoaded = { percentage ->
            // Replaces the current state with the latest progress
            result.value = LoadingState.Percentage(percentage)
        },
        onCompletion = {
            // Replaces the current state with the completion state
            result.value = LoadingState.Done
        },
        onFailure = {
            // Replaces the current state with the failure state
            result.value = LoadingState.Failed
        }
    )
    // Exposes the loading state as a read-only StateFlow
    return result.asStateFlow()
}

// Defines a callback-based API that downloads data asynchronously
object DownloadManager {
    // Starts loading the url asynchronously
    fun startLoading(
        url: String,
        onPercentageLoaded: (Int) -> Unit,
        onCompletion: () -> Unit,
        onFailure: (Throwable) -> Unit
    ) {
        // Uses GlobalScope for illustrative purposes only,
        // to keep this example self-contained
        GlobalScope.launch {
            val failureChancePerStep = 1 - java.lang.Math.pow(0.99, 10.0)

            repeat(10) { step ->
                if (Random.nextDouble() < failureChancePerStep) {
                    onFailure(IOException("Failed to load!"))
                    return@launch
                }
                onPercentageLoaded((step + 1) * 10)
                delay(10.milliseconds)
            }
            onCompletion()
        }
    }
}

suspend fun main() {
    loadBlob("https://example.com/").onEach { state ->
        when (state) {
            is LoadingState.Started -> {
                // Waits for progress updates
            }
            is LoadingState.Percentage ->
                println("Loaded ${state.percents}...")
            is LoadingState.Failed ->
                println("Loading failed.")
            is LoadingState.Done ->
                println("Finished loading!")
        }
    }.takeWhile { it !is LoadingState.Terminal }.collect()
}
//sampleEnd

В этом примере GlobalScope используется только для краткости API на основе обратных вызовов. В собственных приложениях передавайте CoroutineScope функции, запускающей работу, например startLoading() в этом примере, и запускайте корутину в этой области видимости, чтобы вызывающий код мог отменить работу, когда она больше не нужна.

Поскольку StateFlow — это горячий поток, сбор не завершается самостоятельно. В этом примере оператор .takeWhile() останавливает сбор, когда загрузка достигает конечного состояния.

StateFlow выдаёт обновление только в том случае, если новое значение отличается от текущего.

Не храните изменяемые объекты в StateFlow. Изменение самого объекта не заменяет текущее значение, поэтому сборщики не получают обновление.

Также можно обновить StateFlow, вычислив новое состояние на основе текущего. Для таких обновлений используйте функцию .update(). Функция .update() обновляет значение атомарно, что полезно, когда несколько корутин обновляют один и тот же MutableStateFlow.

Если нужно только обновлять общее значение и не требуется наблюдать за изменениями состояния с течением времени, используйте API атомарных переменных Kotlin, например AtomicInt или AtomicReference.

Вот пример, в котором число отметок «Нравится» хранится в StateFlow, а каждое новое состояние вычисляется на основе предыдущего:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random

//sampleStart
class Post(val id: Long) {
    // Stores the current number of likes as a StateFlow
    private val _numberOfLikes = MutableStateFlow<Int>(
        // Sets the initial number of likes
        0
    )

    // Exposes a read-only StateFlow with the current number of likes
    val numberOfLikes: StateFlow<Int>
        get() = _numberOfLikes.asStateFlow()


    // Adds a like
    fun like() {
        // Increments the number of likes atomically for concurrent and multithreaded calls
        _numberOfLikes.update { it + 1 }
    }
}

suspend fun drawUpdatedNumberOfLikes(likes: Int) {
    // Displays the latest number of likes
    println("${Clock.System.now()}: the number of likes is $likes")
}

suspend fun main() {
    withContext(Dispatchers.Default) {
        val post = Post(15)
        val notifyingJob = launch {
            post.numberOfLikes.collect {
                drawUpdatedNumberOfLikes(it)
            }
        }
        // Simulates users who like the post
        coroutineScope {
            repeat(10) {
                launch {
                    delay(Random.nextInt(100).milliseconds)
                    post.like()
                }
            }
        }
        // Cancels collection after all simulated users finish
        notifyingJob.cancelAndJoin()
    }
}
//sampleEnd

В этом примере функция .update() атомарно увеличивает число отметок «Нравится». Это предотвращает потерю обновлений, когда несколько корутин одновременно вызывают функцию like().

Хранение накопленного состояния в StateFlow

Иногда подписчикам нужно получать результат всех предыдущих выдач, а не только последнее выданное значение.

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

Такое поведение можно смоделировать с помощью StateFlow.

Для этого храните всю историю сообщений как текущее значение с помощью StateFlow<List<Message>> вместо того, чтобы передавать каждое сообщение чата как отдельное событие с помощью SharedFlow<Message>:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random

//sampleStart
data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

class Chatroom {
    // Stores the full message history
    private val _messageHistory = MutableStateFlow<List<Message>>(emptyList())

    // Exposes a read-only StateFlow with the current message history
    val messageHistory: StateFlow<List<Message>>
        get() = _messageHistory.asStateFlow()

    // Sends a message to all subscribers of the messageHistory flow
    suspend fun sendMessageToEveryone(message: Message) {
        // Adds the new message to the current history atomically
        _messageHistory.update {
            it + message
        }
    }
}

suspend fun main() {
    val nUsers = 3
    val chatroom = Chatroom()
    withContext(Dispatchers.Default) {
        // Starts a message reader for each user
        val messageReaders = List(nUsers) { userId ->
            launch(start = CoroutineStart.UNDISPATCHED) {
                chatroom.messageHistory.collect { currentHistory ->
                    println("User $userId sees the history as $currentHistory")
                }
            }
        }
        // Sends a greeting from each user
        repeat(nUsers) { userId ->
            chatroom.sendMessageToEveryone(
                Message(
                    userId,
                    Clock.System.now(),
                    "Hello from $userId!"
                )
            )
        }
        // Delays to make sure users have enough time to receive updates
        delay(100.milliseconds)
        // Cancels readers because StateFlow collection doesn't finish by itself
        messageReaders.forEach { it.cancel() }
    }
}
//sampleEnd

В этом примере messageHistory хранит полный список предыдущих сообщений в качестве текущего состояния. Когда отправляется новое сообщение, функция .update() атомарно создаёт новый список на основе предыдущей истории и добавляет в него новое сообщение.

Обновление неизменяемых коллекций путём создания новых коллекций может занимать больше времени по мере роста коллекции. Для повышения эффективности обновления неизменяемых коллекций можно создавать персистентные коллекции с помощью библиотеки Experimental kotlinx.collections.immutable.

Поскольку messageHistory является StateFlow, подписчики получают текущую историю сообщений, когда начинают сбор. Затем они получают новый список при каждой отправке сообщения, изменяющей историю чата.

Преобразование холодных потоков в горячие

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

Следующая упрощённая версия .shareIn() демонстрирует эту идею: она собирает холодный поток один раз, выдаёт его значения в MutableSharedFlow и предоставляет его как доступный только для чтения SharedFlow:

import kotlinx.coroutines.flow.*
import kotlinx.coroutines.*

fun <T> Flow<T>.simpleShareIn(scope: CoroutineScope): SharedFlow<T> {
    val sharedFlow = MutableSharedFlow<T>()
    scope.launch {
        this@simpleShareIn.collect {
            sharedFlow.emit(it)
        }
    }
    return sharedFlow.asSharedFlow()
}

suspend fun main() { 
    
}

В этом примере simpleShareIn() запускает новую корутину в предоставленной области видимости. Чтобы прекратить сбор исходного потока, отмените область видимости, в которой выполняется собирающая корутина.

Если исходный поток выбрасывает исключение, эта собирающая корутина завершается с ошибкой. Используйте такие операторы, как .catch() или .retry(), перед тем как делиться потоком, чтобы обрабатывать исключения исходного потока до сбоя собирающей корутины.

Встроенная функция .shareIn() реализует этот шаблон, избавляя от необходимости самостоятельно создавать MutableSharedFlow. Она также предоставляет параметры для управления началом и остановкой сбора исходного потока, а также числом предыдущих выдач, получаемых новыми подписчиками.

Чтобы использовать встроенную функцию .shareIn(), укажите следующие аргументы:

  • Область видимости корутин, в которой собирается исходный поток.

  • Стратегию SharingStarted, управляющую началом и остановкой сбора исходного потока. Например, SharingStarted.Eagerly немедленно начинает сбор исходного потока в предоставленной области видимости, до того как сбор начнёт какой-либо подписчик.

  • Необязательное значение replay, задающее число предыдущих выдач, которые получают новые подписчики.

Функция .shareIn() собирает исходный поток в предоставленной области видимости корутин и передаёт его выдачи подписчикам.

Вот пример, в котором .shareIn() преобразует холодный поток в горячий, совместно использующий сериализованные сообщения чата несколькими подписчиками:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random

//sampleStart
data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

class Chatroom {
    // Stores the message flow
    private val _messages = MutableSharedFlow<Message>()

    // Exposes a read-only SharedFlow with emitted messages
    // New subscribers don't receive already emitted messages
    val messages: SharedFlow<Message>
        get() = _messages.asSharedFlow()
  
    // Sends a message to all subscribers of the messages flow
    suspend fun sendMessageToEveryone(message: Message) {
        _messages.emit(message)
    }
}

suspend fun main() {
    val nUsers = 3
    val chatroom = Chatroom()
    withContext(Dispatchers.Default) {
        // Creates a child scope of the currently running coroutine
        val derivedFlowsScope = CoroutineScope(
            currentCoroutineContext() + Job(currentCoroutineContext()[Job])
        )
        // Shares serialized messages between subscribers
        val serializedMessages: SharedFlow<String> =
            chatroom
                .messages
                .map {
                    // Serializes each message once for the shared flow
                    "senderId: ${it.senderId}, time: ${it.time}, text: " +
                        Base64.Default.encode(it.text.encodeToByteArray())
                }
                .shareIn(
                    // Starts the sharing coroutine in this scope.
                    // The upstream flow, including .map(), runs in that coroutine
                    derivedFlowsScope,

                    // Starts collecting the upstream flow immediately,
                    // before the first subscriber appears
                    SharingStarted.Eagerly,

                    // Doesn't replay previous serialized messages to new subscribers
                    replay = 0,
                )

        // Starts a message reader for each user
        val messageReaders = List(nUsers) { userId ->
            launch(start = CoroutineStart.UNDISPATCHED) {
                serializedMessages.collect { serializedMessage ->
                    println("User $userId observes the message $serializedMessage")
                }
            }
        }
        // Sends a greeting from each user
        repeat(nUsers) { userId ->
            chatroom.sendMessageToEveryone(
                Message(
                    userId,
                    Clock.System.now(),
                    "Hello from $userId!"
                )
            )
        }
        // Delays to make sure users have enough time to receive updates
        delay(100.milliseconds)
        // Cancels readers because SharedFlow collection doesn't finish by itself
        messageReaders.forEach { it.cancel() }
        // Cancels the scope that runs the derived hot flow
        derivedFlowsScope.cancel()
    }
}
//sampleEnd

В этом примере оператор .map() создаёт холодный поток, сериализующий каждое сообщение. Без функции .shareIn() каждый сборщик выполнял бы сериализацию отдельно. Функция .shareIn() совместно использует один сбор исходного потока, поэтому каждое сообщение сериализуется один раз, а затем передаётся всем подписчикам.

Поскольку SharingStarted.Eagerly немедленно начинает сбор исходного потока, производный горячий поток начинает собирать chatroom.messages сразу после вызова .shareIn().

Аналогично, чтобы преобразовать холодный поток в StateFlow, используйте функцию .stateIn().

В отличие от .shareIn(), для .stateIn() требуется начальное значение, поскольку StateFlow всегда должен иметь текущее значение.

Например:

val lastUpdateFlow: StateFlow<Instant?> =
    chatroom
        .messageHistory
        .map { currentHistory -> currentHistory.lastOrNull()?.time }
        .stateIn(
            // Starts the sharing coroutine in this scope
            // The upstream flow, including .map(), runs in that coroutine
            derivedFlowsScope,

            // Starts collecting when the first subscriber appears
            // and stops when the last subscriber disappears
            SharingStarted.WhileSubscribed(),
            // Sets the initial state before the first upstream emission
            null,
        )

Отмена горячих потоков

Горячие потоки не останавливаются при отмене подписчика.

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

У самих горячих потоков нет операции отмены. Чтобы отменить горячий поток, отмените корутину или область видимости, в которой для него создаются значения.

Горячие потоки, созданные с помощью функций-расширений .shareIn() или .stateIn(), продолжают собирать исходный поток до отмены корутины совместного использования. Чтобы прекратить сбор исходного потока, отмените область видимости, в которой выполняется корутина совместного использования.

Также можно автоматически останавливать сбор исходного потока при отсутствии подписчиков с помощью SharingStarted.WhileSubscribed().

Вот пример, в котором отмена области видимости, переданной в .stateIn(), останавливает сбор новых значений производным горячим потоком:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.*
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.Duration.Companion.seconds
import kotlin.io.encoding.*
import java.io.IOException
import kotlin.random.Random

data class Message(
    val senderId: Int,
    val time: Instant,
    val text: String,
)

class Chatroom {
    // Stores the message history
    private val _messageHistory = MutableStateFlow<List<Message>>(emptyList())

    // Exposes a read-only StateFlow with the current message history
    val messageHistory: StateFlow<List<Message>>
        get() = _messageHistory.asStateFlow()

    // Sends a message to all subscribers of the messageHistory flow
    suspend fun sendMessageToEveryone(message: Message) {
        _messageHistory.update {
            it + message
        }
    }
}

//sampleStart
suspend fun main() {
    val chatroom = Chatroom()
    withContext(Dispatchers.Default) {
        // Creates a child scope of the currently running coroutine
        val derivedFlowsScope = CoroutineScope(
            currentCoroutineContext() + Job(currentCoroutineContext()[Job])
        )
        val totalMessages = chatroom.messageHistory
            .map { currentHistory ->
                currentHistory.size
            }.onEach {
                println("There are currently $it messages")
            }.stateIn(
                // Starts the sharing coroutine in this scope
                derivedFlowsScope
            )
        // Updates messageHistory
        chatroom.sendMessageToEveryone(
            Message(0, Clock.System.now(), "We are shutting down soon!")
        )
        delay(100.milliseconds)
        // Cancels the scope that runs the derived hot flow
        derivedFlowsScope.cancel()
        // Updates messageHistory, but totalMessages no longer receives the update
        chatroom.sendMessageToEveryone(
            Message(0, Clock.System.now(), "We have shut down.")
        )
        println("Last collected history size: ${totalMessages.value}")
        println("Actual history size: ${chatroom.messageHistory.value.size}")
    }
}
//sampleEnd

В этом примере при вызове функции derivedFlowsScope.cancel() функция totalMessages прекращает собирать обновления из messageHistory.

Функция sendMessageToEveryone() по-прежнему обновляет messageHistory, поскольку вызывающая её корутина не была отменена. В результате totalMessages.value сохраняет размер последнего собранного значения, а chatroom.messageHistory.value.size показывает фактическое число сообщений.

Обработка исключений в горячих потоках

В холодных потоках исключения источника передаются вызывающему коду collect(), если предварительно не обработать их с помощью такого оператора, как .catch().

Горячие потоки не передают исключения от производителей подписчикам. Если код, выдающий значения в MutableSharedFlow или обновляющий MutableStateFlow, выбрасывает исключение, обработайте его в корутине, выполняющей этот код. Если подписчик выбрасывает исключение во время сбора, обработайте его в собирающей корутине.

Горячие потоки, созданные с помощью функций-расширений .shareIn() или .stateIn(), собирают исходный поток в корутине совместного использования. Если исходный поток выбрасывает исключение, оно отменяет корутину совместного использования:

import kotlinx.coroutines.flow.*
import kotlinx.coroutines.*

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        launch {
            flow<Int> {
                error("An upstream failure")
            }.stateIn(
                this@launch
            )
        }
    }
}
//sampleEnd

После сбоя можно перезапустить сбор исходного потока. Для этого разместите оператор .retry() перед .shareIn() или .stateIn():

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.time.Duration.Companion.milliseconds

//sampleStart
suspend fun main() {
    coroutineScope {
        launch {
            var currentAttempt = 0

            val stateFlow = flow {
                delay(10.milliseconds)

                if (currentAttempt++ < 5) {
                    println("An error happened!")
                    error("An upstream failure")
                } else {
                    println("Success.")
                    emit(10)
                }
            }
                // Restarts the upstream flow after recoverable failures
                .retry(retries = 5)
                .stateIn(
                    // Starts the sharing coroutine in this scope
                    this@launch
                )

            stateFlow.collect {
                println("Observed $it")

                // Cancels collection and the sharing coroutine
                this@launch.cancel()
            }
        }
    }
}
//sampleEnd

В этом примере поток пять раз завершается с ошибкой, прежде чем выдаёт значение. Поскольку .retry() выполняется перед .stateIn(), он обрабатывает каждую ошибку исходного потока до того, как она достигнет корутины совместного использования.

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

13 июля 2026 г.
Отмена и тайм-аутыОператоры потоков

© 2010–2026 JetBrains s.r.o. and Kotlin Programming Language contributors
Licensed under the Apache License, Version 2.0.
https://kotlinlang.org/docs/coroutines-flow.html

Spec-Zone.ru

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