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

如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 01:20:28