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
相关产品推荐
相关产品推荐

