Flowable.concatArray失效:Rx+Room实现Repository模式遇问题
嘿,我一眼就看到了问题所在——这是Room的Flowable特性和RxJava concatArray 操作符配合时的经典坑!
核心原因
Room中用@Query返回的Flowable<List<Category>>是一个持续订阅的热流:它不会自动终止,只要数据库里的数据发生变化,就会持续发射最新的列表。而concatArray的工作逻辑是:依次订阅每个流,只有当前一个流完全终止(调用onComplete)后,才会订阅下一个流。
你的本地流永远不会终止,所以concatArray会一直卡在第一个本地流上,永远不会执行后面的远程请求代码——这就是为什么你只看到本地日志,远程请求完全没动静的原因!
解决方案
根据你的Repository模式需求(先显示本地缓存,再后台请求远程更新缓存),我给你两种合适的处理方案:
方案1:修正concatArray的使用(适合需要明确分两次发射数据的场景)
把本地流改成只取一次数据就终止,这样concatArray就能顺利执行到远程流。只需要给本地流加上take(1)操作符:
override fun getCategories(): Flowable<List<Category>> { return Flowable.concatArray( // 只取本地一次数据,取完后终止流 categoryDao.getCategories() .take(1) .doOnNext { Timber.w("Retrieving Categories from Local DB!") }, getCategoriesFromRemote() .doOnNext { Timber.w("Retrieving Categories from API!") } ).subscribeOn(scheduleProvider.io()) } private fun getCategoriesFromRemote(): Flowable<List<Category>> { // 去掉多余的fromFlowable,直接把Single转成Flowable return categoryApi.getCategories() .toFlowable() .flatMap { holders -> Flowable.just(holders.firstOrNull()?.categories ?: emptyList()) } .doOnNext { saveCategoriesInDb(it) } } // 确保数据库操作在IO线程执行 private fun saveCategoriesInDb(categories: List<Category>) { Completable.fromAction { categoryDao.saveAllCategories(categories) }.subscribeOn(scheduleProvider.io()) .subscribe() .addTo(disposables) // 记得管理订阅,避免内存泄漏 }
这样执行后,ViewModel会先收到本地的空列表,然后收到远程获取的列表;同时远程数据保存到数据库后,如果你需要持续监听数据库变化,可以去掉take(1),改用其他方式组合流。
方案2:更贴合Repository模式的实现(推荐)
这种方式是先返回本地缓存的持续监听流,同时后台发起远程请求更新缓存——当远程数据保存到数据库后,本地的Flowable会自动发射最新数据,ViewModel就能自动收到更新,完全符合“先显示缓存,再后台更新”的设计初衷:
override fun getCategories(): Flowable<List<Category>> { // 后台发起远程请求,更新本地缓存(不需要订阅它的发射数据,只需要执行更新逻辑) getCategoriesFromRemote() .subscribeOn(scheduleProvider.io()) .subscribe( {}, { Timber.e("Failed to fetch remote categories: ${it.message}") } ) .addTo(disposables) // 持续监听本地数据库的变化,任何更新都会自动发射给ViewModel return categoryDao.getCategories() .doOnNext { Timber.w("Retrieving Categories from Local DB!") } .subscribeOn(scheduleProvider.io()) } private fun getCategoriesFromRemote(): Completable { // 改用Completable,因为我们只关心是否更新成功,不需要返回数据 return categoryApi.getCategories() .flatMapCompletable { holders -> val categories = holders.firstOrNull()?.categories ?: emptyList() Completable.fromAction { categoryDao.saveAllCategories(categories) } } }
这种实现更简洁,也完全符合你想要的Repository模式:ViewModel订阅后立刻拿到本地缓存(空列表),远程请求完成后缓存更新,ViewModel会自动收到最新的分类列表。
额外注意点
- 确保所有数据库操作和网络请求都在IO线程执行,避免阻塞主线程
- 记得用
disposables管理所有RxJava订阅,避免内存泄漏
内容的提问来源于stack exchange,提问作者M Rajoy

