公众号
CSDN原文

什么是数据流

数据流是一种数据处理方式,数据被异步 连续传输和处理,像水从高处流到低处,流的过程中可以处理水,过滤加糖 加热等的,下流是消费端,直接消费喝掉就行。

Kotlin的 Flow 是什么

在Kotlin里 对数据流建模合适的类型就是Flow,从概念上讲,可以异步计算的数据流称之为Flow。Flow 是Kotlin的协程的一部分,提供对数据流的产生 变换 组合 和消费的强大支持。通过Flow可以轻松处理异步的数据流。

Flow可以连续的发出多个值,而不像传统函数,调一次就是一个值。还有区别在于,Flow是使用挂起函数和异步来消费和生产的。
Kotlin Flow 并不是唯一的数据流,但它是协程的一部分,所以和协程配合的很好。类似Rxjava

数据流有三个节点

  • 生产者 :把数据放入流中
  • 加工者 :中间数据加工的节点
  • 消费者 :流的末端消费最终的数据

和LiveData区别在哪

LiveData能保存数据,能感知生命周期,利用观察者模式在可用的生命周期范围内将最新的数据通知给观察者。初中部是为了简化开发,上手简单,但是不适用于复杂场景,比如只能在主线程操作数据,操作符不够强大,不支持切线程(可以切,自己不带切的功能);
Kotlin的协程是Kotlin的并发编程工具, 同步的方式写出异步的代码,简化异步逻辑代码,通过提供结构化并发模式,让异步代码更直观,更好理解,协程也可以恢复和暂停,可以挂起到满足一定条件再执行。这些特性能更方便处理复杂逻辑

Flow 是liveData的补充 增强版,弥补了LiveData的局限性;
对比:

LiveDataFlow
支持Java和Kotlin,使用简单仅Kotlin,java使用困难
无需协程协程环境执行
主线程处理数据协程的线程池里,不阻塞
转换运算符在主线程上执行转换是挂起,在不同线程执行
默认感知生命周期scope来感知

Flow 比rxjava的优势

相比RxJava,Kotlin Flow有什么优势?
可能很多人都用过RxJava,前几年,在纯Java的Android项目中,RxJava对于响应式编程非常友好。但是上手难度还是比较大的,要学很久才知道怎么用,以及怎么才能用好。而后面,大家开始用Kotlin,然后用Kotlin协程,又有了Kotlin Flow,上手难度比RxJava轻松了不少。那么,相比RxJava,Kotlin Flow有哪些优势呢?

更自然的协程支持:Kotlin Flow是集成在Kotlin协程里面的,能更好地利用协程的特性,而且不需要额外引入其他的库。
更简单的语法和易用性:Kotlin Flow的API设计更加简洁,避免了RxJava中复杂的操作符,它利用了扩展函数和lambda表达式,使代码更直观易读。
内存安全与上下文一致性:Kotlin Flow中,数据流的上下文和生命周期是由协程管理的,这意味着可以更容易地处理内存泄漏和取消操作。相比之下,RxJava 需要手动处理订阅的管理和内存泄漏问题。
冷流与热流:Kotlin Flow默认是冷流,即只有在有收集器时才开始执行。这与 RxJava 的 Observable 类似,但更符合大多数使用场景。RxJava 中则需要使用不同的类型(如 Observable 和 Flowable)来区分冷流和热流。
背压处理:Kotlin Flow的冷流特性天然支持背压处理,因为生产者只有在有收集器请求数据时才会产生数据。RxJava 的 Flowable 虽然也支持背压,但需要额外配置和处理,增加了复杂性。
更好的错误处理:Kotlin Flow依赖于 Kotlin 协程的异常处理机制,使得错误处理更加直观。RxJava 中则需要使用 onErrorReturn、onErrorResumeNext 等操作符来处理错误,语法相对复杂。
轻量级和性能:Kotlin Flow相对 RxJava 更轻量,因为它不需要包含 RxJava 的所有操作符和特性。对于大多数常见的异步数据流处理场景,Kotlin Flow 提供了足够的功能,且性能通常更好。
更好的与Kotlin标准库集成:Kotlin Flow是 Kotlin 标准库的一部分,因此与其他 Kotlin 标准库功能(如集合操作、标准函数)无缝集成。这使得开发者可以更自然地使用 Kotlin 语言特性,减少了学习曲线。

