0Pricing
Kotlin Academy · 课时

SharedFlow:事件总线与一次性事件

配置 SharedFlow 的 replay 和 extraBufferCapacity,以广播事件。

SharedFlow:事件总线与一次性事件 是 CoddyKit 上的免费 Kotlin Academy 课时。 这是第 2 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 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 参数

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()
}

一次性事件(界面导航)

对于导航或显示提示条等一次性界面事件,请使用 replay=0 的 SharedFlow——事件不会在重组时重复播放。

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(默认)、丢弃最旧项或丢弃最新项。

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 是一种热事件广播器。一次性事件使用重放值 0,较晚订阅的订阅者使用大于 0 的重放值。在非挂起上下文中使用 tryEmit。事件应选择 SharedFlow,状态应选择 StateFlow。

常见问题解答

「SharedFlow:事件总线与一次性事件」课时是免费的吗?

是的 — 「SharedFlow:事件总线与一次性事件」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Kotlin Academy 课程的其余内容,请升级到 CoddyKit PRO。 Kotlin Academy 课程共包含 4 节课。

「SharedFlow:事件总线与一次性事件」这节课中我会学到什么?

配置 SharedFlow 的 replay 和 extraBufferCapacity,以广播事件。 你通过在浏览器中直接运行的动手代码来练习 Kotlin Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 Kotlin Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 Kotlin Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 2 节课,共 4 节。

「SharedFlow:事件总线与一次性事件」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 Kotlin Academy 课中编写并运行代码吗?

能。每节 Kotlin Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. StateFlow:用于界面的热状态容器
  2. SharedFlow:事件总线与一次性事件
  3. 使用 shareIn 与 stateIn 将冷 Flow 转换为热 Flow
  4. 使用 Turbine 测试 StateFlow 与 SharedFlow
← 返回 Kotlin Academy