Kotlin中RxJava3如何让Flowable仅首次调用后停止发射数据?
问题描述
我正在使用RxJava3,需要获取来自Android Paging的Flowable分页数据,同时首次需从另一个API获取初始数据。希望使用zip操作符仅首次同时调用两者,之后仅调用分页数据,以下是我的代码:
init { fetchData() } fun fetchData() { val nextChallenges = rxRepo.getData().cachedIn(viewModelScope) var classSubscriptions: Flowable<ServiceCallResult<BasePaginatedModel<MySubscriptionEntity>, MTFailure>> = gymRepository.getMySubscriptions(MySubscriptionsRequest()).toFlowable() Flowable.zip( nextChallenges, classSubscriptions, ) { pagingData1: PagingData<NextSubscriptionsEntity>, pagingData2: ServiceCallResult<BasePaginatedModel<MySubscriptionEntity>, MTFailure> -> Pair(pagingData1, pagingData2) }.subscribeOn(AndroidSchedulers.mainThread()).observeOn(Schedulers.io()).subscribeBy { nextSubscriptions.postValue(it.first) classSubscriptions = Flowable.empty() }.addTo(compositeDisposable) }
请问是否有办法在classSubscriptions成功获取一次后,停止两者调用,仅每次调用nextChallenges?或者是否应该避免使用zip操作符?
解决方案
你的当前写法存在逻辑问题:classSubscriptions是个变量,但Flowable.zip订阅的是初始化时的那个Flowable实例,后续在subscribeBy里把变量改成Flowable.empty()不会影响已经订阅的流。而且zip的特性是必须等待两个流都发射数据才会触发下游,分页流会持续发射新数据,这会导致第一次配对后,后续的分页数据会一直等待初始API流的下一个元素(永远不会来),完全无法触发下游回调。
绝对不建议用zip实现这个需求,因为zip的核心是“元素配对”,和你“首次并行、后续仅分页”的场景完全不匹配。推荐两种更合理的实现方式:
方式一:并行触发首次请求,后续仅接收分页数据
如果要首次同时触发两个API调用,之后只处理分页数据,可以把初始API流限制为仅发射一次,再和分页流合并:
fun fetchData() { val nextChallenges = rxRepo.getData().cachedIn(viewModelScope) // 初始订阅流仅调用一次API,发射一次结果后结束 val oneTimeSubscriptions = gymRepository.getMySubscriptions(MySubscriptionsRequest()) .toFlowable() .take(1) .share() // 确保API仅被调用一次,避免重复订阅触发多次请求 // 合并两个流:首次会同时收到分页数据和初始订阅数据,之后只有分页数据 Flowable.merge( nextChallenges.map { Pair(it, null) }, oneTimeSubscriptions.map { Pair(null, it) } ) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribeBy { (pagingData, subscriptionData) -> // 处理分页数据 pagingData?.let { nextSubscriptions.postValue(it) } // 处理仅一次的初始订阅数据 subscriptionData?.let { // 这里写你的初始数据逻辑,比如更新UI或存入本地 } } .addTo(compositeDisposable) }
merge会同时订阅两个流,初始API流因为take(1)在发射一次数据后就会终止,后续只会收到分页流的新数据,完美符合你的需求。
方式二:先获取初始数据,再启动分页流
如果允许先完成初始API调用,再开始接收分页数据,可以用flatMap串联两个流:
fun fetchData() { val nextChallenges = rxRepo.getData().cachedIn(viewModelScope) gymRepository.getMySubscriptions(MySubscriptionsRequest()) .toFlowable() .take(1) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .doOnNext { subscriptionResult -> // 先处理初始订阅数据 } .flatMap { nextChallenges } // 初始数据获取完成后,订阅分页流 .subscribeBy { pagingData -> nextSubscriptions.postValue(pagingData) } .addTo(compositeDisposable) }
这种方式逻辑更简单,适合对初始数据依赖较强的场景。
内容的提问来源于stack exchange,提问作者ayman omara
相关产品推荐
相关产品推荐

