如何在callbackFlow中添加延迟实现Firebase读取超时处理
Firebase回调转Flow的超时实现问题
我将Firebase的回调机制转换为Flow实现,以下是具体情况:
- OK代码:未添加延迟且保留
channel.close()时,代码可正常运行。 - FAIL代码:注释掉回调中的
channel.close()并添加10秒延迟后,程序无法执行到delay行,输出仅为:
Read-Before Await Read-Success
我的需求是:如果10秒内回调函数未被调用,就发送Resource.Error(Exception("Timeout"))并关闭channel。
OK代码
override fun read(id: String): Flow<Resource<DataSnapshot>> = callbackFlow { database.child(id).get().addOnSuccessListener { Timber.d("Read-Success") trySend(Resource.Success(it)) channel.close() }.addOnFailureListener { Timber.d("Read-Failure") trySend(Resource.Error(it)) channel.close() } Timber.d("Read-Before Await") awaitClose { Timber.d("Read-In Await") } Timber.d("Read-after Await") }
FAIL代码
override fun read(id: String): Flow<Resource<DataSnapshot>> = callbackFlow { database.child(id).get().addOnSuccessListener { Timber.d("Read-Success") trySend(Resource.Success(it)) //channel.close() }.addOnFailureListener { Timber.d("Read-Failure") trySend(Resource.Error(it)) //channel.close() } Timber.d("Read-Before Await") delay(10000) trySend(Resource.Error(Exception("timeout"))) channel.close() awaitClose { Timber.d("Read-In Await") } Timber.d("Read-after Await") }
问题原因
callbackFlow中的awaitClose是挂起函数,会暂停当前协程直到channel被关闭。在FAIL代码中,你注释掉了回调里的channel.close(),导致awaitClose一直处于挂起状态,后面的delay代码根本没有执行机会。
解决方案
通过launch启动独立协程处理超时逻辑,让它和回调逻辑并行执行;同时在回调触发时取消超时协程,避免重复发送错误事件:
override fun read(id: String): Flow<Resource<DataSnapshot>> = callbackFlow { // 启动超时任务协程 val timeoutJob = launch { delay(10000) trySend(Resource.Error(Exception("Timeout"))) channel.close() } database.child(id).get() .addOnSuccessListener { timeoutJob.cancel() // 回调成功,取消超时任务 Timber.d("Read-Success") trySend(Resource.Success(it)) channel.close() } .addOnFailureListener { timeoutJob.cancel() // 回调失败,取消超时任务 Timber.d("Read-Failure") trySend(Resource.Error(it)) channel.close() } Timber.d("Read-Before Await") awaitClose { timeoutJob.cancel() // Flow被取消时,清理超时任务 Timber.d("Read-In Await") } }
代码说明
- 独立超时协程:用
launch启动的协程不受awaitClose阻塞,能正常执行延迟逻辑。 - 取消超时任务:无论回调成功还是失败,都要取消超时协程,防止后续发送超时错误。
- Flow取消时清理:在
awaitClose中取消超时任务,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Metropol_Tilkisi
相关产品推荐
相关产品推荐

