【问题标题】:How to detect completion of all contained Observables in Observable<Observable<Object>>如何检测 Observable<Observable<Object>> 中所有包含的 Observable 的完成情况
【发布时间】:2014-01-31 10:01:49
【问题描述】:

我希望我的 api 上的方法返回 Observable> 但我希望该方法中的代码在所有包含的 Observable 完成后知道,以便它可以关闭某些东西。最好的方法是什么?

更明确地说,我是在完成这个方法之后:

public static <T> Observable<Observable<T>> doWhenAllComplete(
        final Observable<Observable<T>> original, Action0 action) {
  ...
}

【问题讨论】:

  • 向我们展示更多代码以更好地了解您的问题
  • 这真的取决于你的 API 方法是如何创建包含的 observables 的。你能从你的方法中发布一些产生Observable&lt;Observable&lt;Object&gt;&gt;的代码吗?然后我们就会知道这些内部可观察对象的来源以及该方法跟踪其完成情况的最佳方式。

标签: java system.reactive rx-java


【解决方案1】:

抱歉,我的答案是在 .NET 中(就像 system.reactive 标记一样);我相信你可以翻译它!

如果您的IObservable&lt;IObservable&lt;Object&gt;&gt;source 给出,则:

source.Merge()
      .Subscribe(_  => {}, /* not interested in onNext */
                 () => /* onCompleted action here, called when all complete */);

注意:如果任何流错误(导致合并流在该点终止),这将崩溃,因此您也可以这样做以吞下各个流上的错误:

source.SelectMany(x => x.Catch(Observable.Empty<Object>()))
      .Subscribe(_  => {}, /* not interested in onNext */
                 () => /* onCompleted action here, called when all complete */);

【讨论】:

  • 您假设 Observable> 可以合并而不会失去意义。不幸的是,这个假设不能成立。
  • 将您的答案视为提示,尽管我认为合并方法将是实现目标的关键
  • 我没有做这样的假设——你只要求知道所有流何时完成,所以这就是我给你的。听起来您可能有一些额外的信息要添加到您的问题中?要做的一件好事是添加解决方案应该通过的一个或多个测试。看你添加了什么,为什么需要通过源码?两个单独的订阅不会这样做吗?它还可以避免副作用。
  • 我明白,两个订阅可以做到这一点,但会破坏我认为在下面的回答中实现的封装。是的,我同意您没有做出这样的假设,谢谢。我可能需要明确表示我希望动作发生时不会对源产生副作用(超出动作本身)。您的回答使我得到了答案,谢谢。
【解决方案2】:

我相信这个方法的实现似乎没有副作用:

public static <T> Observable<Observable<T>> doWhenAllComplete(
        final Observable<Observable<T>> original, final Action0 action) {
    return Observable.create(new OnSubscribeFunc<Observable<T>>() {

        @Override
        public Subscription onSubscribe(Observer<? super Observable<T>> o) {
            ConnectableObservable<Observable<T>> published = original
                    .publish();
            Subscription sub1 = Observable.merge(published)
                    .doOnCompleted(action).subscribe();
            Subscription sub2 = published.subscribe(o);
            Subscription sub3 = published.connect();
            return Subscriptions.from(sub1, sub2, sub3);
        }
    });
}

【讨论】:

    【解决方案3】:

    对我来说,这是可行的:

    bothSources = source1.Cast<Object>().Merge (source2.Cast<Object>());
    

    就我而言,我只需要等待 2 个源,但您可以创建一个函数来接收源列表并合并所有源。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-06-03
      • 1970-01-01
      • 2018-10-06
      • 2020-03-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多