Spec-Zone.ru › Kotlin 2

Операторы Flow

Операторы Flow позволяют преобразовывать и обрабатывать значения в цепочке потока. В Kotlin есть два основных типа операторов Flow:

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

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

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

В следующих разделах приведены примеры собственных реализаций наряду с соответствующими встроенными операторами.

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

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

Промежуточные операторы можно разделить на следующие категории по назначению:

  • Операторы преобразования преобразуют значения перед передачей вниз по потоку.

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

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

  • Операторы объединения собирают значения из нескольких восходящих потоков и передают их в один нисходящий поток.

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

Операторы преобразования

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

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

Оператор .transform() — это универсальный оператор преобразования, на основе которого можно создавать более специализированные операторы, например .map() и .filter().

В этом примере оператор .transform() используется для передачи каждого значения из восходящего потока столько раз, каково это значение:

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

// A simplified custom implementation of the default .transform() operator
inline fun <T, R> Flow<T>.myTransform(
    // Accepts a suspending lambda that can emit values downstream
    crossinline transform: suspend FlowCollector<R>.(value: T) -> Unit
): Flow<R> = flow {
    // Collects values from the upstream flow
    this@myTransform.collect { value ->
        // Applies the transformation and emits values to the downstream flow
        this@flow.transform(value)
    }
}

// Uses the default .transform() operator
suspend fun main() = withContext(Dispatchers.Default) {
    val flow = (0..4).asFlow().transform { value ->
        // Emits each value as many times as its value
        repeat(value) {
            emit(value)
        }
    }
    println(flow.toList())
    // [1, 2, 2, 3, 3, 3, 4, 4, 4, 4]
}

Оператор .map() позволяет преобразовать каждое значение из восходящего потока в одно значение нисходящего потока.

В этом примере оператор .map() используется для умножения каждого значения на четыре:

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

// A simplified custom implementation of the default .map() operator
inline fun <T, R> Flow<T>.myMap(
    crossinline transform: suspend (value: T) -> R
): Flow<R> = transform { value ->
    emit(transform(value))
}

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Multiplies each upstream value by four
    val flow = (0..4).asFlow().map { it * 4 }
    println(flow.toList())
    // [0, 4, 8, 12, 16]
}
//sampleEnd

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

В этом примере передаются значения, при делении которых на 3 остаток равен 1:

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

// A simplified custom implementation of the default .filter() operator
inline fun <T> Flow<T>.myFilter(
    crossinline predicate: suspend (value: T) -> Boolean
): Flow<T> = transform { value ->
    // Emits only values that match the condition
    if (predicate(value))
        emit(value)
}

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Emits only values where dividing by 3 leaves a remainder of 1
    val flow = (0..10).asFlow().filter { it % 3 == 1 }
    println(flow.toList())
    // [1, 4, 7, 10]
}
//sampleEnd

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

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

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

// A simplified custom implementation of the default .mapNotNull() operator
inline fun <T, R: Any> Flow<T>.myMapNotNull(
    crossinline transform: suspend (value: T) -> R?
): Flow<R> = transform { value ->
    transform(value)?.let { transformed ->
        emit(transformed)
    }
}

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Converts each string to Double and skips values that can't be converted
    val flow = flowOf("1.2", "10", "11", "error", "0.000")
        .mapNotNull { it.toDoubleOrNull() }
    
    println(flow.toList())
    // [1.2, 10.0, 11.0, 0.0]
}
//sampleEnd

Операторы фильтрации и ограничения размера

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

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

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

// A simplified custom version of the default .distinctUntilChanged() operator
fun <T> Flow<T>.myDistinctUntilChanged(): Flow<T> = flow {
    var lastEmitted: Any? = Any() // A value that's equal only to itself
    this@myDistinctUntilChanged.collect { value ->
        if (lastEmitted != value) {
            this@flow.emit(value)
            lastEmitted = value
        }
    }
}

suspend fun main() = withContext(Dispatchers.Default) {
    // Removes repeated consecutive values from the upstream flow
    val flow = flowOf(1, 2, 3, 3, 3, 4, 5, 5, 1).distinctUntilChanged()
    println(flow.toList())
    // [1, 2, 3, 4, 5, 1]
}

Оператор .drop() позволяет пропустить первые значения, переданные восходящим потоком. Например, .drop(2) пропускает первые два значения и передаёт оставшиеся значения дальше:

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

// A simplified custom version of the default .drop() operator
fun <T> Flow<T>.myDrop(count: Int): Flow<T> = flow {
    require(count >= 0)
    var elementsAlreadyDropped = 0
    this@myDrop.collect { value ->
        if (elementsAlreadyDropped == count) {
            this@flow.emit(value)
        } else {
            ++elementsAlreadyDropped
        }
    }
}

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Skips the first two values from the upstream flow    
    val flow = flowOf(1, 2, 3, 4, 5).drop(2)
    println(flow.toList())
    // [3, 4, 5]
}
//sampleEnd

