Kotlin Academy · Урок

SharedFlow: шины событий и одноразовые события

Настраивайте replay и extraBufferCapacity у SharedFlow для широковещательной передачи событий.

Урок 2 из 413 шагов

«SharedFlow: шины событий и одноразовые события» — бесплатный урок Kotlin Academy на CoddyKit. Это урок 2 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Kotlin Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Kotlin Academy содержит 4 уроков всего.

Что такое SharedFlow?

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

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val events = MutableSharedFlow<String>()
    launch { events.collect { println("Collector 1: $it") } }
    launch { events.collect { println("Collector 2: $it") } }
    delay(50)
    events.emit("UserLoggedIn")
    events.emit("DataRefreshed")
    delay(50)
    coroutineContext.cancelChildren()
}

Параметр повторной выдачи

replay хранит в буфере последние N выдач. Новые подписчики сразу получают до N предыдущих событий. Значение по умолчанию — 0, то есть повторной выдачи нет.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val flow = MutableSharedFlow<Int>(replay = 2)
    flow.emit(1); flow.emit(2); flow.emit(3)
    // Late subscriber gets last 2:
    flow.collect { print("$it ") }  // prints 2 3
}

extraBufferCapacity

extraBufferCapacity добавляет место в буфере сверх replay. Отправители могут выдавать значения без приостановки, пока буфер не заполнится.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val flow = MutableSharedFlow<Int>(
        replay = 0,
        extraBufferCapacity = 10
    )
    // Can emit up to 10 times without a collector ready:
    repeat(10) { flow.tryEmit(it) }
    launch { flow.collect { print("$it ") } }
    delay(50)
    coroutineContext.cancelChildren()
}

Выдача без приостановки с помощью tryEmit

tryEmit(value) выдаёт значение без приостановки и возвращает false, если буфер заполнен. Используйте его в контекстах без приостановки: в обратных вызовах и обработчиках нажатий.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val flow = MutableSharedFlow<String>(extraBufferCapacity = 5)
    val ok = flow.tryEmit("click")  // non-suspending
    println("Emitted: $ok")
    launch { flow.collect { println(it) } }
    delay(50)
    coroutineContext.cancelChildren()
}

Одноразовые события (навигация интерфейса)

Используйте SharedFlow с replay=0 для одноразовых событий интерфейса, таких как навигация или показ всплывающего уведомления: события не выдаются повторно при повторной композиции.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
class NavViewModel {
    private val _events = MutableSharedFlow<String>()
    val events = _events.asSharedFlow()
    fun navigateTo(route: String) {
        // viewModelScope.launch:
        kotlinx.coroutines.GlobalScope.launch { _events.emit(route) }
    }
}
// Collect in UI:
// viewModel.events.collect { route -> navController.navigate(route) }

SharedFlow и StateFlow

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

import kotlinx.coroutines.flow.*
// StateFlow: always has a value; new subscribers get current value
val uiState = MutableStateFlow("idle")

// SharedFlow: no stored value; events fire and are gone (unless replay>0)
val singleEvents = MutableSharedFlow<String>(replay = 0)

fun main() { println("State = what it is; Event = what happened") }

Шина событий с помощью SharedFlow

Реализуйте простую общую для приложения шину событий с помощью единственного экземпляра SharedFlow, заменив шаблоны с RxJava PublishSubject.

import kotlinx.coroutines.flow.*
object EventBus {
    private val _events = MutableSharedFlow<Any>(extraBufferCapacity = 100)
    val events = _events.asSharedFlow()
    fun post(event: Any) { _events.tryEmit(event) }
}
sealed class AppEvent {
    object UserLoggedOut : AppEvent()
    data class ShowError(val msg: String) : AppEvent()
}
// Usage: EventBus.post(AppEvent.UserLoggedOut)

subscriptionCount

sharedFlow.subscriptionCount — это StateFlow, отслеживающий число активных подписчиков. Он полезен для запуска и остановки вышестоящих производителей данных.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val flow = MutableSharedFlow<Int>()
    println("Subscribers: ${flow.subscriptionCount.value}")  // 0
    val job = launch { flow.collect { } }
    delay(50)
    println("Subscribers: ${flow.subscriptionCount.value}")  // 1
    job.cancel()
    delay(50)
    println("Subscribers: ${flow.subscriptionCount.value}")  // 0
}

resetReplayCache

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

import kotlinx.coroutines.flow.*
fun main() {
    val flow = MutableSharedFlow<Int>(replay = 3)
    flow.tryEmit(1); flow.tryEmit(2); flow.tryEmit(3)
    println("Replay cache: ${flow.replayCache}")  // [1, 2, 3]
    flow.resetReplayCache()
    println("Replay cache: ${flow.replayCache}")  // []
}

Обработка обратного давления в SharedFlow

Если сборщики работают медленно, используйте onBufferOverflow, чтобы выбрать стратегию: SUSPEND (по умолчанию), DROP_OLDEST или DROP_LATEST.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.channels.BufferOverflow
fun main() = runBlocking {
    val flow = MutableSharedFlow<Int>(
        replay = 0,
        extraBufferCapacity = 3,
        onBufferOverflow = BufferOverflow.DROP_OLDEST
    )
    repeat(10) { flow.tryEmit(it) }
    launch {
        flow.collect { println(it) }  // only sees most recent 3
    }
    delay(50)
    coroutineContext.cancelChildren()
}

Сбор с ограничением времени

Собирайте SharedFlow с ограничением времени, чтобы обработать конечное число событий, а затем остановиться; это удобно для тестирования или обработки с ограничением объёма.

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val events = MutableSharedFlow<String>()
    launch { repeat(5) { delay(50); events.emit("Event $it") } }
    withTimeoutOrNull(200) {
        events.collect { println(it) }
    }
    println("Done collecting")
}

Быстрая проверка

Какое значение повторной выдачи следует использовать для одноразовых событий интерфейса, например навигации?

Повторение

SharedFlow — это горячий широковещательный поток событий. Используйте replay=0 для одноразовых событий и replay>0 для поздних подписчиков. Используйте tryEmit в контекстах, где нельзя приостанавливать выполнение. Выбирайте SharedFlow для событий, а StateFlow — для состояния.

Можно начать бесплатно

Изучай Kotlin с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
51
Уроки
203

Часто задаваемые вопросы

Урок «SharedFlow: шины событий и одноразовые события» бесплатный?

Да — полный текст урока «SharedFlow: шины событий и одноразовые события» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Kotlin Academy, подпишись на CoddyKit PRO. Курс Kotlin Academy содержит 4 уроков всего.

Чему я научусь в уроке «SharedFlow: шины событий и одноразовые события»?

Настраивайте replay и extraBufferCapacity у SharedFlow для широковещательной передачи событий. Ты практикуешь Kotlin Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Kotlin Academy?

Предыдущий опыт не требуется. Kotlin Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 2 из 4.

Сколько времени занимает урок «SharedFlow: шины событий и одноразовые события»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке Kotlin Academy?

Да. Каждый урок Kotlin Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. StateFlow: горячий контейнер состояния для интерфейса
  2. SharedFlow: шины событий и одноразовые события
  3. Преобразование холодного Flow в горячий с shareIn и stateIn
  4. Тестирование StateFlow и SharedFlow с Turbine
← Назад к Kotlin Academy