【问题标题】:Combination of several methods in RXJava2RXJava2中几种方法的组合
【发布时间】:2019-07-30 13:02:41
【问题描述】:

事实上,我需要同时从本地数据库和服务器中提取数据,同时检查与 Internet 的连接。

不检查互联网很容易。但是当我关闭移动数据时,崩溃了。

我不明白如何组合并决定这样做:

private void getCategories() {

    composite.add(getDataFromLocal(context)
            .observeOn(AndroidSchedulers.mainThread()).flatMap(new Function<PromoFilterResponse, ObservableSource<List<FilterCategory>>>() {
                @Override
                public ObservableSource<List<FilterCategory>> apply(PromoFilterResponse promoFilterResponse) throws Exception {
                    if (promoFilterResponse != null) {
                        PreferencesHelper.putObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, promoFilterResponse);
                        return combineDuplicatedCategories(promoFilterResponse);
                    } else {
                        return Observable.empty();
                    }
                }
            })
            .subscribe(new Consumer<List<FilterCategory>>() {
                @Override
                public void accept(List<FilterCategory> categories) throws Exception {
                    if (mView != null) {
                        mView.hideConnectingProgress();
                        if (categories != null && categories.size() > 0) {
                            mView.onCategoriesReceived(categories);
                        }
                    }
                }
            }));

    composite.add(InternetUtil.isConnectionAvailable().subscribe(isOnline -> {
        if (isOnline) {
            composite.add(
                    getDataFromServer(context)
                            .flatMap(new Function<PromoFilterResponse, ObservableSource<List<FilterCategory>>>() {
                                @Override
                                public ObservableSource<List<FilterCategory>> apply(PromoFilterResponse promoFilterResponse) throws Exception {
                                    if (promoFilterResponse != null) {
                                        PreferencesHelper.putObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, promoFilterResponse);
                                        return combineDuplicatedCategories(promoFilterResponse);
                                    } else {
                                        return Observable.empty();
                                    }
                                }
                            })
                            .observeOn(AndroidSchedulers.mainThread())
                            .subscribe(categories -> {
                                if (mView != null) {
                                    mView.hideConnectingProgress();
                                    if (categories != null && categories.size() > 0) {
                                        mView.onCategoriesReceived(categories);
                                    } else {
                                        mView.onCategoriesReceivingFailure(errorMessage[0]);
                                    }
                                }
                            }, throwable -> {
                                if (mView != null) {
                                    if (throwable instanceof HttpException) {
                                        ResponseBody body = ((HttpException) throwable).response().errorBody();

                                        if (body != null) {
                                            errorMessage[0] = body.string();
                                        }
                                    }
                                    mView.hideConnectingProgress();
                                    mView.onCategoriesReceivingFailure(errorMessage[0]);
                                }
                            }));
        } else {
            mView.hideConnectingProgress();
            mView.showOfflineMessage();
        }
    }));
} 


private Single<Boolean> checkNetwork(Context context) {
    return InternetUtil.isConnectionAvailable()
            .subscribeOn(Schedulers.io())
            .doOnSuccess(new Consumer<Boolean>() {
                @Override
                public void accept(Boolean aBoolean) throws Exception {
                    getDataFromServer(context);
                }
            });
}

private Observable<PromoFilterResponse> getDataFromServer(Context context) {
    return RetrofitHelper.getApiService()
            .getFilterCategories(Constants.PROMO_FILTER_CATEGORIES_URL)
            .subscribeOn(Schedulers.io())
            .retryWhen(BaseDataManager.isAuthException())
            .publish(networkResponse ->  Observable.merge(networkResponse,  getDataFromLocal(context).takeUntil(networkResponse)))
            .doOnNext(new Consumer<PromoFilterResponse>() {
                @Override
                public void accept(PromoFilterResponse promoFilterResponse) throws Exception {
                    PreferencesHelper.putObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, promoFilterResponse);
                }
            })
            .doOnError(new Consumer<Throwable>() {
                @Override
                public void accept(Throwable throwable) throws Exception {
                    LogUtil.e("ERROR", throwable.getMessage());
                }
            });

}

private Observable<PromoFilterResponse> getDataFromLocal(Context context) {
    PromoFilterResponse response = PreferencesHelper.getObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, PromoFilterResponse.class);
    if (response != null) {
        return Observable.just(response)
                .subscribeOn(Schedulers.io());
    } else {
        return Observable.empty();
    }
}

如您所见,分别连接本地数据库,同时上网并从服务器上传数据。

但在我看来不太对劲。此外,订阅者是重复的等等。

看了很多教程,里面有描述本地数据库和API的结合,但是我没有看到同时处理与互联网的连接错误。

我想很多人都遇到过这样的问题,你是怎么解决的?

【问题讨论】:

    标签: java android rx-java2 rx-android


    【解决方案1】:

    假设你有两个 Obsevable:一个来自服务器,另一个来自数据库

    您可以将它们合并到一个流中,如下所示:

      public Observable<Joke> getAllJokes() {
    
        Observable<Joke> remote = mRepository.getAllJokes()
                .subscribeOn(Schedulers.io());
    
    
        Observable<Joke> local = mRepository.getAllJokes().subscribeOn(Schedulers.io());
    
          return Observable.mergeDelayError(local, remote).filter(joke -> joke != null);
    }
    

    【讨论】:

    • 是的,我做到了,一切正常,但我仍然有 Observable 可以检查互联网,我应该处理他的答案。我有 3 个 observables,从服务器、本地数据库获取数据并检查与 Internet 的连接
    【解决方案2】:

    我不是安卓开发者,但在我看来方法返回类型应该是这样的:

    //just for demonstration
    static boolean isOnline = false;
    
    static class NoInternet extends RuntimeException {
    }
    
    private static Completable ensureOnline() {
        if (isOnline)
            return Completable.complete();
        else
            return Completable.error(new NoInternet());
    
    }
    
    private static Single<String> getDataFromServer() {
        return Single.just("From server");
    }
    
    private static Maybe<String> getDataFromLocal() {
        return Maybe.just("From local");//or Maybe.never()
    }
    

    我们可以与Observable.merge 并行运行。但是如果发生错误NoIternet 怎么办?合并的 observable 将失败。我们可以使用materialisation - 将所有发射和错误转换为onNext 值。

    private static void loadData() {
    
        Observable<Notification<String>> fromServer = ensureOnline().andThen(getDataFromServer()).toObservable().materialize();
    
        Observable<Notification<String>> fromLocaldb = getDataFromLocal().toObservable().materialize();
    
        Observable.merge(fromLocaldb, fromServer)
                .subscribe(notification -> {
                    if (notification.isOnNext()) {
                        //calls one or two times(db+server || db || server)
                        //show data in ui
                    } else if (notification.isOnError()) {
                        if (notification.getError() instanceof NoInternet) {
                            //show no internet
                        } else {
                            //show another error
                        }
                    } else if (notification.isOnComplete()){
                        //hide progress bar
                    }
    
    
    
                });
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-11-14
      • 2017-06-27
      • 1970-01-01
      • 2014-09-28
      • 1970-01-01
      • 1970-01-01
      • 2018-05-28
      • 2020-10-19
      相关资源
      最近更新 更多