Чтобы отменить сбор после фиксированного количества значений, используйте оператор .take(). В этом примере оператор .take() используется для сбора только первых трёх значений:

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

// A simplified custom 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
        }
    }
}

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Collects only the first three values from the upstream flow
    val flow = (0..1000).asFlow().take(3)

    println(flow.toList())
    // [0, 1, 2]
}
//sampleEnd

Операторы конкурентной обработки

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

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

Один из операторов, добавляющих такой буфер, — оператор .buffer(). Он позволяет настроить ёмкость буфера и поведение при его заполнении, например:

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

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    flow {
        repeat(10) {
            emit(it)
            println("Emitted $it!")
        }
    }
        // Lets the upstream flow emit up to four values ahead of the collector
        .buffer(4)
        .collect {
            println("Processed $it!")
            delay(20.milliseconds)
        }
}
//sampleEnd

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

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

Чтобы отбрасывать значения, а не приостанавливать восходящий поток, задайте параметр onBufferOverflow:

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

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    flow {
        repeat(10) {
            emit(it)
            println("Emitted $it!")
        }
    }
        // Stores up to four values before applying the overflow behavior
        // Drops the oldest buffered value when the buffer is full
        .buffer(4, onBufferOverflow = BufferOverflow.DROP_OLDEST)
        .collect { value ->
            println("Processed $value!")
            delay(20.milliseconds)
        }
}
//sampleEnd

Также можно использовать оператор .conflate() — сокращённую форму записи buffer(1, onBufferOverflow = BufferOverflow.DROP_OLDEST). Используйте его, если нужно обрабатывать только последние значения, пропуская значения, переданные во время сбора предыдущего значения:

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

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    flow {
        repeat(10) {
            emit(it)
            println("Emitted $it!")
        }
    }.conflate().collect {
        println("Processed $it!")
        delay(20.milliseconds)
    }
}
//sampleEnd

Оператор .conflate() влияет только на то, какие значения из буфера обрабатывает сборщик. Он не отменяет уже начатую обработку. Для этого используйте collectLatest().

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

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

Вот упрощённый пример, в котором .flowOn() используется для выполнения восходящего потока в Dispatchers.IO:

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

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    flow {
        repeat(10) {
            emit(it)
        }
        println("Finished emitting!")
    }.flowOn(Dispatchers.IO).collect {
        println("Received $it!")
        delay(10.milliseconds)
    }
}
//sampleEnd

В этом примере оператор .flowOn() может обеспечить конкурентную обработку восходящего потока, но поведение буфера явно не настроено.

Чтобы настроить и контекст корутины для восходящего потока, и поведение буфера, объедините .flowOn() с .buffer() или .conflate(). При совместном использовании эти операторы выполняют слияние операторов и используют один общий буфер.

В этом примере .flowOn(Dispatchers.IO) используется для выполнения восходящего потока в Dispatchers.IO, а .conflate() — для сохранения в буфере самого нового значения:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.Random
import kotlin.time.Duration.Companion.milliseconds
import kotlin.math.round

//sampleStart
fun awaitSensorSignal(): SensorSignal {
    Thread.sleep(10)
    val reading =
        round(Random.nextDouble(25.0, 100.0) * 100.0)/100.0
    println("Measured $reading as the temperature")
    return SensorSignal(temperatureCelsius = reading)
}

data class SensorSignal(
    val temperatureCelsius: Double
)

suspend fun sendLatestTemperature(temperatureCelsius: Double) {
    println("Starting to send $temperatureCelsius...")
    delay(50.milliseconds)
    println("Sent $temperatureCelsius.")
}

suspend fun main() = withContext(Dispatchers.Default) {
    val smartHomeTemperatureFlow = flow {
        while (true) {
            val signal = awaitSensorSignal()
            emit(signal.temperatureCelsius)
            println("Emitted $signal")
        }
    }
        // Runs the upstream flow in Dispatchers.IO
        .flowOn(Dispatchers.IO)
        // Keeps the newest buffered value and drops older ones
        .conflate()
        // Collects the first two values from the upstream flow
        .take(2)
        .collect { temperature ->
            println("Received $temperature!")
            sendLatestTemperature(temperature)
        }
}
//sampleEnd

Операторы объединения

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

Чтобы объединять в пары значения из двух восходящих потоков, используйте оператор .zip(). Он объединяет первые значения каждого потока, затем вторые и так далее. Итоговый поток завершается, как только завершается один из восходящих потоков.