自我理解:1 更简单易用,轻,Kotlin的标准库,2 内存安全上更好,避免内存泄露,3 自带冷热流和背压 无需特殊配置,4 错误处理同协程的异常处理,一致易用

使用

创建和消费

fun main(): Unit = runBlocking {
    // 创建3个Flow,生产数据
    val firstFlow = flowOf(1, 2)
    val secondFlow = flow {
        emit(3)
        emit(4)
    }
    val thirdFlow = listOf(5, 6).asFlow()

    // 挨个收集,消费者
    firstFlow.collect {
        println(it)
    }
    secondFlow.collect {
        println(it)
    }
    thirdFlow.collect {
        println(it)
    }
}

中间操作符–接受前转换

val firstFlow = flowOf(1, 2)

// 将数据做 +2 处理
firstFlow.map {
    it + 2
}.collect {
    println(it)
}

转换操作符–发送前转换

map:对每个元素应用一个函数,并返回一个新的 Flow。(和集合的map一样)

flowOf(1, 2, 3).map { it * 2 }

filter:过滤出符合条件的元素。(和集合的filter一样)

flowOf(1, 2, 3).filter { it % 2 == 0 }

transform:对每个元素应用一个自定义的转换,可以发射多个值(这是它和map的区别)。

flowOf(1, 2, 3).transform { value ->
    emit(value * 2)
    emit(value * 3)
}

take:只取前 n 个元素。(和集合的take一样)

flowOf(1, 2, 3, 4).take(2)

组合操作符

Zip
combine

zip:合并两个 Flow 的元素,形成一个新的 Flow。(和集合的zip差不多)

flowOf(1, 2).zip(flowOf("A", "B")) { a, b -> "$a -> $b" }

combine:合并两个 Flow 的最新值。组合最新的值是什么意思呢?两个Flow中任意一个Flow有新的数据来了,那么就需要与另外一个Flow的最新的值进行组合。比如flow2的最新值是A,那么flow1一旦emit发射一个新值1,那么A就会和1结合,flow1再emit发射一个新值2,还是和flow2的最新值A进行结合,这样组合出来的个数就不一定是某个flow数据流的个数。

 val flow1 = flow {
     emit(1)
     delay(100)
     emit(2)
     delay(100)
     emit(3)
 }

 val flow2 = flow {
     emit("A")
     delay(500)
     emit("B")
     emit("C")
 }

 val combinedFlow = flow1.combine(flow2) { a, b -> "$a$b" }

 combinedFlow.collect { println(it) }
 
 // 输出
 // 1A
 // 2A
 // 3A
 // 3B
 // 3C

flatMapConcat:串行地展开一个 Flow。
flatMapMerge:并行地展开一个 Flow。
flatMapLatest:只保留最新展开的 Flow。

3.2.3 末端操作符
collect:收集流的元素并执行给定的动作。

flowOf(1, 2, 3).collect { println(it) }

toList:将 Flow 转换为 List。

val list = flowOf(1, 2, 3).toList()

first:获取第一个元素并终止流的收集。

val first = flowOf(1, 2, 3).first()

上下文相关

flowOn:改变 Flow 的执行上下文。flowOn能改变上游的数据流的执行上下文,collect内部执行的上下文是collect调用处的上下文

flow {
    for (i in 1..3) {
        println("flow  ${currentCoroutineContext()}")
        emit(i)
    }
}.flowOn(Dispatchers.Default)
    .map {
        println("map  ${currentCoroutineContext()}")
        it.toString()
    }
    .flowOn(Dispatchers.IO)
    .collect {
        withContext(Dispatchers.IO) {
            println("collect withContext ${currentCoroutineContext()}")
        }
        println("collect ${currentCoroutineContext()}")
        println(it)
    }

