【问题标题】:What is the proper way of calling d.dispose() or s.cancel() method?调用 d.dispose() 或 s.cancel() 方法的正确方法是什么?
【发布时间】:2019-01-29 04:59:42
【问题描述】:

在 RxJava2 中,Observer 接口和 Subscriber 接口中引入了一个新方法,命名为 -

interface Subscriber<T>{             
 @Override
 public void onSubscribe(Subscription s)
 {
      s.cancel();
      s.request(5);
 }
....
}

interface Observer<T>{             
     @Override
     public void onSubscribe(Disposable d)
     {
          d.dispose();
     }
    ....
    }

我看到 onSubscribe() 方法总是在 onNext(T t) 之前第一次调用,我也阅读了该文档,发现它的用途是如果您的工作是使用特定的 Observable 完成的,则用于处置资源。

问题是我们如何才能知道我们的工作已经完成并disposecancel来源或来源与消费者之间的联系? 那么调用 d.dispose()s.cancel()s.request(7) 的更好方法是什么? p>

【问题讨论】:

  • 我不确定我是否理解您的要求。 how we can know initially that our job is done 什么意思?
  • 我的意思是当我可以确定这是调用 d.dispose() 的正确时间时,因为在刚刚调用之后,与 observable 的连接将会丢失。那么我应该什么时候调用它,在什么条件下?

标签: rx-java rx-android rx-java2


【解决方案1】:

流可以通过两种方式终止:

  1. 错误

  2. 完成

据我所知,在这两种情况下,您都不需要调用 dispose/cancel。确实是reactive stream contract says

当 Observable 向其发出 OnError 或 OnComplete 通知时 观察者,这将结束订阅。

当然,您可以在任何时候停止您的流,在它以错误结束之前或因为它完成。在这些情况下,您必须调用 dispose/cancel。使用:

  • dispose()Observable
  • cancel()Flowable

关于request()方法,如果你想建立一个“reactive pull”,你需要它,我认为它与cancel无关。您可以找到更多信息here

【讨论】:

    【解决方案2】:

    您很少需要调用这些方法,因为take 和类似的其他运算符会为您限制流。此外,DisposableObserverDisposableSubscriber 助手类为您管理 Disposable/Subscription

    在一些特殊的消费者中,您可能希望从onSubscribe() 调用Subscription.request(1),然后从onNext() 调用,但没有实际理由从onError() 或@987654332 调用request()cancel() @ 也没有效果。

    例如,以下代码将通过request(1) 应用背压,因为它使用一个序列并将其转发到异步后处理逻辑,否则 RxJava 无法识别:

    ExecutorService exec = Executors.newSingleThreadedExecutor();
    
    source.subscribe(new Subscriber<Data>() {
        Subscription upstream;
        @Override public void onSubscribe(Subscription s) {
            upstream = s;
            s.request(1);
        }
    
        @Override public void onNext(Data t) {
            exec.submit(() -> {
               if (t.isValid()) {
                   process(t.details);
                   upstream.request(1);
               } else {
                   upstream.cancel();
                   exec.shutdown();
               }
            });
        }
    
        @Override public void onError(Throwable ex) {
            ex.printStackTrace();
            exec.shutdown();
        }
    
        @Override public void onComplete() {
            exec.shutdown();
        }
    });
    

    同样,这通常很少见。在常规的Subscribers 上,您只需调用s.request(Long.MAX_VALUE),因为调用堆栈阻塞特性确保最后一个上游阶段不会压倒onNext

    source.subscribe(new Subscriber<Data>() {
        Subscription upstream;
        @Override public void onSubscribe(Subscription s) {
            upstream = s;
            s.request(Long.MAX_VALUE);
        }
    
        @Override public void onNext(Data t) {
           if (t.isValid()) {
               process(t.details);
           } else {
               upstream.cancel();
           }
        }
    
        @Override public void onError(Throwable ex) {
            ex.printStackTrace();
        }
    
        @Override public void onComplete() {
        }
    });
    

    总之,当没有更多可用数据时调用onComplete,当您想指示即使有更多可用数据也不应进行进一步处理时,您调用cancel

    【讨论】:

      猜你喜欢
      • 2011-11-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-09-29
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多