【问题标题】:RxJava Combining Multiple Observer after filterRxJava在过滤后组合多个观察者
【发布时间】:2018-11-09 15:57:28
【问题描述】:

以下是我当前的代码

  private final List<Disposable> subscriptions = new ArrayList<>();

  for (Instrument instrument : instruments) {
    // Waiting for OrderBook to generate Reliable results.
    GenericBook Book =
        service
            .getBook(instrument.getData())
            .filter(gob -> onBookUpdate(gob))
            .blockingFirst();

    subscriptions.add(
        service
            .getBook(instrument.getData())
            .subscribe(
                gob -> {
                  try {
                    onBookUpdate(gob);
                  } catch (Exception e) {
                    logger.error("Error on subscription:", e);
                  }
                },
                e -> logger.error("Error on subscription:", e)));
  }

所以它的作用是对它首先阻塞的每个仪器等待,直到onBookUpdate(gob) 的输出变为真。 onBookUpdate(gob) 返回布尔值。 一旦我们将第一个 onBookUpdate 设置为 true,我就会将该订阅者推送到订阅变量中。

这会变慢,因为我必须等待每个乐器,然后再继续下一个乐器。

我的目标是并行运行所有这些,然后等待所有完成并将它们推送到订阅变量。

我试过 zip 但没用

  List<Observable<GenericOrderBook>> obsList = null;
  for (Instrument instrument : instruments) {
    // This throws nullException.
   obsList.add(service
            .getBook(instrument.getData())
            .filter(gob -> onBookUpdate(gob))
            .take(1));
    }
  }
// Some how wait over here until all get first onBookUpdate as true.
String o = Observable.zip(obsList, (i) -> i[0]).blockingLast();

【问题讨论】:

    标签: java rx-java observable rx-java2


    【解决方案1】:

    当使用 observables 等时,应该全心全意地拥抱它们。拥抱的前提之一是将管道的配置和构建与其执行分开。

    换句话说,先配置您的管道,然后在数据可用时通过它发送数据。

    此外,拥抱 observables 意味着避免 for 循环。

    我不是 100% 你的用例是什么,但我的建议是创建一个将仪器作为输入并返回订阅的管道......

    类似

    service.getBook(instrument.getData())
     .flatMap(gob -> {
       onBookUpdate(gob);
       return gob;
    });
    

    这将返回一个Observable,您可以订阅它并将结果添加到订阅中。

    然后创建一个可观察的种子,将仪器对象泵入其中。

    不确定您的 API 的某些细节,如果不清楚或者我做出了错误的假设,请与我联系。

    【讨论】:

      【解决方案2】:

      我假设instruments 是一个列表。如果是,那么您可以这样做,

      Observable
          .fromIterable(instruments)
          // Returns item from instrument list one by one and passes it to getBook()
          .flatmap(
              instrument -> getBook(instrument.getData())
          )
          .filter(
              gob -> onBookUpdate(gob)
          )
          // onComplete will be called if no items from filter 
          .switchIfEmpty(Observable.empty())
          .subscribe(
              onBookUpdateResponse -> // Do what you want,
              error -> new Throwable(error)
          );
      

      希望这会有所帮助。

      【讨论】:

        猜你喜欢
        • 2018-06-11
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-04-13
        • 1970-01-01
        • 2016-09-14
        相关资源
        最近更新 更多