如何将Room DAO Flow与手动网络请求合并为单一Flow
合并缓存Flow与网络请求结果的解决方案
问题背景
我现有仓库包含两个方法:
// 供ViewModel监听,但无法获取网络请求失败信息 fun getCachedItems(): Flow<Result<List<Item>>> { // 返回Result.cacheUpdate包装的Room DAO Flow return itemDao.items().map { Result.cacheUpdate(it) } } // 用户手动刷新触发,非Flow类型 suspend fun getLatestItemsFromNetwork(): Result<List<Item>> { val result = remoteSource.getItems() // 返回Result.success或Result.error if (result is Result.Success) { updateDatabaseCache(result.data) } return result }
需要实现一个getItemUpdates(): Flow<Result<List<Item>>>方法,要求同时输出三种Result类型:
Result.cacheUpdate:缓存更新(包括用户编辑标题触发的更新)Result.success:网络请求成功Result.error:网络请求失败
核心痛点:不能仅靠缓存更新触发Flow(网络失败时无缓存更新,会丢失错误信息),必须把网络请求的成功/失败状态也通过Flow传递给UI。
使用场景
- ViewModel启动时先加载缓存数据
- 自动发起网络请求,结果需通过Flow返回成功/失败状态
- 用户可手动触发刷新,结果同样通过Flow返回
- 用户编辑标题触发的缓存更新要纳入Flow输出
实现思路
用SharedFlow专门传递网络请求的结果(success/error),再将它与缓存的Flow合并,同时处理初始缓存的加载逻辑,覆盖所有场景。
完整代码实现
1. 仓库层修改
class ItemRepository( private val itemDao: ItemDao, private val remoteSource: RemoteSource ) { // 用于发送网络请求结果的SharedFlow private val networkResultFlow = MutableSharedFlow<Result<List<Item>>>(replay = 0) // 原缓存Flow,返回cacheUpdate类型 private fun getCachedItems(): Flow<Result<List<Item>>> { return itemDao.items().map { Result.CacheUpdate(it) } } // 手动刷新方法,执行后将结果发送到networkResultFlow suspend fun getLatestItemsFromNetwork(): Result<List<Item>> { val result = remoteSource.getItems() if (result is Result.Success) { updateDatabaseCache(result.data) } // 发送网络结果到Flow networkResultFlow.emit(result) return result } // 核心方法:合并缓存Flow与网络结果Flow fun getItemUpdates(): Flow<Result<List<Item>>> { return flow { // 第一步:先发送初始缓存数据 val initialCache = itemDao.getItemsNonFlow() // 假设你有这个非Flow的缓存查询方法 emit(Result.CacheUpdate(initialCache)) // 第二步:合并缓存更新Flow和网络结果Flow emitAll( merge( getCachedItems(), networkResultFlow ) ) } } // 自动发起网络请求的方法(可在ViewModel初始化时调用) suspend fun fetchInitialNetworkData() { getLatestItemsFromNetwork() } // 其他辅助方法(比如更新数据库) private suspend fun updateDatabaseCache(items: List<Item>) { itemDao.insertItems(items) } // 用户编辑标题触发数据库更新的方法 suspend fun updateItemTitle(itemId: String, newTitle: String) { itemDao.updateTitleById(itemId, newTitle) } } // 假设Result密封类定义如下 sealed class Result<out T> { data class Success<out T>(val data: T) : Result<T>() data class Error(val message: String) : Result<Nothing>() data class CacheUpdate<out T>(val data: T) : Result<T>() }
2. ViewModel层使用示例
class ItemViewModel(private val repository: ItemRepository) : ViewModel() { val itemUpdates = repository.getItemUpdates() init { // ViewModel启动时自动发起网络请求 viewModelScope.launch { repository.fetchInitialNetworkData() } } // 用户手动刷新 fun onRefresh() { viewModelScope.launch { repository.getLatestItemsFromNetwork() } } // 用户编辑标题后更新数据库(触发缓存Flow更新) fun updateItemTitle(itemId: String, newTitle: String) { viewModelScope.launch { repository.updateItemTitle(itemId, newTitle) } } }
方案说明
- 初始缓存加载:在
getItemUpdates()的flow块中先发送一次初始缓存,确保ViewModel启动时UI能立刻拿到数据。 - 网络结果传递:通过
MutableSharedFlow将网络请求的success/error状态发送出去,合并到主Flow中,保证网络失败时UI能收到错误信息。 - 缓存更新监听:原有的缓存Flow会在数据库更新(包括用户编辑标题、网络成功更新缓存)时自动发送
CacheUpdate事件,无需额外处理。 - 手动刷新:调用
getLatestItemsFromNetwork()时,结果会自动发送到networkResultFlow,主Flow会收到对应的success/error事件。
这样就能完美覆盖所有需求场景,同时保证三种Result类型都能通过getItemUpdates()传递给UI。
内容的提问来源于stack exchange,提问作者me.at.coding
相关产品推荐
相关产品推荐

