如何分流Kotlin SharedFlow?现有实现仅获取最后创建的流数据
你的代码出现仅能获取最后创建的SharedFlow数据的问题,核心原因有两个:
- 重复共享原始冷Flow:每次调用
evens()/odds()时,都对原始冷Flow执行shareIn(),这会创建多个独立的SharedFlow实例,且原始冷Flow会被多次订阅执行。加上你在Flow内部使用runBlocking(delay(100))阻塞线程,导致先启动的订阅被阻塞,后启动的订阅会重新触发原始Flow从头执行,最终只有最后一个订阅能正常输出。 - 错误使用
runBlocking:Flow的收集器本身处于挂起上下文,直接使用挂起函数delay()即可,runBlocking会阻塞线程,破坏协程的挂起机制,加剧执行顺序问题。
基础优化方案(复用共享源)
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* fun main() = runBlocking { launch { Flows.evens(this).collect { println("a: $it") } } launch { Flows.odds(this).collect { println("b: $it") } } } object Flows { private val flow: Flow<Int> = flow { for (i in 0..1000) { emit(i) delay(100) // 替换runBlocking,直接使用挂起delay } } private val sharing = SharingStarted.Eagerly // 只对原始Flow做一次共享,生成唯一多播源 private fun sharedSource(scope: CoroutineScope): Flow<Int> = flow.shareIn(scope, sharing) // 基于共享源做分流过滤,无需再次shareIn fun evens(scope: CoroutineScope) = sharedSource(scope).filter { it % 2 == 0 } fun odds(scope: CoroutineScope) = sharedSource(scope).filter { it % 2 == 1 } }
更优预定义方案(适合固定分流场景)
如果你的分流规则是固定的(比如游戏控制器的特定按钮/摇杆),可以提前预定义好分流后的Flow,避免每次调用函数重复处理:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* fun main() = runBlocking { launch { Flows.evens.collect { println("a: $it") } } launch { Flows.odds.collect { println("b: $it") } } } object Flows { // 定义共享Flow的协程作用域,根据实际场景调整(如Android可用viewModelScope) private val sharedScope = CoroutineScope(Dispatchers.Default + SupervisorJob()) private val sharing = SharingStarted.Eagerly private val flow: Flow<Int> = flow { for (i in 0..1000) { emit(i) delay(100) } } // 预创建唯一的共享数据源 private val sharedSource = flow.shareIn(sharedScope, sharing) // 预定义分流后的Flow,直接供外部使用 val evens = sharedSource.filter { it % 2 == 0 } val odds = sharedSource.filter { it % 2 == 1 } }
优化说明
- 单一共享源:仅对原始冷Flow执行一次
shareIn(),所有分流操作都基于这个唯一的SharedFlow,避免原始Flow被多次订阅(对应你的游戏控制器场景,就是避免重复调用轮询API)。 - 移除
runBlocking:使用挂起函数delay()替代阻塞式的runBlocking,符合协程的挂起逻辑,避免线程阻塞问题。 - 按需共享过滤结果:如果过滤后的Flow需要被多个下游订阅,可以在过滤后再调用
shareIn(),但大部分场景下,基于SharedFlow的过滤操作返回的Flow本身就是热流,无需再次共享。
内容的提问来源于stack exchange,提问作者Curtis
相关产品推荐
相关产品推荐

