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

Room与RxJava2(Flowable、Maybe)结合Repository模式使用问题求助

解决Room+Retrofit+RxJava2的本地缓存优先逻辑问题

Hey there! 作为刚入门RxJava的开发者,遇到这种「本地缓存优先,无缓存再请求网络并存库」的经典场景太正常了,我刚学RxJava的时候也在Maybe/Flowable的搭配上卡了好久,下面就给你拆解具体实现思路和代码示例:

核心逻辑梳理

我们要实现的流程是:

  • 先从Room数据库查询目标数据
  • 如果数据库有数据,直接把数据发射给下游(比如UI层)
  • 如果数据库没有数据,发起Retrofit网络请求
  • 拿到网络响应后,把数据存入Room数据库
  • 最后把数据(优先用数据库缓存,保证一致性)发射给下游

第一步:配置Room Dao的返回类型

根据你的需求选择Maybe或Flowable:

用Maybe实现单次查询(无数据触发onComplete)

适合只需要单次获取数据,不需要监听数据库变化的场景:

@Dao
interface YourEntityDao {
    // 查询数据,无数据时会触发Maybe的onComplete
    @Query("SELECT * FROM your_entity WHERE id = :entityId")
    fun getEntityById(entityId: String): Maybe<YourEntity>

    // 插入/更新数据,只关心成功失败,用Completable
    @Insert(onConflict = OnConflictStrategy.REPLACE)
    fun insertEntity(entity: YourEntity): Completable
}

用Flowable实现实时监听(数据库更新自动发射数据)

如果需要UI实时感知数据库变化(比如存库后自动刷新UI),就用Flowable:

@Dao
interface YourEntityDao {
    // 会持续发射数据,包括数据库插入/更新后
    @Query("SELECT * FROM your_entity WHERE id = :entityId")
    fun getEntityById(entityId: String): Flowable<YourEntity>

    @Insert(onConflict = OnConflictStrategy.REPLACE)
    fun insertEntity(entity: YourEntity): Completable
}

第二步:配置Retrofit接口

网络请求一般是单次操作,用Single返回(确保一定会拿到一个结果或错误):

interface YourApiService {
    @GET("api/entities/{id}")
    fun getEntityById(@Path("id") entityId: String): Single<YourEntity>
}

第三步:在Repository层组合流(核心部分)

这里用Dagger注入Dao和Api实例,然后用RxJava操作符串起整个逻辑:

基于Maybe的实现(单次查询)

class YourRepository @Inject constructor(
    private val entityDao: YourEntityDao,
    private val apiService: YourApiService
) {
    fun getEntity(entityId: String): Maybe<YourEntity> {
        return entityDao.getEntityById(entityId)
            // 当数据库Maybe触发onComplete(无数据)时,切换到网络请求流
            .switchIfEmpty(
                apiService.getEntityById(entityId)
                    // 拿到网络数据后,存入数据库(Completable只关心结果)
                    .flatMapCompletable { entity -> entityDao.insertEntity(entity) }
                    // 存库完成后,再次查询数据库,保证下游拿到的是本地缓存数据
                    .andThen(entityDao.getEntityById(entityId))
            )
            // 数据库和网络操作都放在IO线程
            .subscribeOn(Schedulers.io())
            // 切换回主线程给UI层用
            .observeOn(AndroidSchedulers.mainThread())
    }
}

基于Flowable的实现(实时监听)

如果用Flowable,我们可以利用它自动监听数据库变化的特性,网络请求存库后,数据库的Flowable会自动发射新数据:

class YourRepository @Inject constructor(
    private val entityDao: YourEntityDao,
    private val apiService: YourApiService
) {
    fun getEntity(entityId: String): Flowable<YourEntity> {
        // 数据库流:持续监听数据变化
        val databaseFlow = entityDao.getEntityById(entityId)
        // 网络流:请求并存库,完成后不发射数据(交给数据库流发射更新后的数据)
        val networkFlow = apiService.getEntityById(entityId)
            .flatMapCompletable { entity -> entityDao.insertEntity(entity) }
            .andThen(Flowable.empty())

        // 合并两个流,先发射数据库的数据,网络请求完成后数据库流会自动发射新数据
        return databaseFlow
            .mergeWith(networkFlow)
            // 如果只需要获取一次数据,就加firstElement()转Maybe再转回Flowable;如果要持续监听就去掉
            .firstElement()
            .toFlowable()
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
    }
}

给RxJava入门者的小提示

  • 区分RxJava类型:Maybe(可能有数据/无数据的单次操作)、Single(必返回一个结果的单次操作)、Completable(只关心成功失败的操作)、Flowable(持续发射数据+背压支持),别搞混了!
  • 线程调度别忘:数据库和网络操作一定要放在Schedulers.io(),UI更新必须切回AndroidSchedulers.mainThread(),不然会崩溃或者卡顿。
  • 管理订阅避免泄漏:在Activity/Fragment里用CompositeDisposable管理所有订阅,在onDestroy()里调用dispose(),防止内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:23:07