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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 05:14:58