【问题标题】:Can I trace consuming of events in RxJava Subscriber?我可以跟踪 RxJava 订阅者中的事件消耗吗?
【发布时间】:2017-05-29 21:44:46
【问题描述】:

我想跟踪订阅者何时开始消费事件以及何时完成。 是否有适用于所有 Observables/Subscribers 的通用方法?

【问题讨论】:

  • 有一个 onSubscribe(disposable) 方法会在 Observable 开始发出事件时被调用。当onComplete()/onError 被调用意味着它终止了。
  • 谢谢,但我正在寻找像 RxJavaHooks 这样的通用解决方案。所以我不需要访问每个可观察的。对于当前 jvm 实例中的任何 observable,它都应该像间谍一样工作。
  • 我明白了。你可以在 RxJavaPlugin 中使用setOnObservableSubscribe
  • 对不起,我不太明白如何使用它。此挂钩仅适用于“订阅”事件。我需要跟踪订阅者何时开始使用事件以及何时完成。

标签: rx-java rx-java2


【解决方案1】:

是 1.x 还是 2.x?对于 2.x,它可能会变得非常复杂,因为必须考虑所有内部协议不会意外地取消优化您的流程。

否则,它可以像写一个Observer 一样简单,在真正的Observer 和运算符之间填充:

import io.reactivex.Observer;

RxJavaPlugins.setOnObservableSubscribe((observable, observer) -> {
    if (!observable.getClass().getName().toLowerCase().contains("map")) {
        return observer;
    }

    System.out.println("Started");

    class SignalTracker implements Observer<Object>, Disposable {
        Disposable upstream;
        @Override public void onSubscribe(Disposable d) {
            upstream = d;
            // write the code here that has to react to establishing the subscription
            observer.onSubscribe(this);
        }
        @Override public void onNext(Object o) {
            // handle onNext before or aftern notifying the downstream
            observer.onNext(o);
        }
        @Override public void onError(Throwable t) {
            // handle onError
            observer.onError(t);
        }
        @Override public void onComplete() {
            // handle onComplete
            System.out.println("Completed");
            observer.onComplete();
        }
        @Override public void dispose() {
            // handle dispose
            upstream.dispose();
        }
        @Override public boolean isDisposed() {
            return upstream.isDisposed();
        }
    }
    return new SignalTracker();
  });

  Observable<Integer> observable = Observable.range(1, 5)
      .subscribeOn(Schedulers.io())
      .observeOn(Schedulers.computation())
      .map(integer -> {
        try {
          TimeUnit.SECONDS.sleep(1);
        } catch (InterruptedException e) {
          e.printStackTrace();
        }
        return integer * 3;
      });

  observable.subscribe(System.out::println);

  Thread.sleep(6000L);

打印:

Started
3
6
9
12
15
Completed

编辑: RxJava 1 版本需要更多的 lambdas 但可行:

RxJavaHooks.setOnObservableStart((observable, onSubscribe) -> {
    if (!onSubscribe.getClass().getName().toLowerCase().contains("map")) {
        return onSubscribe;
    }

    System.out.println("Started");

    return (Observable.OnSubscribe<Object>)observer -> {
        class SignalTracker extends Subscriber<Object> {
            @Override public void onNext(Object o) {
                // handle onNext before or aftern notifying the downstream
                observer.onNext(o);
            }
            @Override public void onError(Throwable t) {
                // handle onError
                observer.onError(t);
            }
            @Override public void onCompleted() {
                // handle onComplete
                System.out.println("Completed");
                observer.onCompleted();
            }
            @Override public void setProducer(Producer p) {
                observer.setProducer(p);
            }
        }
        SignalTracker t = new SignalTracker()
        observer.add(t);
        onSubscribe.call(t);
    };
  });

  Observable<Integer> observable = Observable.range(1, 5)
      .subscribeOn(Schedulers.io())
      .observeOn(Schedulers.computation())
      .map(integer -> {
        try {
          TimeUnit.SECONDS.sleep(1);
        } catch (InterruptedException e) {
          e.printStackTrace();
        }
        return integer * 3;
      });

  observable.subscribe(System.out::println);

  Thread.sleep(6000L);

【讨论】:

  • 我在说明中添加了示例。此解决方案不适用于它。
  • 在订阅者完成所有事件的处理之前打印“已完成”。
  • 钩子添加了对所有操作符的跟踪:range、subscribeOn、observeOn、map,因此你得到了 4 个“开始”。此外,如果你睡 6 秒,你可以看到 4 完成。
  • @smalafeev 这正是我告诉你需要为你的 observable 设置条件的原因。每个运算符都会从中生成一个新的 Observable,例如:ObservableMap。您的上游 observable 收到了 onComplete 事件,但您的下游没有。因此,您会在地图之前看到 onComplete。
  • 现在建链有问题。我找不到生成的 observables 之间的关系。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2014-01-22
  • 2016-11-25
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多