RxJava中.flatMapCompletable()后.andThen()始终未触发的问题
问题分析与解决方案
问题原因
核心矛盾在于**Observable.interval的无限序列特性与andThen的触发规则不兼容**:
Observable.interval会每隔5秒持续发射递增的Long值,是一个永远不会终止的Observable流。flatMapCompletable会将每个interval事件转换为Completable,但整个Completable流会一直处理新的interval事件,永远不会进入「完成」状态。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
相关产品推荐
相关产品推荐

