Операторы 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().
В этом примере эти операторы используются для вывода сообщения перед началом сбора и перед передачей каждого значения вниз по потоку:
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
© 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