【问题标题】:Rxjava retryWhen called instantlyRxjava 重试时立即调用
【发布时间】:2017-10-25 16:33:10
【问题描述】:

我对 rxjava 有一个非常具体的问题或误解,希望有人能提供帮助。

我正在运行 rxjava 2.1.5 并且有以下代码 sn-p:

public static void main(String[] args) {

    final Observable<Object> observable = Observable.create(emitter -> {
        // Code ... 
    });

    observable.subscribeOn(Schedulers.io())
            .retryWhen(error -> {
                System.out.println("retryWhen");
                return error.retry();
            }).subscribe(next -> System.out.println("subscribeNext"),
                         error -> System.out.println("subscribeError"));

}

执行后,程序打印:

retryWhen

Process finished with exit code 0

我的问题,我不明白的是:为什么在订阅 Observable 时立即调用 retryWhen? observable 什么都不做。

我想要的是在发射器上调用 onError 时调用 retryWhen。我是否误解了 rx 的工作原理?

谢谢!

添加新的sn-p:

public static void main(String[] args) throws InterruptedException {

    final Observable<Object> observable = Observable.create(emitter -> {
        emitter.onNext("next");
        emitter.onComplete();
    });

    final CountDownLatch latch = new CountDownLatch(1);
    observable.subscribeOn(Schedulers.io())
            .doOnError(error -> System.out.println("doOnError: " + error.getMessage()))
            .retryWhen(error -> {
                System.out.println("retryWhen: " + error.toString());
                return error.retry();
            }).subscribe(next -> System.out.println("subscribeNext"),
                         error -> System.out.println("subscribeError"),
                         () -> latch.countDown());

    latch.await();
}

发射器 onNext 和完成被调用。 DoOnError 永远不会被调用。输出是:

retryWhen:io.reactivex.subjects.SerializedSubject@35fb3008 订阅下一个

进程以退出代码 0 结束

【问题讨论】:

    标签: java rx-java


    【解决方案1】:

    retryWhenObserver 订阅它时调用提供的函数,因此您有一个主序列伴随一个序列,该序列发出主序列失败的Throwable。你应该在Observable 上编写一个逻辑到你在这个Function 中,所以最后,一个Throwable 会在另一端产生一个值。

    Observable.error(new IOException())
        .retryWhen(e -> {
             System.out.println("Setting up retryWhen");
             int[] count = { 0 };
             return e
                .takeWhile(v -> ++count[0] < 3)
                .doOnNext(v -> { System.out.println("Retrying"); });
        })
        .subscribe(System.out::println, Throwable::printStackTrace);
    

    由于 e -&gt; { } 函数体是针对每个单独的订阅者执行的,因此您可以安全地拥有每个订阅者的状态,例如重试计数器。

    使用e -&gt; e.retry() 无效,因为输入错误流永远不会调用其onError

    【讨论】:

    • 是的,我得出了这个结论。有时,写一个问题似乎可以帮助你解决它。谢谢!
    【解决方案2】:

    一个问题是,您不会再收到任何结果,因为您正在使用 retryWhen() 创建线程,但您的应用似乎已完成。要查看该行为,您可能需要一个 while 循环来保持您的应用程序运行。

    这实际上意味着您需要在代码末尾添加类似的内容:

    while (true) {}
    

    另一个问题是您不会在示例中发出任何错误。您需要发出至少一个值来调用onNext(),否则它不会重复,因为它正在等待它。

    这是一个工作示例,其中一个值,然后它发出错误并重复。你可以使用

    .retryWhen(errors -> errors)
    

    相同
    .retryWhen(errors -> errors.retry())
    

    工作样本:

     public static void main(String[] args) {
            Observable
                    .create(e -> {
                        e.onNext("test");
                        e.onError(new Throwable("test"));
                    })
                    .retryWhen(errors -> errors.retry())
                    .subscribeOn(Schedulers.io())
                    .subscribe(
                            next -> System.out.println("subscribeNext"),
                            error -> System.out.println("subscribeError"),
                            () -> System.out.println("onCompleted")
                    );
    
            while (true) {
    
            }
        }
    

    您需要发出结果的原因是,Observable 需要发出一个值,否则它会等到它收到一个新值。

    这是因为 onError 只能调用 onec(在订阅中),但 onNext 会发出 1..* 值。

    您可以使用 doOnError() 检查此行为,它会在每次重试 Observable 时为您提供错误。

    Observable
                .create(e -> e.onError(new Exception("empty")))
                .doOnError(e -> System.out.println("error received " + e))
                .retryWhen(errors -> errors.retry())
                .subscribeOn(Schedulers.io())
                .subscribe(
                        nextOrSuccess -> System.out.println("nextOrSuccess " + nextOrSuccess),
                        error -> System.out.println("subscribeError")
                );
    

    【讨论】:

    • 我的 sn-p 只是我的问题的再现。真正的代码调用发射器 onNext 错误并完成......问题是发射器上的 onError 永远不会被调用。 retryWhen 在调用订阅时被调用,我觉得这很奇怪和错误。鉴于我发布的代码,它是可运行的。是否应该调用 retryWhen 中的函数设置?这是正确的行为吗?
    • 在两者之间放置一个 .doOnError(),你会看到幕后真的发生了 :-)
    • 我正在发布一个新示例。
    • 等一下,启动 netbeans
    • 我认为我的误解是 retryWhen 将一个带有可观察值的函数作为输入而不是实际错误。
    猜你喜欢
    • 2016-08-16
    • 2023-03-11
    • 2013-02-06
    • 1970-01-01
    • 2014-07-04
    • 1970-01-01
    • 1970-01-01
    • 2013-06-01
    • 2011-12-09
    相关资源
    最近更新 更多