You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何分流Kotlin SharedFlow?现有实现仅获取最后创建的流数据

Kotlin SharedFlow 分流优化方案

你的代码出现仅能获取最后创建的SharedFlow数据的问题,核心原因有两个:

  1. 重复共享原始冷Flow:每次调用evens()/odds()时,都对原始冷Flow执行shareIn(),这会创建多个独立的SharedFlow实例,且原始冷Flow会被多次订阅执行。加上你在Flow内部使用runBlocking(delay(100))阻塞线程,导致先启动的订阅被阻塞,后启动的订阅会重新触发原始Flow从头执行,最终只有最后一个订阅能正常输出。
  2. 错误使用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 }
}

优化说明

  1. 单一共享源:仅对原始冷Flow执行一次shareIn(),所有分流操作都基于这个唯一的SharedFlow,避免原始Flow被多次订阅(对应你的游戏控制器场景,就是避免重复调用轮询API)。
  2. 移除runBlocking:使用挂起函数delay()替代阻塞式的runBlocking,符合协程的挂起逻辑,避免线程阻塞问题。
  3. 按需共享过滤结果:如果过滤后的Flow需要被多个下游订阅,可以在过滤后再调用shareIn(),但大部分场景下,基于SharedFlow的过滤操作返回的Flow本身就是热流,无需再次共享。

内容的提问来源于stack exchange,提问作者Curtis

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.22 20:27:30