【问题标题】:RxJava: how to compose multiple Observables with dependencies and collect all results at the end?RxJava:如何组合多个具有依赖关系的 Observable 并在最后收集所有结果?
【发布时间】:2014-03-07 02:39:18
【问题描述】:

我正在学习 RxJava,作为我的第一个实验,我尝试重写 this code 中第一个 run() 方法中的代码(在 Netflix's blog 上引用为 RxJava 可以帮助解决的问题)以提高其异步性RxJava,即它不会等待第一个 Future (f1.get()) 的结果,然后继续执行其余代码。

f3 依赖于f1。我知道如何处理这个问题,flatMap 似乎可以解决问题:

Observable<String> f3Observable = Observable.from(executor.submit(new CallToRemoteServiceA()))
    .flatMap(new Func1<String, Observable<String>>() {
        @Override
        public Observable<String> call(String s) {
            return Observable.from(executor.submit(new CallToRemoteServiceC(s)));
        }
    });

接下来,f4f5 依赖于 f2。我有这个:

final Observable<Integer> f4And5Observable = Observable.from(executor.submit(new CallToRemoteServiceB()))
    .flatMap(new Func1<Integer, Observable<Integer>>() {
        @Override
        public Observable<Integer> call(Integer i) {
            Observable<Integer> f4Observable = Observable.from(executor.submit(new CallToRemoteServiceD(i)));
            Observable<Integer> f5Observable = Observable.from(executor.submit(new CallToRemoteServiceE(i)));
            return Observable.merge(f4Observable, f5Observable);
        }
    });

这开始变得很奇怪(mergeing 他们可能不是我想要的......)但最终允许我这样做,而不是我想要的:

f3Observable.subscribe(new Action1<String>() {
    @Override
    public void call(String s) {
        System.out.println("Observed from f3: " + s);
        f4And5Observable.subscribe(new Action1<Integer>() {
            @Override
            public void call(Integer i) {
                System.out.println("Observed from f4 and f5: " + i);
            }
        });
    }
});

这给了我:

Observed from f3: responseB_responseA
Observed from f4 and f5: 140
Observed from f4 and f5: 5100

这是所有数字,但不幸的是我在单独的调用中得到了结果,所以我不能完全替换原始代码中的最终 println:

System.out.println(f3.get() + " => " + (f4.get() * f5.get()));

我不明白如何在同一行访问这两个返回值。我认为这里可能缺少一些函数式编程。我怎样才能做到这一点?谢谢。

【问题讨论】:

  • 你想要toList吗?

标签: java functional-programming reactive-programming rx-java


【解决方案1】:

看起来你真正需要的只是多一点鼓励和对如何使用 RX 的看法。我建议您更多地阅读文档以及大理石图(我知道它们并不总是有用的)。我还建议查看lift() 函数和运算符。

  • observable 的全部意义在于将数据流和数据操作连接到一个对象中
  • 调用mapflatMapfilter 的目的是为了操作数据流中的数据
  • 合并的重点是合并数据流
  • 操作符的目的是让您能够中断稳定的可观察数据流并定义您自己对数据流的操作。例如,我编写了一个移动平均运算符。这总结了 n doubles 在双打的 Observable 以返回移动平均线流。代码看起来像这样

    Observable movingAverage = Observable.from(mDoublesArray).lift(new MovingAverageOperator(frameSize))

许多您认为理所当然的过滤方法都有lift(),这让您松了一口气。

话虽如此;合并多个依赖项所需要的只是:

  • 使用mapflatMap 将所有传入数据更改为标准数据类型
  • 将标准数据类型合并到流中
  • 如果一个对象需要等待另一个对象,或者如果您需要对流中的数据进行排序,则使用自定义运算符。注意:这种方法会减慢流速度
  • 用于列出或订阅以收集所有数据

【讨论】:

  • 感谢您的回复。我希望有一天能尽快回到使用 RxJava。
【解决方案2】:

编辑:有人将我作为问题编辑添加的以下文本转换为答案,我对此表示赞赏,并且理解这可能是正确的做法,但是 我不认为这是一个答案,因为这显然不是正确的方法。我永远不会使用此代码,也不会建议任何人复制它。 欢迎其他/更好的解决方案和 cmets!


我能够通过以下方式解决此问题。我没有意识到你可以不止一次地flatMap 一个 observable,我认为结果只能被消费一次。所以我只是 flatMap f2Observable 两次(对不起,我在我的原始帖子之后重命名了代码中的一些东西),然后在所有 Observables 上使用 zip,然后订阅它。由于类型杂耍,zip 中的 Map 聚合值是不可取的。 欢迎其他/更好的解决方案和 cmets!full code is viewable in a gist。谢谢。

Future<Integer> f2 = executor.submit(new CallToRemoteServiceB());
Observable<Integer> f2Observable = Observable.from(f2);
Observable<Integer> f4Observable = f2Observable
    .flatMap(new Func1<Integer, Observable<Integer>>() {
        @Override
        public Observable<Integer> call(Integer integer) {
            System.out.println("Observed from f2: " + integer);
            Future<Integer> f4 = executor.submit(new CallToRemoteServiceD(integer));
            return Observable.from(f4);
        }       
    });     

Observable<Integer> f5Observable = f2Observable
    .flatMap(new Func1<Integer, Observable<Integer>>() {
        @Override
        public Observable<Integer> call(Integer integer) {
            System.out.println("Observed from f2: " + integer);
            Future<Integer> f5 = executor.submit(new CallToRemoteServiceE(integer));
            return Observable.from(f5);
        }       
    });     

Observable.zip(f3Observable, f4Observable, f5Observable, new Func3<String, Integer, Integer, Map<String, String>>() {
    @Override
    public Map<String, String> call(String s, Integer integer, Integer integer2) {
        Map<String, String> map = new HashMap<String, String>();
        map.put("f3", s);
        map.put("f4", String.valueOf(integer));
        map.put("f5", String.valueOf(integer2));
        return map;
    }       
}).subscribe(new Action1<Map<String, String>>() {
    @Override
    public void call(Map<String, String> map) {
        System.out.println(map.get("f3") + " => " + (Integer.valueOf(map.get("f4")) * Integer.valueOf(map.get("f5"))));
    }       
});     

这会产生我想要的输出:

responseB_responseA => 714000

【讨论】:

  • 我向你学习这项技术!它运作良好,但我想知道是否有一种“更清洁”的方式来处理这种常见情况。
  • 我有 2 个 observable 相互依赖,observable 2 依赖于 observable 1,所以 observable 1 应该在 observable 2 之前执行,然后我需要结合两个 observable 的结果。我们有这方面的运营商吗?可以 zip 完成这项工作。我不想使用 flatMap,因为它将一个流转换为另一个流,但在这里我需要设置依赖关系,然后压缩结果。请回复。
【解决方案3】:

我认为您正在寻找的是 switchmap。我们遇到了类似的问题,我们有一个会话服务来处理从 api 获取新会话,我们需要该会话才能获取更多数据。我们可以添加到返回 sessionToken 的 session observable 中,以便在我们的数据调用中使用。

getSession 返回一个 observable;

public getSession(): Observable<any>{
  if (this.sessionToken)
    return Observable.of(this.sessionToken);
  else if(this.sessionObservable)
    return this.sessionObservable;
  else {
    // simulate http call 
    this.sessionObservable = Observable.of(this.sessonTokenResponse)
    .map(res => {
      this.sessionObservable = null;
      return res.headers["X-Session-Token"];
    })
    .delay(500)
    .share();
    return this.sessionObservable;
  }
}

getData 获取该 observable 并附加到它。

public getData() {
  if (this.dataObservable)
    return this.dataObservable;
  else {
    this.dataObservable = this.sessionService.getSession()
      .switchMap((sessionToken:string, index:number) =>{
        //simulate data http call that needed sessionToken
          return Observable.of(this.dataResponse)
          .map(res => {
            this.dataObservable = null;
            return res.body;
          })
          .delay(1200)
        })
        .map ( data => {
          return data;
        })
        .catch(err => {
          console.log("err in data service", err);
         // return err;
        })
        .share();
    return this.dataObservable;
  }
}

你仍然需要一个平面图来组合不依赖的 observables。

Plunkr:http://plnkr.co/edit/hiA1jP?p=info

我从哪里得到了使用 switch 地图的想法:http://blog.thoughtram.io/angular/2016/01/06/taking-advantage-of-observables-in-angular2.html

【讨论】:

  • 您好,感谢您的回答!自从我上次看 RxJava 已经很久了,我真的不知道我是否应该把它作为公认的答案。不过绝对赞成。
猜你喜欢
  • 2014-04-12
  • 2016-08-06
  • 1970-01-01
  • 1970-01-01
  • 2018-07-23
  • 1970-01-01
  • 2021-09-13
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多