You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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")
    }
}

代码说明

  1. 独立超时协程:用launch启动的协程不受awaitClose阻塞,能正常执行延迟逻辑。
  2. 取消超时任务:无论回调成功还是失败,都要取消超时协程,防止后续发送超时错误。
  3. Flow取消时清理:在awaitClose中取消超时任务,避免资源泄漏。

内容的提问来源于stack exchange,提问作者Metropol_Tilkisi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 20:05:09