【问题标题】:RxJava - Observers don't get onNext() items after a whileRxJava - 一段时间后观察者不会获得 onNext() 项目
【发布时间】:2016-06-08 21:49:24
【问题描述】:

我有一个由多个相互连接的组件组成的系统。一段时间内一切正常,但经过一段时间后,一些观察者停止接收可观察对象的 onNext() 发送的项目。

一个简化的场景是这样的:我有 Component1.start() -> 使用Observable.create(...).subscribeOn().observeOn().publish() 创建一个 ConnectableObservable,并订阅 Component2。之后,它连接()s。这个 observable 在循环中发出一些项目,然后在完成时调用 s.onComplete()

Component2 实现了观察者。此外,它还有一个运行 while(true) 循环的 ConnectableObservable。当它在由 Component1 调用的 onNext() 中获得一个值时,它会使用自己的 ConnectableObservable 通知 Component0。 (注意,我还使用 PublishSubject 实现了它们,并且发生了同样的情况)。

Component1.start() //Creates Component1's ConnectableObservable, subscribes Component2 and starts running with connect();

Component1.connectableObservable -> onNext() ---> Component2 

Component2.connectableObservable -> onNext() ---> Component0

当 Component0.onNext() 获得一个特定项目时(经过 100 次迭代),它会停止 Component1.observable,使其退出循环并调用 onComplete()。

一段时间后,Component0 调用 Component1.start(),一切重新开始。

我看到的是,一切正常 Component1.observable.onNext() 调用rx.internal.operators.OperatorSubscribeOn.......subscriber.onNext()

rx.internal.operators.OperatorSubscribeOn

@Override
public void call(final Subscriber<? super T> subscriber) {
    final Worker inner = scheduler.createWorker();
    subscriber.add(inner);

    inner.schedule(new Action0() {
        @Override
        public void call() {
            final Thread t = Thread.currentThread();

            Subscriber<T> s = new Subscriber<T>(subscriber) {
                @Override
                public void onNext(T t) {
                    subscriber.onNext(t);

subscriber.onNext() 是内部类private static final class ObserveOnSubscriber&lt;T&gt;,这里最终调用了 schedule():

@Override
public void onNext(final T t) {
    if (isUnsubscribed() || finished) {
        return;
    }
    if (!queue.offer(on.next(t))) {
        onError(new MissingBackpressureException());
        return;
    }
    schedule();
}

schedule() 是

protected void schedule() {
    if (counter.getAndIncrement() == 0) {
        recursiveScheduler.schedule(this);
    }
}

counter 为 0,所以 recursiveScheduler.schedule(this); 被调用并且 Component2 获取项目。

现在,当它停止工作时,计数器不再为 0,实际上每次调用都会增加它。因此,recursiveScheduler.schedule(this); 永远不会被调用,而 Component2 也不会得到任何东西。

这可能是什么原因?为什么计数器 0 并且在某些时候开始增加?

更新:在源代码中挖掘我看到了以下内容:在调用 schedule() 之后,有一个计划任务调用下面的代码,当项目还没有时减少计数器错过:

private static final class ObserveOnSubscriber<T> extends Subscriber<T> implements Action0 {

    // only execute this from schedule()
    @Override
    public void call() {
        ...
            emitted = currentEmission;
            missed = counter.addAndGet(-missed);
            if (missed == 0L) {
                break;
            }

据此,由于项目丢失,计数器增加,然后后续项目也丢失。

遗漏物品的原因可能是什么?

我注意到了一些奇怪的事情。如果我从程序中删除任何其他(未提及的)可观察项,则不会遗漏任何项目。他们将 Component0 作为观察者并在他们自己的 subscribeOn() 线程中生成他们的项目,所以我看不出它们如何影响这种情况。

更新 2:我一直试图找出发生了什么。当我执行Component1.connectableObservable.connect() 时,它最终会调用private static final class ObserveOnSubscriber&lt;T&gt; extends Subscriber&lt;T&gt; implements Action0 -&gt; init()

这里调用了 schedule():

void init() {
            // don't want this code in the constructor because `this` can escape through the 
            // setProducer call
            Subscriber<? super T> localChild = child;

            localChild.setProducer(new Producer() {

                @Override
                public void request(long n) {
                    if (n > 0L) {
                        BackpressureUtils.getAndAddRequest(requested, n);
                        schedule();

在 schedule() 之后,正确的行为使 OperatorObserveOn.counter = 0。当它不再工作时,scheduler() 将 OperatorObserveOn.counter 的值增加 +1。

【问题讨论】:

    标签: java rx-java reactive-programming


    【解决方案1】:

    获取每个订阅的观察者并创建一个侦听器以告诉您何时有一个观察者取消订阅可观察对象,然后您就可以理解为什么会发生这种情况。

    无论如何,在你的情况下,我会看看 Relay,因为你不必取消订阅你的 observable,它更安全,而且你可以确定它永远不会停止发射事件。 p>

    看看这个例子。

    https://github.com/politrons/reactive/blob/master/src/test/java/rx/relay/Relay.java

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2014-10-10
    • 1970-01-01
    • 2012-05-11
    • 1970-01-01
    • 2018-06-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多