【问题标题】:How to ignore error and continue infinite stream?如何忽略错误并继续无限流?
【发布时间】:2015-05-12 06:16:22
【问题描述】:

我想知道如何忽略异常并继续无限流(在我的情况下是位置流)?

我正在获取当前用户位置(使用 Android-ReactiveLocation),然后将它们发送到我的 API(使用 Retrofit)。

在我的情况下,当网络调用期间发生异常(例如超时)时,onError 方法被调用并且流自行停止。如何避免?

活动:

private RestService mRestService;
private Subscription mSubscription;
private LocationRequest mLocationRequest = LocationRequest.create()
            .setPriority(LocationRequest.PRIORITY_HIGH_ACCURACY)
            .setInterval(100);
...
private void start() {
    mRestService = ...;
    ReactiveLocationProvider reactiveLocationProvider = new ReactiveLocationProvider(this);
    mSubscription = reactiveLocationProvider.getUpdatedLocation(mLocationRequest)
            .buffer(50)
            .flatMap(locations -> mRestService.postLocations(locations)) // can throw exception
            .subscribeOn(Schedulers.newThread())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe();
}

休息服务:

public interface RestService {
    @POST("/.../")
    Observable<Response> postLocations(@Body List<Location> locations);
}

【问题讨论】:

标签: android rx-java rx-android


【解决方案1】:

您可能想要使用error handling operators 之一。

  • onErrorResumeNext( ) — 指示 Observable 在遇到错误时发出一系列项目
  • onErrorReturn( ) — 指示 Observable 在遇到错误时发出特定项目
  • onExceptionResumeNext( ) — 指示 Observable 在遇到异常后继续发射项目(但不是其他种类的 throwable)
  • retry( ) — 如果源 Observable 发出错误,请重新订阅它,希望它能够正确完成
  • retryWhen( ) — 如果源 Observable 发出错误,则将该错误传递给另一个 Observable 以确定是否重新订阅源

特别是 retryonExceptionResumeNext 在你的情况下看起来很有希望。

【讨论】:

  • 当我在flatMap(locations -&gt; mRestService.postLocations(locations)) onCompleted 被调用并且流结束之后添加onExceptionResumeNext(Observable.empty())
  • 不要在flatMap之后添加,而是在flatMap里面。
  • 谢谢! onErrorResumeNext() - 非常有用的后备构造。
【解决方案2】:

mRestService.postLocations(locations) 发出一项,然后完成。 如果发生错误,则发出错误,从而完成流。

当您在 flatMap 中调用此方法时,错误会继续出现在您的“主”流中,然后您的流停止。

