Kotlin Flow遇Collector异常时如何继续发射?异常处理诉求
问题:Flow发射时捕获Collector异常并继续发射后续值
当Collector中抛出异常时,emitAll会将异常抛回,尝试捕获后继续发射后续值时触发了Flow的禁止异常。需求是在库侧处理所有来自Collector的异常,实现异常后继续发射,且无法控制Collector的实现。
相关代码
库侧代码
// library code val items = listOf(1, 2, 3) val flow = flow<Int> { for (item in items) { val childFlow = listOf(item, item, item) .asFlow() .catch { println(it.message) } .onCompletion { error -> println("completed child flow, error=${error?.message}") } try { emitAll( childFlow ) } catch (ex: Exception) { if (ex.message?.contains("not allowed") == true) { println("transient exception when emitting ${ex.message}") continue } throw ex } } } .catch { println(it.message) } .onEach { }
客户端代码
// client code flow .onEach { item -> if (item == 2) { throw Exception("not allowed") } println(item) } .collect()
触发的异常
Previous 'emit' call has thrown exception java.lang.Exception: not allowed, but then emission attempt of value '3' has been detected. Emissions from 'catch' blocks are prohibited in order to avoid unspecified behaviour, 'Flow.catch' operator can be used instead. For a more detailed explanation, please refer to Flow documentation.
解决方案
核心原因
Flow的Collector一旦抛出异常,就会进入失败状态,此时上游再尝试发射值会触发安全检查,抛出上述异常。直接在try/catch块中继续发射的方式违反了Flow的状态约束。
解决方法
使用Flow的onErrorContinue操作符,它可以在捕获元素处理过程中的异常后,继续发射后续元素,同时不会让Collector进入失败状态。我们在库侧添加这个操作符,统一处理来自下游Collector的异常:
修改后的库侧代码
// library code val items = listOf(1, 2, 3) val flow = flow<Int> { for (item in items) { val childFlow = listOf(item, item, item) .asFlow() .catch { println(it.message) } .onCompletion { error -> println("completed child flow, error=${error?.message}") } emitAll(childFlow) } } // 添加onErrorContinue处理下游异常,继续发射后续元素 .onErrorContinue { ex, value -> if (ex.message?.contains("not allowed") == true) { println("transient exception when emitting ${ex.message} for value $value") } else { // 非预期异常重新抛出,避免吞掉错误 throw ex } } .catch { println(it.message) }
代码说明
onErrorContinue会捕获下游处理元素时抛出的异常,第一个参数是异常对象,第二个参数是触发异常的元素值。- 针对预期的"not allowed"异常,我们记录日志后直接跳过,流会继续发射后续元素。
- 针对非预期异常,重新抛出以避免隐藏错误,确保异常能被上层感知。
- 原有的
try/catch块可以移除,因为onErrorContinue已经统一处理了下游异常。
内容的提问来源于stack exchange,提问作者oguzh4n
相关产品推荐
相关产品推荐