/*
输出:
flow  [ProducerCoroutine{Active}@3b6f6746, Dispatchers.Default]
flow  [ProducerCoroutine{Active}@3b6f6746, Dispatchers.Default]
flow  [ProducerCoroutine{Active}@3b6f6746, Dispatchers.Default]
map  [ScopeCoroutine{Active}@4b60f5ce, Dispatchers.IO]
map  [ScopeCoroutine{Active}@4b60f5ce, Dispatchers.IO]
map  [ScopeCoroutine{Active}@4b60f5ce, Dispatchers.IO]
collect withContext [DispatchedCoroutine{Active}@8945c64, Dispatchers.IO]
collect [ScopeCoroutine{Active}@6fdb1f78, BlockingEventLoop@51016012]
1
collect withContext [DispatchedCoroutine{Active}@437a60dc, Dispatchers.IO]
collect [ScopeCoroutine{Active}@6fdb1f78, BlockingEventLoop@51016012]
2
collect withContext [DispatchedCoroutine{Active}@7b238e10, Dispatchers.IO]
collect [ScopeCoroutine{Active}@6fdb1f78, BlockingEventLoop@51016012]
3
*/


背压 --buffer

生产和消费速度不一致的时候需要做的背压
buffer 本质上就是通过弄两个协程上下文来完成的,buffer前一个协程,后一个协程

// 先看一下有缓冲区的情况
flowOf("A","B","C","D","E")
  .onEach { println("Woman matchmaker emits: $it") }
  .buffer()
  .collect {
      println("Girl appointment with: $it")
      delay(1000)
  }

//输出
Woman matchmaker emits: A
Woman matchmaker emits: B
Woman matchmaker emits: C
Woman matchmaker emits: D
Woman matchmaker emits: E
Girl appointment with: A
Girl appointment with: B
Girl appointment with: C
Girl appointment with: D
Girl appointment with: E


// 无缓冲区的情况
flowOf("A","B","C","D","E")
  .onEach { println("Woman matchmaker emits: $it") }
  .collect {
      println("Girl appointment with: $it")
      delay(1000)
  }

// 输出
Woman matchmaker emits: A
Girl appointment with: A
Woman matchmaker emits: B
Girl appointment with: B
Woman matchmaker emits: C
Girl appointment with: C
Woman matchmaker emits: D
Girl appointment with: D
Woman matchmaker emits: E
Girl appointment with: E

conflate 只发最新值,collectLatest只接受最新的

flowOf(1, 2, 3).conflate()

错误处理操作符

catch 补和处理异常

flow {
    emit(1)
    throw RuntimeException("RuntimeException")
}.catch { e ->
    emit(-1)
}.collect {
    println(it)
}

retry 重复的流 最多重试次数

flow {
   emit(1)
   throw RuntimeException("RuntimeException")
}.retry(3).collect {
   println(it)
}

// 输出
1
1
1
1
Exception in thread "main" java.lang.RuntimeException: RuntimeException

两者组合起来用,retry几次后不成功再catch下来

冷流和热流

冷 Flow

  • 数据生产和消费是绑定的:冷流是懒惰的,数据只有在消费者 collect的时候才会生产,意味着一个新的消费者都会出发一个新的数据流
  • 多个消费者独立消费:每个消费都独立于其他消费者,每个消费者都是从头接受数据,每个消费者都是拥有自己独立的数据流
  • 适合单播场景:冷流更适合单独的场景,每个消费者都独立消费数据,不受其他消费者的影响

热 StateFlow SharedFlow

  • 数据生产者独立于消费者,热流会创建后立即开始生产,不管是否有消费者,说明生产和消费是独立的。
  • 多个消费者共享一个数据流:新加入的只能接收到最新的数据,而不是从头开始
  • 适合多播场景:一个数据被多个消费者同时消费的

StateFlow SharedFlow MutableShareFlow MutableStateFlow

StateFlow 和 SharedFlow 是热流,Mutable是可读可写版本
stateFlow 和livedata类似 :

  • 提供可读可写 仅可读的两个版本 (stateFlow MutableStateFlow)
  • 值唯一
  • 允许多个观察者
  • 只会发最新的值重现给订阅者
  • 支持DataBinding
    不同的:
  • 必须初始值
  • value空安全
  • 防抖 默认防抖(发送会判断和目前值相同,相同就不发送)

SharedFlow MutableSharedFlow
区别:

  • 无初始值
  • 不防抖 可以连续发生相同数据
    场景:
    stateFlow 适合状态管理 最后是什么状态,改成最终的状态即可,一样的状态不需要更新
    sharedFlow 适合通知做事情的,每次通知都要执行的

stateIn 和 ShareIn

