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

RxJava中等待两个Observable完成的最佳实现及相关疑问

Hey there! Let's break down your questions one by one and fix up that loading page logic to meet your needs.

1. Can we avoid using BiFunction and returning null?

Absolutely! The core issue here is that zip requires a merge function because it's designed to combine emissions from multiple Observables. But since you don't care about the combined result—you just need to know when both requests finish (success or failure)—we can switch to using Completable instead.

Completable doesn't emit any data; it only signals when an operation completes (successfully or with an error, if we let it). By converting each of your request Observables to Completable and using Completable.mergeArray, we can eliminate the need for a BiFunction entirely, no null returns required.

2. How does zip behave when one request fails?

This is a key gotcha with your current setup: zip will immediately terminate and emit an error as soon as any of the source Observables emits an error. It won't wait for the other request to finish, and your onAllCallsComplete() method will never run if either request fails.

Your doOnError handlers for A and B will still execute (since they're part of each individual Observable's chain), but the zip operator will short-circuit before calling your merge function. The subscriber's onError will also be triggered immediately, which might not be what you want if you need to wait for both requests to wrap up before navigating.

3. Is your current implementation the optimal solution?

Nope—for the reason mentioned above: it doesn't guarantee onAllCallsComplete() runs if either request fails. Let's fix that with a better approach.

Use Completable to handle both requests in parallel, ensuring we wait for both to finish (success or failure) before triggering the page jump. Here's how to adjust your code:

private Disposable retrieveBothThings() { 
    return Completable.mergeArray(
        getThingACompletable(),
        getThingBCompletable()
    )
    .subscribeOn(Schedulers.io()) 
    .observeOn(AndroidSchedulers.mainThread()) 
    .subscribe(this::onAllCallsComplete, Logger::e); 
} 

private Completable getThingACompletable() { 
    return SessionManager.getInstance().getApi()
        .getA() 
        .subscribeOn(Schedulers.io()) 
        .observeOn(AndroidSchedulers.mainThread()) 
        .doOnNext(this::onACompleted) 
        .doOnError(this::onAFailed)
        .ignoreElements() // Convert Observable to Completable (discard data)
        .onErrorComplete(); // Treat errors as "completed" so we don't terminate early
} 

private Completable getThingBCompletable() { 
    return SessionManager.getInstance().getApi()
        .getB() 
        .subscribeOn(Schedulers.io()) 
        .observeOn(AndroidSchedulers.mainThread())
        .toObservable() 
        .doOnNext(this::onBSuccess) 
        .doOnError(this::onBFailure)
        .ignoreElements()
        .onErrorComplete(); 
} 

Why this works:

  • ignoreElements() converts each Observable to a Completable, since we don't need the response data.
  • onErrorComplete() tells each Completable to emit a "complete" event instead of an error, so mergeArray will wait for both operations to finish—even if one or both fail.
  • Completable.mergeArray runs both requests in parallel and only emits its own "complete" event when both source Completables have finished. This ensures onAllCallsComplete() always runs after both requests are done.

Bonus: If you need to track success/failure status

If you later want to know which requests succeeded or failed, you can wrap each result in a status object instead of using Completable. For example:

// A simple class to hold request status
class RequestResult<T> {
    final T data;
    final Throwable error;

    RequestResult(T data, Throwable error) {
        this.data = data;
        this.error = error;
    }
}

private Disposable retrieveBothThings() { 
    return Observable.zip(
        getThingAResult(),
        getThingBResult(),
        (resultA, resultB) -> {
            onAllCallsComplete();
            return null; // We still don't need the combined result, but zip requires it
        }
    )
    .subscribeOn(Schedulers.io()) 
    .observeOn(AndroidSchedulers.mainThread()) 
    .subscribe(o -> {}, Logger::e); 
} 

private Observable<RequestResult<A>> getThingAResult() { 
    return SessionManager.getInstance().getApi()
        .getA() 
        .subscribeOn(Schedulers.io()) 
        .observeOn(AndroidSchedulers.mainThread()) 
        .map(data -> new RequestResult<>(data, null))
        .doOnNext(result -> {
            if (result.data != null) onACompleted(result.data);
        })
        .onErrorReturn(throwable -> {
            onAFailed(throwable);
            return new RequestResult<>(null, throwable);
        });
} 

This way, zip will always receive emissions from both requests (success or failure), so your merge function runs reliably. But for your original use case (just waiting for both to finish), the Completable approach is cleaner and simpler.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 09:16:53