【问题标题】:RxJava: Dynamic set of observablesRxJava:动态的可观察对象集
【发布时间】:2016-10-14 08:57:57
【问题描述】:

我有一个中心班,叫它Central。它可以添加 1-N 个 observables。这些我需要动态添加,然后知道最终的 onComplete() 何时执行。

如何做到这一点?

代码示例:

public class Central {
  public void addObservable(Observable o){
    // Add the new observable to subscriptions
  }
}

更新。

我一直在研究这个。使用@DaveMoten 的答案,我已经接近了。

这是我随意添加一个新的可观察对象的方法,并在它们全部完成时收到通知(伪代码):

class central {
  PublishSubject<Observable<T>> subject = PublishSubject.create();
  int commandCount = 0;
  int concatCount = 0;
  addObservable(Observable newObs){
    commandCount++;
    concatCount++;
    subject.concatMap(o -> 
          o.doOnCompleted(
                  () -> { 
                    concatCount--;
                    LOG.warn("inner completed, concatCount: " + concatCount);
                  })
          ).doOnNext(System.out::println)
            .doOnCompleted(() -> { 
                System.out.println("OUTER completed");
            } ) .subscribe();

    onNext(newObs);
    subject.onCompleted(); // If this is here, it ends prematurely (after first onComplete)
    // if its not, OUTER onComplete never happens

    newObs.subscribe((e) ->{},
    (error) -> {},
    () -> {
        commandCount--;
        LOG.warn("inner end: commandCount: " + commandCount);
    });
  }
}
// ... somewhere else in the app:
Observable ob1 = Observable.just(t1);
addObservable.addObservable(ob1);
// ... and possibly somewhere else:
Observable ob2 = Observable.just(t2, t3);
addObservable(ob2);

// Now, possibly somewhere else:
ob1.onNext(1);
// Another place:
ob2.onNext(2);

日志如下所示:

19:51:16.248 [main] WARN  execute: commandCount: 1
19:51:17.340 [main] WARN  execute: commandCount: 2
19:51:23.290 [main] WARN  execute: commandCount: 3
9:51:26.969 [main] WARN   inner completed, concatCount: 2
19:51:27.004 [main] WARN  inner end: commandCount: 2
19:51:27.008 [main] WARN  inner completed, concatCount: 1
19:51:27.009 [main] WARN  inner end: commandCount: 1
19:51:51.745 [ProcessIoCompletion0] WARN  inner completed, concatCount: 0
19:51:51.750 [ProcessIoCompletion0] WARN  inner completed, concatCount: -1
19:51:51.751 [ProcessIoCompletion0] WARN  inner end: commandCount: 0
19:51:51.752 [ProcessIoCompletion0] WARN  inner completed, concatCount: -2
19:51:51.753 [ProcessIoCompletion0] WARN  inner completed, concatCount: -3

更新:我添加了一些计数器,这表明我不明白 concatMap 发生了什么。您可以看到观察者本身的 lambda 订阅正确地倒计时到 0,但 concatMap oncomplete 下降到 -3! OUTER complete 永远不会发生。

【问题讨论】:

  • 过早地说是什么意思?我没有看到问题。我进行了一个测试,它打印了1, inner completed, 2, 3, inner completed, outer completed。那只是票。
  • @DaveMoten 问题是这是在一个新的可观察对象数量不确定的方法中。我不能这样调用 onComplete ,因为它会导致主题完成,即使新的可观察对象有被添加。更新了示例以尝试更清楚。
  • 没关系。永远不要打电话给subject.onCompleted()
  • 顺便说一句,您不能拨打ob1.onNextob2.onNext 电话。这只能在Subject 上实现。您确定不希望内部 observables 也成为 Subjects(并使用 flatMap 而不是 concatMap)吗?
  • 我正在测试,如果没有调用 subject.onCompleted(),则 OUTER 完成似乎永远不会发生(系统相当复杂,但似乎确实如此!)

标签: java rx-java


【解决方案1】:

使用PublishSubject

PublishSubject<Observable<T>> subject = 
    PublishSubject.create();
subject
    // protect against concurrent calls to subject (optional)
    .serialize()
    .concatMap(o -> 
      o.doOnCompleted(() -> System.out.println("inner completed")))
    .doOnNext(System.out::println)
    .doOnCompleted(() -> System.out.println("completed"))
    .subscribe(subscriber);

subject.onNext(Observable.just(t1));
subject.onNext(Observable.just(t2, t3));
subject.onCompleted();

  .

【讨论】:

  • 这几乎可以满足我的要求...不是的部分,doOnCompleted 只有在我手动使用subject.onCompleted() 调用它时才会触发。当源 observable 触发 onComplete 时,我想获得一个 onComplete。
  • 我几乎可以使用observable.doOnCompleted(subject::onCompleted); 到达那里,但这仅适用于单个可观察源。如何为多个来源做到这一点?
  • concatMap 中添加doOnCompletedo 即可轻松完成。我修改了答案以反映这一点。
  • 我接受了这个,因为我认为就是这样......但我正在寻找的是在所有内部可观察对象都执行 onComplete 时得到通知......这可能吗?
  • 在所有内部可观察对象完成之前,不会调用上面的最终doOnCompleted
猜你喜欢
  • 1970-01-01
  • 2015-01-12
  • 1970-01-01
  • 1970-01-01
  • 2015-10-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-02-06
相关资源
最近更新 更多