Flow 里的拓展函数。可能冷流转成热流

stateIn 就是冷转成 StateFlow。 shareIn就是转成ShareFlow

stateIn 要穿三个参数:

  • scope:作用域

  • started: 启动策略

    • SharingStarted.Lazily : 第一个订阅者出现的时候,开始运转,当scope取消的时候才停止。
    • SharingStarted.Eagerly : 立即启动,当scope取消的时候才停止。
    • SharingStarted.WhileSubscribed(stopTimeoutMillis: Long,replayExpirationMillis: Long) : 当至少有一个订阅者的时候启动,最后一个订阅者停止订阅之后还能继续保持stopTimeoutMillis时间的活跃,之后才停止。replayExpirationMillis直接翻译过来是重播过期时间,默认是Long.MAX_VALUE,当取消协程之后,这个缓存的值需要保留多久,如果是0,表示立马就过期,并把shareIn运算符的缓存值设置为initialValue初始值。
  • initialValue 初始值

官方建议:对于那些一次性的操作来说,你可以使用Lazily、Eagerly,但是对于需要观察其他的Flow的情况来说,更推荐用WhileSubscribed。WhileSubscribed在最后一个订阅者停止订阅之后还能继续保持stopTimeoutMillis时间的活跃,之后才停止,这有个好处,比如用户将app切换到后台,此时,上游的流就没必要继续产生并发射数据了,app都没在前台了,发射数据有点浪费资源了。但是,当app只是从竖屏切换到横屏状态时(lifecycle.repeatOnLifecycle那里我们传入的是STARTED,所以会被取消,重新STARTED时会重新collect),这种就没必要取消上游的生产者生成数据了,所以有个stopTimeoutMillis的时间值在那里。官方表示,合适的时间是5000毫秒。

shareIn

shareIn和stateIn参数差不多,但是没有初始值,多了一个replay参数:假设:之前有订阅者,并且已经有3个流数据了,replay=1,这时再来一个订阅者,那么就会发射最新的那个值给这个新的订阅者,而不会发射给这个新的订阅者早先的第一个和第二个数据。

val shareInFlow = flow {
        emit(1)
        delay(300L)
        emit(2)
        delay(300L)
        emit(3)
    }.shareIn(viewModelScope, SharingStarted.WhileSubscribed(5000L), 1)//1改成3 就会接收到123

fun testShareIn() {
    viewModelScope.launch {
        launch {
            shareInFlow.collect {
                log("订阅者1 shareInFlow data $it")
            }
        }
        delay(1000L)
        launch {
            shareInFlow.collect {
                log("订阅者2 shareInFlow data $it")
            }
        }
    }
}

// 输出:
订阅者1 shareInFlow data 1
订阅者1 shareInFlow data 2
订阅者1 shareInFlow data 3
订阅者2 shareInFlow data 3

回调里怎么能到数据转换的flow. callback

类似Kotlin里的 suspendCancellableCoroutine,将callback 转换成协程风格,那么回调的callback拿到的数据都怎么转成flow?? callbackFlow 是Flow的构造器,允许回调中发射数据

Flow 感知生命周期–lifecycleScope里非destroy都会收集,所以要处理–launchWhenStarted

默认是不具备生命周期的感知能力,又几个招可以,asLiveData方法-
repeatOnLifecycle 可以限制到什么状态

// Update the uiState
lifecycleScope.launch {
    lifecycle.repeatOnLifecycle(Lifecycle.State.STARTED) {
        viewModel.uiState
            .onEach { uiState = it }
            .collect()
    }
}

你可能会问了?我都在lifecycleScope里面进行collect了,为啥还需要考虑生命周期的问题,不是自带感知吗? 答案是:这种方式会在ui非DESTROY时一直可以collect,即使app在后台,即onStop的状态时,也会进行收集,然后对ui进行更新,除非你确实有这个需求,那么不然就有点浪费资源了。

那么,在lifecycleScope.launchWhenStarted里面收集Flow数据应该没问题了吧?还是有一点问题,在ui层达到onStop状态之后,但并未destroy时,比如按home键,app此时在后台活着,这个时候Flow的管道还继续存在,且Flow的生产方还可以继续生产并emit。

FLow处理配置变更

用StateFlow 不变就不发

