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

RxJava中组合Observer而非Observable的实现方案求助

解决方案:用RxJava的zip操作符合并两个UseCase的数据流

刚接触RxJava遇到这种场景很正常,我们可以通过把两个UseCase的执行逻辑包装成Observable,再用zip操作符实现“两者都拿到数据后再统一展示”的需求。下面是贴合你现有代码结构的完整示例:

// 1. 将mGetPotatoes包装为Observable流
Observable<List<Potatoes>> potatoesObservable = Observable.create(emitter -> {
    mGetPotatoes.execute(new DisposableObserver<List<Potatoes>>() {
        @Override
        public void onNext(List<Potatoes> potatoes) {
            emitter.onNext(potatoes);
        }

        @Override
        public void onComplete() {
            emitter.onComplete();
        }

        @Override
        public void onError(Throwable e) {
            emitter.onError(e);
        }
    });
});

// 2. 将mGetBurger包装为Observable流
Observable<Burger> burgerObservable = Observable.create(emitter -> {
    mGetBurger.execute(new DisposableObserver<Burger>() {
        @Override
        public void onNext(Burger burger) {
            emitter.onNext(burger);
        }

        @Override
        public void onComplete() {
            emitter.onComplete();
        }

        @Override
        public void onError(Throwable e) {
            emitter.onError(e);
        }
    });
});

// 3. 定义合并后的观察者,统一处理结果
DisposableObserver<Pair<List<Potatoes>, Burger>> combinedObserver = new DisposableObserver<>() {
    @Override
    public void onNext(Pair<List<Potatoes>, Burger> dataPair) {
        // 同时拿到土豆和汉堡数据,统一展示
        List<Potatoes> potatoes = dataPair.first;
        Burger burger = dataPair.second;
        
        getMvpView().showPotatoes(mPotatoesMapper.mapPotatoesToViewModels(potatoes));
        getMvpView().showBurger(mBurgerMapper.mapBurgerToViewModel(burger));
        getMvpView().hideProgress();
    }

    @Override
    public void onComplete() {}

    @Override
    public void onError(Throwable e) {
        // 任何一个请求出错都会走到这里,统一处理错误
        getMvpView().hideProgress();
        getMvpView().showErrorMessage(e.getMessage());
    }
};

// 4. 用zip合并两个Observable并订阅
Observable.zip(
        potatoesObservable,
        burgerObservable,
        Pair::new // 把两个数据打包成Pair传递给下游
).subscribe(combinedObserver);

// 记得在页面生命周期结束时取消订阅(比如onDestroy),避免内存泄漏
// compositeDisposable.add(combinedObserver);

关键细节说明

  • 包装UseCase为Observable:通过Observable.create把execute的回调逻辑转换成RxJava的标准流,这样就能使用各种操作符处理数据流。
  • zip操作符的核心作用:zip会等待所有输入的Observable都发射至少一次数据后,才会把对应位置的数据组合起来发送给下游,完美满足你“两者都触发onNext才展示”的需求。
  • 统一错误处理:原来的代码两个请求各自处理错误,现在合并后只要其中一个请求出错,都会触发统一的onError回调,减少重复代码。
  • 资源管理:建议使用CompositeDisposable管理所有订阅,在页面销毁时调用clear()或dispose(),防止内存泄漏。

简化优化(可选)

如果不想用Pair传递数据,也可以直接在zip的组合函数里完成展示逻辑:

Observable.zip(
        potatoesObservable,
        burgerObservable,
        (potatoes, burger) -> {
            // 直接在这里处理展示逻辑
            getMvpView().showPotatoes(mPotatoesMapper.mapPotatoesToViewModels(potatoes));
            getMvpView().showBurger(mBurgerMapper.mapBurgerToViewModel(burger));
            return null; // 不需要向下游传递数据,返回任意值即可
        }
).subscribe(new DisposableObserver<Object>() {
    @Override
    public void onNext(Object o) {
        getMvpView().hideProgress();
    }

    @Override
    public void onComplete() {}

    @Override
    public void onError(Throwable e) {
        getMvpView().hideProgress();
        getMvpView().showErrorMessage(e.getMessage());
    }
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:32:55