【问题标题】:How to cancel Observable.timer with rxjava2?如何使用 rxjava2 取消 Observable.timer?
【发布时间】:2017-05-30 22:52:52
【问题描述】:

我的代码中有一个重试机制,我使用下面的行来执行我的重试逻辑。例如,我生成一个随机毫秒来延迟我的执行。当计时器滴答到 30 * 1000 毫秒时,我想取消此计时器。如何取消此计时器并立即执行我的逻辑。

//register retryWhen
Observable.retryWhen(new RetryWhenException());

//retry code.
public class RetryWhenException implements Function<Observable<? extends Throwable>, Observable<?>>{       
public Observable<?> apply(final Observable<? extends Throwable> observable) throws Exception {
            return observable.zipWith(Observable.range(1, count + 1), new BiFunction<Throwable, Integer, Wrapper>() {
                @Override
                public Wrapper apply(Throwable throwable, Integer integer) throws Exception {
                    return new Wrapper(throwable, integer);
                }
            }).flatMap(new Function<Wrapper, Observable<?>>() {
                @Override
                public Observable<?> apply(Wrapper wrapper) throws Exception {
                    long delay = 60 * 1000;
                    //How can I add some code here to cancel this time and execute this api call immediately if I receive a event like network gets back?
                    return Observable.timer(delay, TimeUnit.MILLISECONDS);
                }
            });
        }
}

提前致谢。

【问题讨论】:

    标签: android rx-java reactive-programming retrofit2 rx-java2


    【解决方案1】:

    如果我正确理解您的问题,您希望能够将多个条件组合到重试中,这意味着重试最多会在一定时间后发生(计时器),甚至在某些事件发生后更早(网络已连接) )。
    在这种情况下,您需要结合这两个事件,首先,您需要一些Observable 来通知有关网络事件(如何创建它是不同的讨论,用 Observable 包装系统广播事件应该不是问题) ,那么你可以这样做:

    private Observable<NetworkState> networkStateObservable;
    
    public class RetryWhenException implements Function<Observable<? extends Throwable>, Observable<?>> {
        public Observable<?> apply(
                final Observable<? extends Throwable> observable) throws Exception {
            return observable.zipWith(Observable.range(1, count + 1), Wrapper::new)
                    .flatMap(wrapper -> {
                        long delay = 60 * 1000;
                        Observable<NetworkState> networkConnectedEvents =
                                networkStateObservable.filter(networkState -> networkState.isConnected())
                                   take(1);
                        Observable<Long> timer = Observable.timer(delay, TimeUnit.MILLISECONDS);
    
                        return Observable.amb(Arrays.asList(networkConnectedEvents, timer));
                    });
        }
    }
    

    网络状态Observable 被过滤以仅在连接时获得通知,take(1) 也是为了确保在第一次收到通知后我们将取消订阅它(无需进一步收听)。
    @ 987654325@ 操作符在这里看起来很合适,因为它将选择发出 firs 的 Observable,并取消订阅另一个,这意味着如果网络 Observable 在定时器之前发出,定时器 Observable 将被取消订阅( = 计时器将被取消)。

    编辑:

    移除了错误的 takeUntil(networkConnectedEvents),因为 amb 会在需要时取消订阅。

    【讨论】:

    • 嗨 yosriz,你能帮忙看看我应该如何正确发出 networkStateObservable 吗?我在下面发布了我的代码。非常感谢。
    【解决方案2】:

    非常感谢您的回复。这正是我一直在寻找的。正如你所描述的,我创建了一个网络 Observable 包装器, 似乎在调用 onNetworkChanged() 后 networkStateObservable 没有发出(state.isConnected() 为真)。 如何使 networkStateObservable 正确发射?非常感谢。 我在

    中调用了以下行
    e.onNext(state);  
     e.onComplete();
    

    发出 networkStateObservable 我错了吗?

    Observable<NetworkState> networkStateObservable = getNetworkObservableWrapper();
    public Observable<NetworkState> getNetworkObservableWrapper() {
            return Observable.create(new ObservableOnSubscribe<NetworkState>() {
                @Override
                public void subscribe(final ObservableEmitter<NetworkState> e) throws Exception {
                    final NetworkChangedListener listener = new NetworkChangedListener() {
                        @Override
                        public void onNetworkChanged(NetworkState state) {
                            //Below lines were called, but Observable.amb(Arrays.asList(networkConnectedEvents, timer)); didn't work.                       
                            e.onNext(state);
                            e.onComplete();
                        }
                    };
                    registerNetworkChanged(listener);
                }
            });
        }
    

    【讨论】:

    • 它似乎是由 .takeUntil(networkConnectedEvents) 引起的。如果我删除了这一行,它工作正常。
    • 是的,其实我放错了,我把它从答案中删除了,amb会取消订阅(如果需要,从计时器取消订阅)
    • 猜测问题:如果您的 registerNetworkChanged() 覆盖现有的网络更改侦听器,可能会发生此 Observable 订阅了两次(一次在 amb 和一次在 takeUntil),因此第二个 takeUntil 可能会覆盖现有的监听器。
    • 几个 cmets,根据我对根本原因的猜测,a) 当您注册到网络更改事件时,添加一个侦听器而不是设置 b) 除此之外,您应该使用 @987654324 从侦听器逻辑中添加注销@
    • 非常感谢您的解释。
    猜你喜欢
    • 1970-01-01
    • 2019-04-14
    • 2019-10-17
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-09
    • 1970-01-01
    • 2017-10-17
    相关资源
    最近更新 更多