Android Firestore callbackFlow 等待内部Flow收集完成再返回结果
问题根因
你的代码时序错误来自两个核心问题:
- 初始实现中你手动创建了独立的
CoroutineScope(Dispatchers.IO)启动协程拉取tripDays,这个协程生命周期不受外层flow管控,和外层发送Response.Success(trips)的逻辑完全并行,根本不会等待tripDays赋值完成就提前返回了空数据的列表。 - 你第二次改造的版本已经接近正确,但
getTripDays的callbackFlow在发送完Success结果后没有主动调用close(),导致collect操作会一直挂起,同时你把trips列表定义在launch外面,还是存在状态不一致的风险。 - 额外的坏实践:callbackFlow内部自带绑定生命周期的协程作用域,不需要手动新建CoroutineScope,手动创建的scope不会随flow取消而取消,容易造成内存泄漏。
最小改动修复方案
不需要重构整体结构,只需要补全流关闭逻辑、调整协程作用域、等待所有子任务完成后再发送成功结果即可。
首先修复getTripDays,在数据发送完成/出错后主动关闭流,让collect操作能正常结束:
override fun getTripDays(idTrip: String): Flow<Response<List<TripDay>>> = callbackFlow { val tripDaysRef = tripsRef.document(idTrip).collection(TRIP_DAYS_REF) val tripDays = mutableListOf<TripDay>() tripDaysRef .get() .addOnSuccessListener { for (doc in it.documents) { val tripDay = doc.toObject(TripDay::class.java) if (tripDay != null) { tripDay.id = doc.id tripDays.add(tripDay) } } trySend(Response.Success(tripDays)) close() // 数据发送完成关闭流,结束collect } .addOnFailureListener { trySend(Response.Error(it.message ?: it.toString())) close() // 请求出错也关闭流 } awaitClose { } }
然后修复getTrips,使用callbackFlow内置协程作用域,等所有trip的tripDays赋值完成后再发送最终成功响应:
override fun getTrips(idUser: String): Flow<Response<List<Trip>>> = callbackFlow { tripsRef .whereEqualTo(ID_USER, idUser) .get() .addOnSuccessListener { snapshot -> // 直接使用callbackFlow内置的协程作用域,生命周期和flow绑定 launch { val trips = mutableListOf<Trip>() trySend(Response.Loading) for (doc in snapshot.documents) { if (doc.id == "kJJFEatove5mxw4P70uq") continue val trip = doc.toObject(Trip::class.java) ?: continue trip.id = doc.id // 串行收集每个trip的days数据,因为getTripDays会主动close,这里collect会正常结束 tripDaysRepository.getTripDays(trip.id).collect { resp -> when (resp) { is Response.Loading -> { /* 外层已发送全局Loading,内层加载状态不需要透传 */ } is Response.Success -> trip.tripDays = resp.data is Response.Error -> { trySend(Response.Error(resp.message)) close() return@collect } is Response.Message -> trySend(Response.Message(resp.message)) } } trips.add(trip) } // 所有trip数据组装完成,才发送最终成功结果 trySend(Response.Success(trips)) close() } } .addOnFailureListener { trySend(Response.Error(it.message ?: it.toString())) close() } awaitClose { } }
可选优化(不改动原有架构)
如果需要提升拉取速度,可以把串行拉取tripDays改成并行,利用协程的async/awaitAll等待所有并行任务完成后再组装结果,改动量很小:
// 替换上面getTrips里launch块的逻辑即可 launch { trySend(Response.Loading) val tripDeferred = snapshot.documents .filter { it.id != "kJJFEatove5mxw4P70uq" } .mapNotNull { doc -> val trip = doc.toObject(Trip::class.java) ?: return@mapNotNull null trip.id = doc.id // 每个trip的days拉取作为并行异步任务 async(Dispatchers.IO) { var days = emptyList<TripDay>() tripDaysRepository.getTripDays(trip.id).collect { resp -> when(resp) { is Response.Success -> days = resp.data is Response.Error -> throw RuntimeException(resp.message) else -> {} } } trip.apply { tripDays = days } } } // 等待所有并行拉取任务完成 val finalTrips = tripDeferred.awaitAll() trySend(Response.Success(finalTrips)) close() }
内容的提问来源于stack exchange,提问作者Smooth
相关产品推荐
相关产品推荐