Пример:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlin.random.Random
import kotlin.time.Duration.Companion.milliseconds
import kotlin.time.TimeSource

//sampleStart
suspend fun main() = withContext(Dispatchers.Default) {
    // Emits a ticker value every 100 milliseconds
    val tickerFlow = flow {
        while (true) {
            emit(Unit)
            delay(100.milliseconds)
        }
    }

    val start = TimeSource.Monotonic.markNow()
    tickerFlow
        // Combines each ticker emission with the next number
        .zip(flowOf(1, 2, 3)) { _, value ->
            value
        }.collect {
            println("${start.elapsedNow()}: received $it")
        }
}
//sampleEnd

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

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

//sampleStart
enum class Theme {
    Dark,
    Light,
}

data class UiState(
    val messages: List<String>,
    val theme: Theme,
)

val messagesFlow = MutableStateFlow(
    listOf(
        "Hello!",
        "Is anyone here?",
    )
)

val themeFlow = MutableStateFlow(
    Theme.Light
)

// Combines the latest values from both upstream flows
val uiStateFlow = combine(messagesFlow, themeFlow) { messages, theme ->
    UiState(messages, theme)
}

suspend fun main() {
    withContext(Dispatchers.Default) {
        // Uses UNDISPATCHED to subscribe before the first update happens
        val uiUpdateJob = launch(start = CoroutineStart.UNDISPATCHED) {
            uiStateFlow.collect {
                // Draws the UI
                println(it)
            }
        }
        messagesFlow.update { messages -> messages + "I'll be back!" }
        delay(100.milliseconds)
        
        themeFlow.value = Theme.Dark
        delay(100.milliseconds)
        
        uiUpdateJob.cancel()
    }
}
//sampleEnd

В этом примере combine() создаёт uiStateFlow из последних значений messagesFlow и themeFlow. При обновлении любого из восходящих потоков передаётся новый UiState с последними сообщениями и темой.

Если нужно собирать значения из нескольких потоков конкурентно и передавать их в один нисходящий поток, используйте оператор .merge():

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

//sampleStart
interface UiEvent

class ClickEvent: UiEvent

class RightClickEvent: UiEvent

suspend fun main() {
    withContext(Dispatchers.Default) {

        val clickFlow = MutableSharedFlow<ClickEvent>()
        val rightClickFlow = MutableSharedFlow<RightClickEvent>()

        coroutineScope {
            // Uses UNDISPATCHED to subscribe before the first update happens
            val collectJob = launch(start = CoroutineStart.UNDISPATCHED) {

                // Collects both upstream flows concurrently and emits their values downstream
                merge(clickFlow, rightClickFlow).collect {
                    println("Observed an event: $it")
                }
            }
            clickFlow.emit(ClickEvent())
            delay(100.milliseconds)
            
            clickFlow.emit(ClickEvent())
            delay(100.milliseconds)
            
            rightClickFlow.emit(RightClickEvent())
            delay(100.milliseconds)
            
            collectJob.cancel()
        }
    }
}
//sampleEnd

Операторы жизненного цикла

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

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

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

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

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

// A simplified custom version of the default .onStart() operator
fun <T> Flow<T>.myOnStart(
    action: suspend FlowCollector<T>.() -> Unit
): Flow<T> = flow {
    this@flow.action()
    this@myOnStart.collect(this@flow)
}

suspend fun main() {
    withContext(Dispatchers.Default) {
        flowOf("Page 1", "Page 2", "Page 3").onStart {
            println("Processing pages!")
        }.onEach {
            println("Emitted $it")
        }.collect {
            println("Collected $it")
        }
    }
}

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

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

// A simplified custom version of the default .onCompletion() operator
fun <T> Flow<T>.myOnCompletion(
    action: suspend FlowCollector<T>.(cause: Throwable?) -> Unit
): Flow<T> = flow {
    var exception: Throwable? = null
    try {
        this@myOnCompletion.collect(this@flow)
    } catch (e: Throwable) {
        // Run `action`, but if `action` calls `emit`, throw `e` from it
        FlowCollector<T> { throw e }.action(e)
        throw e
    }
    this@flow.action(null)
}

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        flowOf("Page 1", "Page 2", "Page 3").onCompletion {
            println("Almost done...")
            // Emits an additional value after the upstream flow completes
            emit("Last Page!")
        }.collect {
            println("Collected $it")
        }
    }
}
//sampleEnd

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

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

