【问题标题】:Observable.publish() doesn't call onCompleted() on observers that subscribe after the source Observable is doneObservable.publish() 不会在源 Observable 完成后订阅的观察者上调用 onCompleted()
【发布时间】: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


【解决方案1】:

试试.share(),即.publish().refCount()。

【讨论】:

  • 这是我的第一个方法。问题是,在所有订阅者取消订阅后(或在源 Observable 完成后,它会自动取消订阅所有订阅者)并且新订阅者订阅后,源 Observable 将被再次订阅。
  • 明确地说,使用 PublishSubject 解决了我的问题。我只是觉得 .publish() 的行为很奇怪。
  • Righto,请原谅我没有完整阅读您的问题,我只是在空闲时间把它扔掉了。我认为您应该更新您的 RxJava 版本,因为自该版本以来已经进行了许多重要的修复(我认为现在是 1.0.14)。
  • 我也试过 1.0.14。结果是一样的。我将更新问题以反映它。
【解决方案2】:

这是设计使然。如果您在这种情况下调用connect(),您的订阅者将从一开始就收到所有事件。如果终止的publish 将立即终止其子订阅者,您可能无法观察到值,因为一旦连接,publish 会在没有订阅者的情况下缓慢地标记其源。

【讨论】:

  • 你的意思是 Observable 是 .publish() 的结果,如果源 Observable 完成,不调用 onError() 或 onNext() 立即取消订阅它的订阅者?您能否详细说明“如果没有订阅者,发布会慢慢删除其来源”是什么意思?
  • 换句话说,一旦发布开始并且所有订阅者都消失了,它就会删除源值。
【解决方案3】:

我 99% 确定这是预期的行为。我不确定 RxJava,但在我所知道的发布和订阅模式的大多数实现中,可观察对象的默认行为是将事件发布给订阅者并忘记它们。这意味着通知不是“追溯的”(即订阅者对过去发出的事件一无所知)。

另外,来自Observable Contract (section 'multiple observers') of the RxJava documentation:

If a second observer subscribes to an Observable that is already emitting items to a first observer, it is up to the Observable whether it will thenceforth emit the same items to each observer ... There is no general guarantee that two observers of the same Observable will see the same sequence of items.

【讨论】:

  • 关于 onNext() 通知确实如此。这不是(据我所知)关于 onError() 和 onCompleted() 通知的真实情况。为了确保这一点,我检查了 PublishSubject,它总是向订阅者发出 onCompleted() 通知,即使很久以前调用了 onCompleted()。
【解决方案4】:

Publish 通过构建所有订阅者的列表来工作,然后一旦调用connect(),它就会开始向其订阅者列表中的所有订阅者生成数据。这意味着在调用 connect 之前必须知道所有订阅者。这是您将如何使用publish() 或更优选publish(Func1&lt;Observable&lt;T&gt;, Observable&lt;R&gt;&gt;) 重载的方法。

已知订阅者数量:发布

关闭所有订阅的函数。

observableStream.publish(new Func1<Observable<Integer>, Observable<Integer>>() {
    @Override
    public Observable<Integer> call(Observable<Integer> subject) {
        Observable<Integer> o1 = subject.doOnNext(somework1());
        Observable<Integer> o2 = subject.doOnNext(somework2());
        return Observable.merge(o1, o2);
    }
});

手动调用连接和订阅:

ConnectableObservable<Integer> subject = observableStream.publish();
subject.subscribe(somework1());
subject.subscribe(somework2());
subject.connect();

如果您不知道您将拥有多少订阅者,那么您可以将输入窗口化为可管理的块,然后在您的 Transformers 集合上发布您的输入。

订阅者数量未知:窗口

final Set<Transformer<Integer, String>> transformers = new HashSet<>();
observableStream
    .window(100, TimeUnit.MILLISECONDS, 1000)
    .flatMap(new Func1<Observable<Integer>, Observable<String>>(){
        @Override
        public Observable<String> call(Observable<Integer> window) {
            return window.publish(new Func1<Observable<Integer>, Observable<String>>() {
                @Override
                public Observable<String> call(Observable<Integer> publish) {
                    Observable<Observable<String>> workObservables = Observable.from(transformers)
                        .map(new Func1<Transformer<Integer, String>, Observable<String>>(){
                            @Override
                            public Observable<String> call(Transformer<Integer, String> transformer) {
                                return publish.compose(transformer);
                            }});
                    return Observable.merge(workObservables);
                }});
        }})
    .subscribe();

还有第三种选择。您可以使用observable.cache(),但这会将来自该可观察流的所有输入数据保存在内存中,因此您要小心使用它。在这种情况下,您可能最终还是会通过窗口来控制缓存主题的范围。

【讨论】:

  • 我不想要任何缓存或窗口或任何东西。我的期望是,如果 Observable 不再发出任何项目,它应该向其订阅者发送终止通知。换句话说,Observable 应该可以跳过 onNext() 通知,但不能跳过 onError() 或 onCompleted() 通知。
  • Francisco 的回答是正确的,因为发布主题将事件流式传输给子订阅者。如果在连接之后发生订阅,那么它将不会收到任何事件(包括终端信号)。鉴于在连接之前必须订阅所有观察者的限制,我的回答试图为您提供使用这些 API 的方法。
  • 这样说:'如果在连接之后发生订阅,那么它将不会收到任何事件(包括终端信号)'你的意思是在我的问题中的代码段中吗?因为如果 observable 在订阅后发出项目,观察者确实会收到它们。我用 Observable.interval() 对其进行了测试。
  • 而 publish() 的文档清楚地表明了其他情况(看一下图表):reactivex.io/documentation/operators/publish.html
  • 您链接的弹珠图显示两个观察者在连接之前订阅。欢迎您提交关于该项目的 github 问题,但我相信这是故意行为。 Observable.interval 实际上是一个无限的可观察对象,并且不会自行完成。它们有非常不同的语义。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-04-20
  • 1970-01-01
  • 2017-12-28
  • 2020-08-10
  • 1970-01-01
  • 1970-01-01
  • 2018-06-11
相关资源
最近更新 更多