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

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 未打印,但对应数据正常

问题原因

问题出在线程调度的冲突上:

  1. 你给每个源 Observable 都添加了 observeOn(AndroidSchedulers.mainThread()),这会把每个请求的 onNext 和 onComplete 事件都切换到主线程处理
  2. 结合 delaySubscription 延迟订阅,主线程的消息队列会堆积多个请求的事件
  3. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:49:22