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

RxJava中.flatMapCompletable()后.andThen()始终未触发的问题

问题分析与解决方案

问题原因

核心矛盾在于**Observable.interval的无限序列特性与andThen的触发规则不兼容**:

  1. Observable.interval会每隔5秒持续发射递增的Long值,是一个永远不会终止的Observable流。
  2. flatMapCompletable会将每个interval事件转换为Completable,但整个Completable流会一直处理新的interval事件,永远不会进入「完成」状态。
  3. andThen的触发条件是:只有当前面的Completable完全完成后,才会执行后续的流逻辑。由于前面的Completable永远不会完成,所以getFavoriteCoins()永远不会被调用。

解决方案

根据需求场景,提供两种修复方案:

场景1:每次interval触发时,都先执行启动检查逻辑,再获取币种价格

将flatMapCompletable替换为flatMap,把启动逻辑与后续的获取收藏币种逻辑合并到每个interval事件的处理链中:

public void getCoinPrices() {
    disposable = Observable
            .interval(5, TimeUnit.SECONDS)
            .flatMap(n -> {
                Timber.d("Called flatmap: " + n);
                boolean isFirstTime = sharedPrefManager.isFirstTimeOpeningApp();
                Completable firstTimeLogic;
                
                if (isFirstTime) {
                    Timber.d("Is first time.");
                    firstTimeLogic = insertFavoritesUseCase.insertStartingCoins()
                            .andThen(Completable.fromAction(() -> {
                                sharedPrefManager.setFirstTimeOpeningApp(false);
                            }));
                } else {
                    Timber.d("Not first time.");
                    firstTimeLogic = Completable.complete();
                }
                
                // 执行完启动逻辑后,切换到获取收藏币种的Single并转成Observable
                return firstTimeLogic
                        .andThen(getFavoritesUseCase.getFavoriteCoins())
                        .toObservable();
            })
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .flatMap(listFavorites -> {
                Timber.d("List favorites called...");
                return getCoinsUseCase.getCoinPrices(listFavorites);
            })
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(listCoinsWithPrice -> {
                Timber.d("Called.......");
            }, err -> {
                Timber.e("Failed to get price of coins: " + err);
            });
}

场景2:首次启动逻辑只执行一次,之后每隔5秒直接获取币种价格

将首次启动逻辑从interval循环中剥离,先执行一次启动逻辑,再启动interval循环,避免重复检查:

public void getCoinPrices() {
    // 首次启动逻辑:仅执行一次
    Completable firstTimeSetup = Completable.defer(() -> {
        if (sharedPrefManager.isFirstTimeOpeningApp()) {
            Timber.d("Is first time.");
            return insertFavoritesUseCase.insertStartingCoins()
                    .andThen(Completable.fromAction(() -> {
                        sharedPrefManager.setFirstTimeOpeningApp(false);
                    }));
        } else {
            Timber.d("Not first time.");
            return Completable.complete();
        }
    });

    disposable = firstTimeSetup
            // 首次逻辑完成后,启动无限循环的interval
            .andThen(Observable.interval(5, TimeUnit.SECONDS))
            .flatMap(n -> {
                Timber.d("Interval triggered: " + n);
                return getFavoritesUseCase.getFavoriteCoins().toObservable();
            })
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .flatMap(listFavorites -> {
                Timber.d("List favorites called...");
                return getCoinsUseCase.getCoinPrices(listFavorites);
            })
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(listCoinsWithPrice -> {
                Timber.d("Called.......");
            }, err -> {
                Timber.e("Failed to get price of coins: " + err);
            });
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 20:46:34