【发布时间】:2018-04-04 10:58:52
【问题描述】:
我正在尝试让Observable 与所有订阅者共享它的排放量,这样它就会被 subscribe()d 恰好一次。
我尝试使用Observable.publish(),但如果订阅源Observable 完成后订阅,则发布的Observable 的订阅者似乎不会收到任何终止消息(onCompleted() 和可能的onError())。这是一段代码来证明这一点:
static <T> Observer<T> printObserver(String name) {
return new Observer<T>() {
@Override public void onCompleted() {
System.out.println(name + ": onCompleted()");
}
@Override public void onError(Throwable e) {
System.out.println(name + ": onError( " + e + " )");
}
@Override public void onNext(T value) {
System.out.println(name + ": onNext( " + value + " )");
}
};
}
public void testRxPublishConnect() throws Exception {
Observable<Integer> sourceObservable = Observable.range(1, 5);
ConnectableObservable<Integer> sharedObservable = sourceObservable.publish();
sharedObservable.subscribe(printObserver("Observer #1"));
sharedObservable.connect();
sharedObservable.subscribe(printObserver("Observer #2"));
}
这是打印出来的:
观察者 #1:onNext(1) 观察者 #1: onNext( 2 ) 观察者 #1: onNext( 3 ) 观察者 #1: onNext( 4 ) 观察者 #1: onNext( 5 ) 观察者 #1:onCompleted()请注意,Observer #2 不会收到 onCompleted()。
我不认为这是期望的行为。我错过了什么吗?
我在 RxJava 版本 1.0.8 和 1.0.14 中尝试过,结果相同。
【问题讨论】:
-
出于好奇,我为 Rx.Net 编写了相同的代码。它会打印一条“Observer#2: OnCompleted()”消息(除了 Observer#1 的消息),这是我对 publish()/Publish() 的期望。
标签: java rx-java reactive-programming