【发布时间】: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<T>,这里最终调用了 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<T> extends Subscriber<T> implements Action0 -> 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