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 反馈 — 无需本地设置。
此课程中的所有课时
- StateFlow:用于界面的热状态容器
- SharedFlow:事件总线与一次性事件
- 使用 shareIn 与 stateIn 将冷 Flow 转换为热 Flow
- 使用 Turbine 测试 StateFlow 与 SharedFlow