Flow与produce().receiveAsFlow()输出差异及报错原因咨询
我写了一段代码,想要创建两个Flow,当它们等待的消息出现在输入Flow中时就执行完成。
regularFlow测试的输出符合预期:
1 REC message2 2 REC message2 2 Done 1 REC message1 1 Done
但receiveChannel测试的输出却报错了:
1 REC message2 2 REC message1 Expected at least one element matching the predicate Function2<java.lang.String, kotlin.coroutines.Continuation<? super java.lang.Boolean>, java.lang.Object>
测试代码
@Test fun receiveChannel() { runBlocking { doIt { produce<String> { delay(500) send("message2") delay(1000) send("message1") }.receiveAsFlow() } } } @Test fun regularFlow() { runBlocking { doIt { flow { delay(500) emit("message2") delay(1000) emit("message1") } } } } suspend fun doIt(buildFlow: suspend () -> Flow<String>) { coroutineScope { val commands = buildFlow() val filter1 = async { commands.onEach { println("1 REC " + it) }.first { it == "message1" } println("1 Done") } val filter2 = async { commands.onEach { println("2 REC " + it) }.first { it == "message2" } println("2 Done") } filter1.await() filter2.await() } }
原因解释
核心差异在于普通Flow是冷流,而produce创建的Channel转成的Flow是热流,两者的订阅逻辑完全不同:
普通Flow(
flow { ... })是冷流:
每次调用first这类终端操作时,都会重新执行Flow的构建逻辑,每个订阅者都拥有独立的数据流副本。在doIt函数里,filter1和filter2两个异步任务各自订阅commands,相当于各自触发了一次flow { ... }的执行,所以都能拿到完整的message2和message1,最终都能匹配到目标消息完成执行。produce转Flow(receiveAsFlow())是热流:produce创建的是Channel,Channel属于热流组件——数据一旦被发送,只会被第一个订阅它的消费者拿走,后续订阅者只能获取之后发送的数据。这里filter1先启动订阅,拿走了message2但还在等待message1;之后filter2启动订阅时,message2已经被消费掉了,它只能拿到后续的message1,但filter2需要的是message2,因此永远匹配不到目标,最终抛出“没有符合条件的元素”的异常。
从输出也能验证这点:filter1收到了message2,filter2收到了message1,说明两者共享了同一个Channel的数据流,前面的消息被先订阅的任务消费了。
内容的提问来源于stack exchange,提问作者Jozef Dransfield

