Kotlin Flow 一 Flow 的创建和使用

Flow 和 RxJava 差不多,不过 Flow 是和协程一起使用的 API。

简单的例子

flow {
    emit(1)
    emit(2)
}.collect {
    println(it)
}

在 Flow 中可以使用 emit 发送数据,相当于 RxJava 中的 onNext()

创建 Flow

Flow 可以使用函数,集合,数组等来进行创建。

  • flowOf()

    flowOf(1, 2, 3, 4, 5)
        .onEach {
            delay(it * 100L)
        }
        .collect {
            println(it)
        }
    
    
  • asFlow

    listOf(1, 2, 3).asFlow()
        .collect {
            println(it)
        }
    arrayOf('a', 'b', 'c')
            .asFlow()
            .collect {
                println(it)
            }
    
    

    asFlow 可以将任何使用 Iterator 和数组转换为 Flow

    {
       1 + 1
    }
    .asFlow()
        .collect {
            println(it)
        }
    suspend {
        delay(1000)
        1 + 1
    }
        .asFlow()
        .collect { 
            println(it)
        }
    
    

    asFlow 还可以将函数(挂起函数也可以)的结果转换为 flow

  • channelFlow

    channelFlow {
        for (i in 1..5) {
            delay(i * 100L)
            send(i)
        }
    }.collect {
        println(it)
    }
    
    

    channelFlow 是通过协程中的 channel 来实现的,所以它的发送数据和普通的 flow 不太一样。 channelFlow 是通过 send 来发送的。

    flow 是冷的流。在没有切换协程的情况下,生产者和消费者是同步非阻塞的。

    channel 是热流。 channelFlow 实现了生产者和消费者的异步非阻塞模型。

切换协程

Flow 是专用于 Kotlin 的 API,所以在会对协程进行切换,而不是协程。

flow {
    for (i in 1..5) {
        delay(100L)
        emit(i)
    }
}
    .map {
        it * it
    }
    .flowOn(Dispatchers.IO)
    .collect {
        println(it)
    }

相比于 RxJava,flow 切换协程更加简单,只需要调用 flowOn 即可切换协程。

Flow 的取消

如果 Flow 是在一个挂起函数中被挂起了,那么 Flow 是可以取消的,否则不可以取消。

Terminal flow operator

Flow 的 API 和 Java 的相似,都有 Intermediate Operations、Terminal Operations。

Flow 的 Terminal 运算符可以是 suspend 函数,如 collect、single、reduce、toList 等,也可以是 launchIn 运算符,用户在指定的 CoroutineScope 内使用 flow。

整理一下 Flow 的 terminal operator,如下:

  • collect
  • single/first
  • toList/toSet/toCollection
  • count
  • fold/reduce
  • launchIn/produceIn/broadcastIn

collect

接受 flow 下数据流:

flow {
    emit(1)
    emit(2)
}.collect {
    println(it)
}

single

如果 flow 中只发送一个数据可以通过 single 来返回,但是如果多于一个数据就会抛出 java.lang.IllegalArgumentException: Flow has more than one element

println("------flow terminal single")
val singleValue = flowOf(1)
    .single()
println("single value = $singleValue")

//  flowOf(1, 2)
//      .single()

single 的源码如下:

public suspend fun <T> Flow<T>.single(): T {
    var result: Any? = NULL
    collect { value ->
        require(result === NULL) { "Flow has more than one element" }
        result = value
    }

    if (result === NULL) throw NoSuchElementException("Flow is empty")
    return result as T
}

first

获取 Flow 中的第一个元素。

val firstValue = flowOf(1, 2, 3)
    .first()
println("first value = $firstValue")

first 的源码如下:

public suspend fun <T> Flow<T>.first(): T {
    var result: Any? = NULL
    collectWhile {
        result = it
        false
    }
    if (result === NULL) throw NoSuchElementException("Expected at least one element")
    return result as T
}

toList

toList 把 flow 中的数据添加到集合中。

val list = flowOf(1, 2, 3)
    .toList()
println("list = $list")
val myList = mutableListOf(100, 200)
val myList2 = flowOf(1, 2, 3)
    .toList(myList)
println("add myList as param res = $myList2")

toList 源码如下:

public suspend fun <T> Flow<T>.toList(destination: MutableList<T> = ArrayList()): List<T> = toCollection(destination)
public suspend fun <T, C : MutableCollection<in T>> Flow<T>.toCollection(destination: C): C {
    collect { value ->
        destination.add(value)
    }
    return destination
}

fold

用初始值作为初累计结果,然后再将累计的结果和每一个值做 operator 运算,然后将结果作为累计,用于下一次的计算。

Accumulate value starting with initial value and applying operation current accumulator value and each element.

val sum = flowOf(1, 2, 3)
    .fold(0) { acc, value ->
        acc + value
    }

println("use flow fold to sum, result = $sum")

fold 源码如下:

public suspend inline fun <T, R> Flow<T>.fold(
    initial: R,
    crossinline operation: suspend (acc: R, value: T) -> R
): R {
    var accumulator = initial
    collect { value ->
        accumulator = operation(accumulator, value)
    }
    return accumulator
}

reduce

reduce 和 fold 相似,不同的是 reduce 没有初始值。

val sum = flowOf(1, 2, 3)
    .reduce { accumulator, value ->
        accumulator + value
    }
println("use flow reduce to sum, result = $sum")

reduce 源码如下:

public suspend fun <S, T : S> Flow<T>.reduce(operation: suspend (accumulator: S, value: T) -> S): S {
    var accumulator: Any? = NULL

    collect { value ->
        accumulator = if (accumulator !== NULL) {
            @Suppress("UNCHECKED_CAST")
            operation(accumulator as S, value)
        } else {
            value
        }
    }

    if (accumulator === NULL) throw NoSuchElementException("Empty flow can't be reduced")
    @Suppress("UNCHECKED_CAST")
    return accumulator as S
}

launchIn

launchIn 会在执行的协程内调用 Flow 的 collect() 方法。

val simpleFlow: Flow<Int> = flow {
    for (i in 1..10) {
        emit(i)
        println("launchIn ${currentCoroutineContext()[CoroutineName]?.name} emit $i")
        delay(100)
    }
}
coroutineScope {
    launch(CoroutineName("SimpleCoroutine")) {
        simpleFlow.launchIn(this)
    }
}

launchIn 的源码如下:

public fun <T> Flow<T>.launchIn(scope: CoroutineScope): Job = scope.launch {
    collect() // tail-call
}

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 203,456评论 5 477
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 85,370评论 2 381
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 150,337评论 0 337
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 54,583评论 1 273
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 63,596评论 5 365
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 48,572评论 1 281
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 37,936评论 3 395
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 36,595评论 0 258
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 40,850评论 1 297
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 35,601评论 2 321
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 37,685评论 1 329
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 33,371评论 4 318
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 38,951评论 3 307
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 29,934评论 0 19
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 31,167评论 1 259
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 43,636评论 2 349
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 42,411评论 2 342