Flow 和 LiveData 的实现方式

发送方

// ViewModel
private val _livedata1 = MutableLiveData<String?>()
val livedata1 = _livedata1
fun fetchData1() {
    viewModelScope.launch(Dispatchers.IO) {
        val result = api.listRepos()
        _livedata1.postValue(result?.toString())
    }
}

val flow1 = flow<String?> {
    val result = api.listRepos()
    emit(result.toString())
}.flowOn(Dispatchers.IO)
    .stateIn(viewModelScope, SharingStarted.WhileSubscribed(5000L), null)

接收方

flowViewModel.fetchData1()
flowViewModel.livedata1.observe(this) {
    log("livedata1 数据 $it")
}

lifecycleScope.launch {
    lifecycle.repeatOnLifecycle(Lifecycle.State.STARTED) {
        flowViewModel.flow1.collect {
            log("flow1 数据 $it")
        }
    }
}

//更多例子可以看原文

更多细节

  • null 也能 emit 发送
  • 相同的对象,也可以emit多次,然后collect收集到–flowOf默认不防抖,stateFlow防抖
  • zip合并的flow 多的数据会被丢弃,
  • stateFlow 无法连续emit解决:包裹+时间戳,或者换sharedFlow
  • 多次联系collect导致的问题–放到各自的launch里
  • Flow里没有异常处理,最终异常会走到哪—走到zygoteInit.main 最终崩溃
    -使用flow构造,生产者无法emit 来自不同coroutineContext,因此,不要通过创建新协程或使用 withContext 代码块在不同的 CoroutineContext 中调用 emit 。在这些情况下,可以使用其他流程构建器,例如 callbackFlow
错误代码!!!
fun errorUseFlow1(): Flow<Int> = flow {
   emit(1)
   // 下面这种是错误的用法
   withContext(Dispatchers.IO) {
       emit(2)
   }
}

5.6 Flow和LiveData我到底怎么选?要不要迁移已有的代码到Flow?—copy原文

LiveData仍然是Java项目、Android初学者、简单场景下的最佳选择。

对于除上面以外的其他情况,官方是建议使用Kotlin Flow。但是学习Kotlin Flow需要一些学习时间,但它是Kotlin语言的一部分,谷歌很挺这个东西。我去简单看了下NowInAndroid这个项目(该项目是谷歌官方的一个Android demo实战项目,功能齐全、完全使用Kotlin和Compose。遵循 Android 设计和开发最佳实践,旨在为开发人员提供有用的参考),我发现里面已经完全没有在使用LiveData了,全是Flow和各种Hilt依赖注入。理论上LiveData能做的事,Kotlin Flow也能做;LiveData不能做或者做起来比较困难的事,Kotlin Flow也能做。

我们可以打开LiveData官网,即使是2024年,谷歌也没有将它标记为过时,它仍然是非常棒的选择。现在和将来很长一段时间理论上都不会被标记为过时,毕竟还有那么多Java项目,而且主要是用起来也非常简单方便。从谷歌2022年的一个采访(Architecture: Live Q&A - MAD Skills)里面可以看出,谷歌意思是你想用LiveData就继续用,当然,更推荐你用Kotlin Flow。

至于要不要迁移,我的理解是,老代码就不动它。新代码,喜欢就可以用Flow,不喜欢就还是可以继续用LiveData。

6. Kotlin Flow 小结

好了,文章比较长,咱们再来回忆一下,主要介绍了Kotlin Flow的相关知识,包括基本概念、基本使用、实际应用以及一些需要注意的问题。Kotlin Flow是Kotlin协程的一部分,用于处理异步数据流,它相比LiveData和RxJava具有诸多优势,如更自然的协程支持、简单的语法、内存安全、更好的错误处理等。在基本使用方面,介绍了Flow的创建、消费、操作符、类型以及如何将回调转换为Flow、让Flow具备生命周期感知能力和处理配置变更问题等。在实际应用中,展示了Flow在请求网络、与Room结合使用以及替代LiveData解决问题等场景的用法。接着还有一些细节方面的讨论,如StateFlow无法连续emit的解决办法、多次连续的collect可能导致的问题、Flow中异常的处理、使用flow构建器的注意事项以及Flow和LiveData的选择和迁移问题。

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