// A simplified custom version of the default .onEmpty() operator
fun <T> Flow<T>.myOnEmpty(
    action: suspend FlowCollector<T>.() -> Unit
): Flow<T> = flow {
    var emittedSomething = false
    this@myOnEmpty.collect { value ->
        emittedSomething = true
        this@flow.emit(value)
    }
    if (!emittedSomething) {
        action()
    }
}

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        flowOf("Page 1", "Page 2", "Page 3").onEmpty {
            // Doesn't print anything, because the upstream flow emits values
            println("No pages to load!")
        }.collect()
        flowOf<Int>().onEmpty {
            println("No pages to load!")
            // No pages to load!
        }.collect()
    }
}
//sampleEnd

Терминальные операторы

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

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

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

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        flowOf(1, 2, 3).collect {
            println("Collected $it!")
        }
    }
}
//sampleEnd

Также можно вызвать collect() без лямбды:

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

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        flowOf(1, 2, 3).onEach {
            println("Collected $it!")
        }.collect()
    }
}
//sampleEnd

Если нужно собирать поток, отменяя незавершённую работу при передаче нового значения, используйте оператор collectLatest():

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

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        flow {
            println("Emitting Page 1")
            emit("Page 1")
            delay(50.milliseconds)
            println("Emitting Page 2 in quick succession")
            emit("Page 2")
            delay(200.milliseconds)
            println("Emitting Page 3")
            emit("Page 3")
        }.flowOn(Dispatchers.IO).collectLatest {
            println("Starting to process $it!")
            try {
                delay(100.milliseconds)
            } catch (e: CancellationException) {
                println("Canceled processing $it.")
                throw e
            }
            println("Done processing!")
        }
    }
}
//sampleEnd

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

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


suspend fun main() { 
    withContext(Dispatchers.Default) {
        val firstValue = flowOf(1, 2, 3).first()

        println(firstValue)
        // 1
    }
}

Собранные значения можно поместить в коллекцию с помощью операторов .toList() или .toSet():

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

// A simplified custom implementation of the default .toList() operator
suspend fun <T> Flow<T>.myToList(): List<T> = buildList {
    this@myToList.collect { value ->
        // Adds each emitted value to the resulting list
        add(value)
    }
}

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        val list = flowOf(1, 2, 3).toList()
        println(list)
        // [1, 2, 3]

        val set = flowOf(1, 2, 2, 3).toSet()
        println(set)
        // [1, 2, 3]
    }
}
//sampleEnd

Чтобы объединить переданные значения в один результат, используйте операторы .reduce() или .fold(). Оператор .fold() использует переданное вами значение в качестве начального, а оператор .reduce() — первое переданное значение.

Пример:

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

//sampleStart
suspend fun main() {
    withContext(Dispatchers.Default) {
        // Uses the first emitted value as the starting value
        val reduced = flowOf(1, 2, 3).reduce { accumulator, value ->
            accumulator + value
        }

        // Starts with the provided starting value
        val folded = flowOf(1, 2, 3).fold(2) { accumulator, value ->
            accumulator + value
        }

        println(reduced)
        // 6

        println(folded)
        // 8
    }
}
//sampleEnd

Сбор потока в определённой CoroutineScope

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

Чтобы собирать поток в определённой CoroutineScope, используйте терминальный оператор .launchIn(). Этот оператор возвращает Job собирающей корутины.

В этом примере экран собирает значения из StateFlow и останавливает собирающую корутину при закрытии экрана:

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

// A simplified custom version of the default .launchIn() operator
fun <T> Flow<T>.myLaunchIn(scope: CoroutineScope): Job = scope.launch {
    this@myLaunchIn.collect()
}

//sampleStart
data class Coordinate(val x: Int, val y: Int)

class MyScreen(val scope: CoroutineScope) {
    private val _mousePosition =
        MutableStateFlow<Coordinate>(Coordinate(0, 0))
    val mousePosition get() = _mousePosition.asStateFlow()

    init {
        // Starts collecting the StateFlow in the screen's CoroutineScope
        mousePosition.onEach {
            updateStatusBar()
        }.launchIn(scope)
    }

    fun moveMouse(newCoordinate: Coordinate) {
        _mousePosition.value = newCoordinate
    }

    private fun updateStatusBar() {
        println("Mouse is at ${_mousePosition.value}")
    }
}

suspend fun main() {
    withContext(Dispatchers.Default) {
        val childScope = CoroutineScope(
            currentCoroutineContext() + Job(currentCoroutineContext()[Job])
        )
        val screen = MyScreen(childScope)
        delay(100.milliseconds)
        
        screen.moveMouse(Coordinate(10, 15))
        delay(100.milliseconds)
        
        screen.moveMouse(Coordinate(1, 3))
        delay(100.milliseconds)
        
        childScope.cancel()
    }
}
//sampleEnd
28 июля 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-operators.html

Spec-Zone.ru

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