【问题标题】:RxJava onErrorResumeNext()RxJava onErrorResumeNext()
【发布时间】:2014-10-27 08:59:13
【问题描述】:

我有两个 observables(为简单起见命名为 A 和 B)和一个订阅者。因此,订阅者订阅了 A,如果 A 上出现错误,则 B(这是回退)启动。现在,每当 A 遇到错误时,B 都会被正常调用,但是 A 在订阅者上调用 onComplete(),所以 B 响应即使 B 执行成功,也永远不会到达订阅者。

这是正常行为吗?我认为 onErrorResumeNext() 应该继续流,并在完成后通知订阅者,如文档 (https://github.com/ReactiveX/RxJava/wiki/Error-Handling-Operators#onerrorresumenext) 中所述。

这是我正在做的整体结构(省略了几个“无聊”的代码):

public Observable<ModelA> observeGetAPI(){
    return retrofitAPI.getObservableAPI1()
            .flatMap(observableApi1Response -> {
                ModelA model = new ModelA();

                model.setApi1Response(observableApi1Response);

                return retrofitAPI.getObservableAPI2()
                        .map(observableApi2Response -> {
                            // Blah blah blah...
                            return model;
                        })
                        .onErrorResumeNext(observeGetAPIFallback(model))
                        .subscribeOn(Schedulers.newThread())
            })
            .onErrorReturn(throwable -> {
                // Blah blah blah...
                return model;
            })
            .subscribeOn(Schedulers.newThread());
}

private Observable<ModelA> observeGetAPIFallback(ModelA model){
    return retrofitAPI.getObservableAPI3().map(observableApi3Response -> {
        // Blah blah blah...
        return model;
    }).onErrorReturn(throwable -> {
        // Blah blah blah...
        return model;
    })
    .subscribeOn(Schedulers.immediate());
}

Subscription subscription;
subscription = observeGetAPI.subscribe(ModelA -> {
    // IF THERE'S AN ERROR WE NEVER GET B RESPONSE HERE...
}, throwable ->{
    // WE NEVER GET HERE... onErrorResumeNext()
},
() -> { // IN CASE OF AN ERROR WE GET STRAIGHT HERE, MEANWHILE, B GETS EXECUTED }
);

任何想法我做错了什么?

谢谢!

编辑: 以下是正在发生的事情的大致时间表:

---> HTTP GET REQUEST B
<--- HTTP 200 REQUEST B RESPONSE (SUCCESS)

---> HTTP GET REQUEST A
<--- HTTP 200 REQUEST A RESPONSE (FAILURE!)

---> HTTP GET FALLBACK A
** onComplete() called! ---> Subscriber never gets fallback response since onComplete() gets called before time.
<--- HTTP 200 FALLBACK A RESPONSE (SUCCESS)

这是我制作的一个简单图表的链接,它代表了我想要发生的事情: Diagram

【问题讨论】:

  • 您的时间线显示失败响应的 HTTP 200。是否有其他方法可以从 getObservableAPI2() 发出错误信号?另外,您能否指定哪些 API 请求对应于时间线输出?它看起来像 getObservableAPI1->REQUEST B, getObservableAPI2->REQUEST A, getObservableAPI3->FALLBACK A 但我只是想确定一下。
  • 是的,实际上虽然响应是 200,但有些数据可能为空,所以我在这些情况下抛出错误。是的,这就是时间线请求关系,我会尽快编辑问题以匹配时间线请求。
  • 你的逻辑看起来不错。您应该在 onComplete 之前获得回退响应。你能删除所有的 subscribeOn() 调用,看看会发生什么。它们不应该是必需的,因为 Retrofit 无论如何都会在自己的线程池上执行请求。
  • @kjones 我已经尝试过了,得到了完全相同的输出,onComplete 被调用得太早了。
  • 最好将你的链扁平化而不是嵌套它(使它超级难以阅读、跟踪和调试)。非常不清楚您在这里要做什么,尤其是。在flatMap 块内。请整理好你的方法和变量,不管天气是不是改造

标签: java retrofit rx-java


【解决方案1】:

下面使用的 Rx 调用应该模拟您使用 Retrofit 所做的事情。

fallbackObservable =
        Observable
                .create(new Observable.OnSubscribe<String>() {
                    @Override
                    public void call(Subscriber<? super String> subscriber) {
                        logger.v("emitting A Fallback");
                        subscriber.onNext("A Fallback");
                        subscriber.onCompleted();
                    }
                })
                .delay(1, TimeUnit.SECONDS)
                .onErrorReturn(new Func1<Throwable, String>() {
                    @Override
                    public String call(Throwable throwable) {
                        logger.v("emitting Fallback Error");
                        return "Fallback Error";
                    }
                })
                .subscribeOn(Schedulers.immediate());

stringObservable =
        Observable
                .create(new Observable.OnSubscribe<String>() {
                    @Override
                    public void call(Subscriber<? super String> subscriber) {
                        logger.v("emitting B");
                        subscriber.onNext("B");
                        subscriber.onCompleted();
                    }
                })
                .delay(1, TimeUnit.SECONDS)
                .flatMap(new Func1<String, Observable<String>>() {
                    @Override
                    public Observable<String> call(String s) {
                        logger.v("flatMapping B");
                        return Observable
                                .create(new Observable.OnSubscribe<String>() {
                                    @Override
                                    public void call(Subscriber<? super String> subscriber) {
                                        logger.v("emitting A");
                                        subscriber.onNext("A");
                                        subscriber.onCompleted();
                                    }
                                })
                                .delay(1, TimeUnit.SECONDS)
                                .map(new Func1<String, String>() {
                                    @Override
                                    public String call(String s) {
                                        logger.v("A completes but contains invalid data - throwing error");
                                        throw new NotImplementedException("YUCK!");
                                    }
                                })
                                .onErrorResumeNext(fallbackObservable)
                                .subscribeOn(Schedulers.newThread());
                    }
                })
                .onErrorReturn(new Func1<Throwable, String>() {
                    @Override
                    public String call(Throwable throwable) {
                        logger.v("emitting Return Error");
                        return "Return Error";
                    }
                })
                .subscribeOn(Schedulers.newThread());

subscription = stringObservable.subscribe(
        new Action1<String>() {
            @Override
            public void call(String s) {
                logger.v("onNext " + s);
            }
        },
        new Action1<Throwable>() {
            @Override
            public void call(Throwable throwable) {
                logger.v("onError");
            }
        },
        new Action0() {
            @Override
            public void call() {
                logger.v("onCompleted");
            }
        });

日志语句的输出是:

RxNewThreadScheduler-1 发射 B RxComputationThreadPool-1 flatMapping B RxNewThreadScheduler-2 发出 A RxComputationThreadPool-2 A 完成但包含无效数据 - 抛出错误 RxComputationThreadPool-2 发出回退 RxComputationThreadPool-1 onNext 一个后备 RxComputationThreadPool-1 onCompleted

这似乎是您正在寻找的东西,但也许我错过了一些东西。

【讨论】:

    猜你喜欢
    • 2016-12-02
    • 1970-01-01
    • 1970-01-01
    • 2018-11-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多