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

RxJava2多方法组合疑问:本地+服务器数据拉取与网络校验

优化RxJava本地+远程数据获取与网络校验的实现方案

看了你的代码,核心问题在于重复的订阅逻辑、网络状态与数据请求的耦合不清晰,以及无网场景下错误处理缺失导致崩溃。我们可以通过RxJava的操作符整合逻辑,复用代码,同时让网络校验、本地/远程数据获取的流程更健壮。

核心优化思路

  • 封装通用逻辑:把存储数据、合并分类的代码抽成独立方法,避免重复冗余
  • 统一网络校验:先检查网络状态,再决定是否发起远程请求,同时保证本地数据优先展示
  • 合并数据流:用RxJava的concat合理组合本地和远程数据,避免重复订阅
  • 全局错误处理:统一捕获网络错误、服务器错误等异常,彻底避免无网崩溃

优化后的完整代码

private void getCategories() {
    mView.showConnectingProgress(); // 统一显示加载进度

    // 1. 先获取本地缓存数据并展示
    Observable<List<FilterCategory>> localDataObservable = getDataFromLocal(context)
            .map(this::processFilterResponse) // 复用统一处理逻辑
            .observeOn(AndroidSchedulers.mainThread())
            .doOnNext(categories -> {
                if (mView != null && categories != null && !categories.isEmpty()) {
                    mView.onCategoriesReceived(categories);
                }
            });

    // 2. 网络校验+远程数据请求流
    Observable<List<FilterCategory>> remoteDataObservable = InternetUtil.isConnectionAvailable()
            .subscribeOn(Schedulers.io())
            .flatMapObservable(isOnline -> {
                if (!isOnline) {
                    // 无网时抛出自定义异常,统一在错误回调处理
                    return Observable.error(new NoNetworkException());
                }
                // 有网时发起远程请求,复用处理逻辑
                return getDataFromServer(context)
                        .map(this::processFilterResponse);
            })
            .observeOn(AndroidSchedulers.mainThread());

    // 3. 合并本地+远程数据流,统一处理结果与错误
    Observable.concat(localDataObservable, remoteDataObservable)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(
                    categories -> {
                        // 远程数据返回后自动更新UI,覆盖本地缓存
                        if (mView != null && categories != null && !categories.isEmpty()) {
                            mView.onCategoriesReceived(categories);
                        }
                        mView.hideConnectingProgress();
                    },
                    throwable -> {
                        mView.hideConnectingProgress();
                        if (throwable instanceof NoNetworkException) {
                            mView.showOfflineMessage();
                        } else {
                            // 统一处理其他错误:服务器异常、解析错误等
                            String errorMsg = throwable.getMessage() != null ? throwable.getMessage() : "获取数据失败";
                            if (throwable instanceof HttpException) {
                                ResponseBody body = ((HttpException) throwable).response().errorBody();
                                if (body != null) {
                                    try {
                                        errorMsg = body.string();
                                    } catch (IOException e) {
                                        e.printStackTrace();
                                    }
                                }
                            }
                            mView.onCategoriesReceivingFailure(errorMsg);
                        }
                    }
            );
}

// 封装通用的响应处理逻辑:存储数据+合并重复分类
private List<FilterCategory> processFilterResponse(PromoFilterResponse response) {
    if (response == null) {
        return Collections.emptyList();
    }
    // 存储到本地SharedPreferences
    PreferencesHelper.putObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, response);
    // 执行分类合并逻辑
    try {
        return combineDuplicatedCategories(response).blockingFirst();
    } catch (Exception e) {
        return Collections.emptyList();
    }
}

// 自定义无网异常,方便统一识别处理
private static class NoNetworkException extends Exception {
    public NoNetworkException() {
        super("无网络连接");
    }
}

private Observable<PromoFilterResponse> getDataFromServer(Context context) {
    return RetrofitHelper.getApiService()
            .getFilterCategories(Constants.PROMO_FILTER_CATEGORIES_URL)
            .subscribeOn(Schedulers.io())
            .retryWhen(BaseDataManager.isAuthException())
            .doOnError(throwable -> LogUtil.e("ERROR", throwable.getMessage()));
}

private Observable<PromoFilterResponse> getDataFromLocal(Context context) {
    PromoFilterResponse response = PreferencesHelper.getObject(context, PreferencesKey.FILTER_CATEGORIES_KEY, PromoFilterResponse.class);
    return response != null ? Observable.just(response) : Observable.empty();
}

关键优化点说明

  • 复用逻辑减少冗余:processFilterResponse方法统一处理数据存储和分类合并,避免本地/远程请求中重复写相同代码
  • 优先展示本地缓存:用concat先订阅本地数据流,保证用户先看到缓存数据,再用远程数据更新,提升加载体验
  • 无网场景安全处理:通过自定义异常捕获无网状态,在全局错误回调中统一处理,彻底避免崩溃
  • 统一错误处理入口:所有异常(无网、服务器错误、解析错误)都在同一个onError回调中处理,逻辑更集中易维护
  • 简化订阅管理:不再分别添加多个订阅到CompositeDisposable,合并成一个数据流后更易管理生命周期

这样修改后,既解决了无网崩溃的问题,又让代码结构更清晰,逻辑更健壮,同时兼顾了用户体验。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 06:38:01