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
相关产品推荐
相关产品推荐

