【问题标题】:Wait for multiple async calls to finish in RxJava在 RxJava 中等待多个异步调用完成
【发布时间】:2017-01-23 20:41:52
【问题描述】:

我的场景很简单,但我似乎无法在任何地方找到它。

我有一组要迭代的元素,每个元素调用一个异步函数,然后等待所有元素完成(这又以异步方式发生,在函数逻辑中实现) .我对 RxJava 比较陌生,过去在 NodeJS 中通过将回调传递给函数并在最后等待很容易做到这一点。 这是我需要的伪代码(元素的迭代器不需要同步也不需要排序):

for(line in lines){
 callAsyncFunction(line);
}
WAIT FOR ALL OF THEM TO FINISH

非常感谢您的帮助!

【问题讨论】:

  • 请发布实际代码。 RxJava 有很多具有不同行为的运算符。不知道你使用什么运算符,是不可能提供帮助的。
  • 使用 countdownlatch 类,并在每个可观察到的调用 countDown 方法的 oncomplete 方法中。

标签: java asynchronous rx-java reactive-programming


【解决方案1】:

使用 defer 将“要计算的项目”的 Iterable 转换为 Observable 的 Iterable,然后在 Observable 上使用 zip

更困难的部分将是“等待他们全部完成”。正如您可能已经猜到的那样,响应式扩展是关于“对事物做出反应”,而不是“等待事物发生”。你可以 subscribe 到 observable,如果每个 observable 只有一个 item,它将发出一个 item,然后完成。该订阅者可以在等待后执行您通常会执行的任何操作;这允许您返回并让您的代码做它的事情而不会阻塞它。

【讨论】:

    【解决方案2】:

    从技术上讲,如果您考虑一下,您需要做的是从您的所有元素创建一个 Observable,然后 zip them together 继续执行您的流。

    在伪代码中会给你这样的东西:

    List<Observable<?>> observables = new ArrayList<>();
    for(line in lines){
       observables.add(Observable.fromCallable(callAsyncFunction(line));
    }
    Observable.zip(observables, new Function<...>() { ... }); // kinda like Promise.all()
    

    但是Observable.from() 可以将可迭代对象中的每个元素公开为对象流,这也就不足为奇了,从而消除了对循环的需要。因此,您可以使用Observable.fromCallable() 创建一个在异步操作完成时调用onCompleted() 的新Observable。之后,您可以通过将这些新的 Observables 收集到一个列表中来等待它们。

    Observable.from(lines)
       .flatMap(new Func1<String, Observable<?>>() {
            @Override
            public Observable<?> call(String line) {
                return Observable.fromCallable(callAsyncFunction(line)); // returns Callable
            }
        }).toList()
          .map(new Func1<List<Object>, Object>() {
            @Override
            public Object call(List<Object> ignored) {
                // do something;
            }
        });
    

    我的答案的后半部分主要基于this answer

    【讨论】:

    • 非常感谢。对我来说效果很好。
    【解决方案3】:

    使用 Rx:

    Observable
    .from(lines)
    .flatMap(line -> callAsyncFunctionThatReturnsObservable(line).subscribeOn(Schedulers.io())
    .ignoreElements();
    

    此时,根据您想要做什么,您可以使用 .switchIfEmpty(...) 订阅另一个 observable。

    【讨论】:

    • 不错!超级优雅
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-02-15
    • 1970-01-01
    • 1970-01-01
    • 2019-02-10
    • 2012-05-04
    • 2019-10-30
    • 1970-01-01
    相关资源
    最近更新 更多