您可以做的是将您的错误转换为另一个项目(如此处所述:https://stackoverflow.com/a/28971140/476690),但不是在您的主流上(我想您已经尝试过),而是在mRestService.postLocations(locations) 上。

这样,此调用将发出一个错误,该错误将被转换为一个项目/另一个可观察对象,然后完成。 (无需致电onError)。

在消费者视图中,mRestService.postLocations(locations) 将发出一个项目,然后完成,就像一切都成功一样。

mSubscription = reactiveLocationProvider.getUpdatedLocation(mLocationRequest)
        .buffer(50)
        .flatMap(locations -> mRestService.postLocations(locations).onErrorReturn((e) -> Collections.emptyList()) // can't throw exception
        .subscribeOn(Schedulers.newThread())
        .observeOn(AndroidSchedulers.mainThread())
        .subscribe();

【讨论】:

  • 不幸的是它不是忽略,而是发出空列表。最好不要在错误时发出
  • // can throw exception 是什么意思?这是否意味着即使在onErrorReturn 内部我们也可能再次出现异常?!
  • 我认为这是一个错误,而是:“不能抛出异常”(所以我只是编辑评论。)很好!
【解决方案3】:

如果您只想忽略flatMap 中的错误而不返回元素,请执行以下操作:

flatMap(item -> 
    restService.getSomething(item).onErrorResumeNext(Observable.empty())
);

【讨论】:

  • 这实际上会完成流。不会吗?
【解决方案4】:

只需粘贴来自@MikeN 答案的链接信息,以防丢失:

import rx.Observable.Operator;
import rx.functions.Action1;

public final class OperatorSuppressError<T> implements Operator<T, T> {
    final Action1<Throwable> onError;

    public OperatorSuppressError(Action1<Throwable> onError) {
        this.onError = onError;
    }

    @Override
    public Subscriber<? super T> call(final Subscriber<? super T> t1) {
        return new Subscriber<T>(t1) {

            @Override
            public void onNext(T t) {
                t1.onNext(t);
            }

            @Override
            public void onError(Throwable e) {
                onError.call(e);
            }

            @Override
            public void onCompleted() {
                t1.onCompleted();
            }

        };
    }
}

并在可观察源附近使用它,因为其他操作员可能 在此之前急切地退订。

Observerable.create(connectToUnboundedStream()).lift(new OperatorSuppressError(log()).doOnNext(someStuff()).subscribe();

但是请注意,这会抑制从 资源。如果链中的任何 onNext 在它抛出异常之后,它是 源仍然可能会被取消订阅。

【讨论】:

  • 这可行,但在抑制错误后,我的源 observable 似乎停止工作。有办法重启吗?
  • @Matthias 所以这不能解决问题吗? (与 onErrorResumeNext 相同 - 它完成了 observable)?
【解决方案5】:

这个答案可能有点晚了,但如果有人偶然发现这个问题,可以使用 Jacke Wharton 的即用型 Relay 库,而不是重新发明轮子

https://github.com/JakeWharton/RxRelay

有很好的文档,但本质上,Relay 是一个 Subject,除了不能调用 onComplete 或 onError。

选项有:

行为中继

Relay that emits the most recent item it has observed and all subsequent observed items to each subscribed Observer.
    // observer will receive all events.
    BehaviorRelay<Object> relay = BehaviorRelay.createDefault("default");
    relay.subscribe(observer);
    relay.accept("one");
    relay.accept("two");
    relay.accept("three");

    // observer will receive the "one", "two" and "three" events, but not "zero"
    BehaviorRelay<Object> relay = BehaviorRelay.createDefault("default");
    relay.accept("zero");
    relay.accept("one");
    relay.subscribe(observer);
    relay.accept("two");
    relay.accept("three");

发布中继 中继,一旦观察者订阅,将所有后续观察到的项目发送给订阅者。

    PublishRelay<Object> relay = PublishRelay.create();
    // observer1 will receive all events
    relay.subscribe(observer1);
    relay.accept("one");
    relay.accept("two");
    // observer2 will only receive "three"
    relay.subscribe(observer2);
    relay.accept("three");

重播中继 Relay 缓存它观察到的所有项目,并将它们重播给任何订阅的观察者。

    ReplayRelay<Object> relay = ReplayRelay.create();
    relay.accept("one");
    relay.accept("two");
    relay.accept("three");
    // both of the following will get the events from above
    relay.subscribe(observer1);
    relay.subscribe(observer2);

【讨论】:

  • 这应该是公认的答案。 Repay 对于事件总线类型的广播非常方便
【解决方案6】:

尝试在 Observable.defer 调用中调用其余服务。这样每次调用你都有机会使用它自己的“onErrorResumeNext”,错误不会导致你的主流完成。

reactiveLocationProvider.getUpdatedLocation(mLocationRequest)
  .buffer(50)
  .flatMap(locations ->
    Observable.defer(() -> mRestService.postLocations(locations))
      .onErrorResumeNext(<SOME_DEFAULT_TO_REACT_TO>)
  )
........

该解决方案最初来自此线程-> RxJava Observable and Subscriber for skipping exception?,但我认为它也适用于您的情况。

【讨论】:

    【解决方案7】:

    添加我对这个问题的解决方案:

    privider
        .compose(ignoreErrorsTransformer)
        .subscribe()
    
    private final Observable.Transformer<ResultType, ResultType> ignoreErrorsTransformer =
            new Observable.Transformer<ResultType, ResultType>() {
                @Override
                public Observable<ResultType> call(Observable<ResultType> resultTypeObservable) {
                    return resultTypeObservable
                            .materialize()
                            .filter(new Func1<Notification<ResultType>, Boolean>() {
                                @Override
                                public Boolean call(Notification<ResultType> resultTypeNotification) {
                                    return !resultTypeNotification.isOnError();
                                }
                            })
                            .dematerialize();
    
                }
            };
    

    【讨论】:

    • 具体化操作符是什么?
    【解决方案8】:

    稍微修改解决方案 (@MikeN) 以使有限流能够完成:

    import rx.Observable.Operator;
    import rx.functions.Action1;
    
    public final class OperatorSuppressError<T> implements Operator<T, T> {
        final Action1<Throwable> onError;
    
        public OperatorSuppressError(Action1<Throwable> onError) {
            this.onError = onError;
        }
    
        @Override
        public Subscriber<? super T> call(final Subscriber<? super T> t1) {
            return new Subscriber<T>(t1) {
    
                @Override
                public void onNext(T t) {
                    t1.onNext(t);
                }
    
                @Override
                public void onError(Throwable e) {
                    onError.call(e);
                    //this will allow finite streams to complete
                    t1.onCompleted();
                }
    
                @Override
                public void onCompleted() {
                    t1.onCompleted();
                }
    
            };
        }
    }
    

    【讨论】:

      【解决方案9】:

      这是我用于忽略错误的 kotlin 扩展函数

      fun <T> Observable<T>.ignoreErrors(errorHandler: (Throwable) -> Unit) =
          retryWhen { errors ->
              errors
                  .doOnNext { errorHandler(it) }
                  .map { 0 }
          }
      

      这利用retryWhen 无限期地重新订阅上游,同时仍然允许您以非终端方式处理错误。

      感觉很危险

      【讨论】:

        【解决方案10】:

        使用 Rxjava2,我们可以调用带有 delayErrors 参数的重载平面图: flatmap javadoc

        当传递为真时:

        当前 Flowable 和所有内部 Publisher 的异常都会延迟,直到它们全部终止,如果为 false,则第一个发出异常信号的异常将立即终止整个序列

        【讨论】:

          【解决方案11】:

          您可以使用 onErrorComplete() 方法跳过错误

          mSubscription = reactiveLocationProvider.getUpdatedLocation(mLocationRequest)
              .buffer(50)
              .flatMapMaybe(locations -> Maybe.just(mRestService.postLocations(locations).onErrorComplete()) // skip item
              .subscribeOn(Schedulers.newThread())
              .observeOn(AndroidSchedulers.mainThread())
              .subscribe();
          

          【讨论】:

            猜你喜欢
            • 1970-01-01
            • 2013-03-13
            • 1970-01-01
            • 1970-01-01
            • 1970-01-01
            • 2011-05-07
            • 1970-01-01
            • 1970-01-01
            • 2017-09-23
            相关资源
            最近更新 更多