RxAndroid Observable.zip场景下doOnComplete()偶发不执行问题求助
问题:Observable.zip 结合 doOnComplete 进度更新异常,部分 doOnComplete 未执行
我来帮你分析这个问题:你用 Observable.zip 同时发起多个 Retrofit API 请求,想通过每个请求的 doOnComplete 来更新进度条,但发现部分请求的 doOnComplete 回调没触发,虽然所有 API 都成功返回了数据。
问题代码与日志
你的核心实现是给每个请求 Observable 添加 doOnComplete 来累加进度,同时用 delaySubscription 延迟部分请求,最后通过 zip 组合结果:
private void preFetchData() { ApiInterface apiService1 = ApiClient.getWooRxClient().create(ApiInterface.class); ApiInterface apiService2 = ApiClient.getRxClient().create(ApiInterface.class); Map<String, String> map1 = new HashMap<>(); map1.put("on_sale", "true"); Observable<List<Product>> call1 = apiService1.getProducts1(map1) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .doOnComplete(() -> { progress += 18; Log.e("Progress1", progress + ""); mProgressBar.setProgress(progress); mProgressBar.setProgressText(progress + "%"); }); // 其他 call2-call6 逻辑类似,省略... Observable<CombinedHomePage> combined = Observable.zip(call1, call2, call3, call4, call5, call6, CombinedHomePage::new); disposable = combined.subscribe(this::successHomePage, this::throwableError); }
日志显示:
- 第一次运行:
Progress4未打印,但topSellerProductList数据正常 - 第二次运行:
Progress3未打印,但对应数据正常
问题原因
问题出在线程调度的冲突上:
- 你给每个源 Observable 都添加了
observeOn(AndroidSchedulers.mainThread()),这会把每个请求的onNext和onComplete事件都切换到主线程处理 - 结合
delaySubscription延迟订阅,主线程的消息队列会堆积多个请求的事件 Observable.zip会在主线程等待所有源的onNext事件完成后,立即发射组合结果并触发successHomePage,但此时部分源的onComplete事件可能还没被主线程调度执行,甚至因为消息队列的竞争被“遗漏”
简单说:虽然所有请求都成功返回了数据(onNext 事件被 zip 处理),但部分请求的 onComplete 事件因为主线程调度的时序问题,没有触发对应的 doOnComplete 回调。
解决方案
调整线程调度的方式,避免每个源 Observable 直接占用主线程,统一在 zip 后处理主线程逻辑,同时单独在 doOnComplete 中切换到主线程更新进度:
private void preFetchData() { progress = 0; // 初始化进度 ApiInterface apiService1 = ApiClient.getWooRxClient().create(ApiInterface.class); ApiInterface apiService2 = ApiClient.getRxClient().create(ApiInterface.class); Map<String, String> map1 = new HashMap<>(); map1.put("on_sale", "true"); Observable<List<Product>> call1 = apiService1.getProducts1(map1) .subscribeOn(Schedulers.io()) .doOnComplete(() -> updateProgress(18, "Progress1")); Map<String, String> map2 = new HashMap<>(); map2.put("featured", "true"); Observable<List<Product>> call2 = apiService1.getProducts1(map2) .subscribeOn(Schedulers.io()) .delaySubscription(100, TimeUnit.MILLISECONDS) .doOnComplete(() -> updateProgress(18, "Progress2")); Map<String, String> map3 = new HashMap<>(); map3.put("page", "1"); map3.put("sort", "rating"); map3.put("per_page", "10"); Observable<List<Product>> call3 = apiService2.getCustomProducts1(map3) .subscribeOn(Schedulers.io()) .delaySubscription(200, TimeUnit.MILLISECONDS) .doOnComplete(() -> updateProgress(18, "Progress3")); Map<String, String> map4 = new HashMap<>(); map4.put("page", "1"); map4.put("sort", "popularity"); map4.put("per_page", "10"); Observable<List<Product>> call4 = apiService2.getCustomProducts1(map4) .subscribeOn(Schedulers.io()) .delaySubscription(300, TimeUnit.MILLISECONDS) .doOnComplete(() -> updateProgress(18, "Progress4")); Observable<ResponseBody> call5 = apiService2.getCurrencySymbol() .subscribeOn(Schedulers.io()) .delaySubscription(400, TimeUnit.MILLISECONDS) .doOnComplete(() -> updateProgress(10, "Progress5")); Observable<List<Category>> call6 = apiService1.getAllCategories() .subscribeOn(Schedulers.io()) .delaySubscription(500, TimeUnit.MILLISECONDS) .doOnComplete(() -> updateProgress(18, "Progress6")); // 统一在 zip 后切换到主线程处理结果 Observable<CombinedHomePage> combined = Observable.zip(call1, call2, call3, call4, call5, call6, CombinedHomePage::new) .observeOn(AndroidSchedulers.mainThread()); disposable = combined.subscribe(this::successHomePage, this::throwableError); } // 单独封装进度更新方法,确保在主线程执行 private void updateProgress(int increment, String logTag) { Observable.just(increment) .observeOn(AndroidSchedulers.mainThread()) .subscribe(add -> { progress += add; Log.e(logTag, progress + ""); mProgressBar.setProgress(progress); mProgressBar.setProgressText(progress + "%"); }); } // 原有成功/失败回调保持不变 private void successHomePage(CombinedHomePage o) { Log.e("Response", "Featured List Size " + o.featuredProductList.size()); Log.e("Response", "Sale List Size " + o.saleProductList.size()); Log.e("Response", "Rated List Size " + o.topRatedProductList.size()); Log.e("Response", "Seller List Size " + o.topSellerProductList.size()); Log.e("Response", "Currency " + o.CURRENCY); Log.e("Response", "Category List Size " + o.categoryList.size()); } private void throwableError(Throwable t) { Log.e("Response", "Fail"); }
额外优化建议
因为 Retrofit 请求都是单次发射的,你可以用 Single 代替 Observable,语义更清晰,也能避免 Observable 的潜在问题:
Single<List<Product>> call1 = apiService1.getProducts1(map1) .subscribeOn(Schedulers.io()) .doOnSuccess(products -> updateProgress(18, "Progress1")); // 用 Single.zip 组合多个 Single Single<CombinedHomePage> combined = Single.zip(call1, call2, call3, call4, call5, call6, CombinedHomePage::new) .observeOn(AndroidSchedulers.mainThread());
内容的提问来源于stack exchange,提问作者Abhishek Singh
相关产品推荐
相关产品推荐

