从 List<CompletableFuture> 转换为 CompletableFuture<List>

新手上路,请多包涵

我正在尝试将 List<CompletableFuture<X>> 转换为 CompletableFuture<List<T>> 。当您有许多异步任务并且您需要获得所有这些任务的结果时,这非常有用。

如果其中任何一个失败,那么最终的未来就会失败。这就是我的实施方式:

 public static <T> CompletableFuture<List<T>> sequence2(List<CompletableFuture<T>> com, ExecutorService exec) {
    if(com.isEmpty()){
        throw new IllegalArgumentException();
    }
    Stream<? extends CompletableFuture<T>> stream = com.stream();
    CompletableFuture<List<T>> init = CompletableFuture.completedFuture(new ArrayList<T>());
    return stream.reduce(init, (ls, fut) -> ls.thenComposeAsync(x -> fut.thenApplyAsync(y -> {
        x.add(y);
        return x;
    },exec),exec), (a, b) -> a.thenCombineAsync(b,(ls1,ls2)-> {
        ls1.addAll(ls2);
        return ls1;
    },exec));
}

运行它:

 ExecutorService executorService = Executors.newCachedThreadPool();
Stream<CompletableFuture<Integer>> que = IntStream.range(0,100000).boxed().map(x -> CompletableFuture.supplyAsync(() -> {
    try {
        Thread.sleep((long) (Math.random() * 10));
    } catch (InterruptedException e) {
        e.printStackTrace();
    }
    return x;
}, executorService));
CompletableFuture<List<Integer>> sequence = sequence2(que.collect(Collectors.toList()), executorService);

如果其中任何一个失败,那么它就失败了。即使有一百万个期货,它也会按预期提供输出。我遇到的问题是:假设有超过 5000 个期货并且其中任何一个失败,我得到一个 StackOverflowError

java.util.concurrent.CompletableFuture.internalComplete(CompletableFuture.java:210) 在 java.util.concurrent.CompletableFutureThenCompose.run(CompletableFuture.java)线程“pool1thread2611”中的异常java.lang.StackOverflowError:1487)java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:193)java.util.concurrent.CompletableFuture.internalComplete(CompletableFuture.java:210)java.util.concurrent.CompletableFutureThenCompose.run(CompletableFuture.java) 线程“pool-1-thread-2611”中的异常 java.lang.StackOverflowError :1487) 在 java.util.concurrent.CompletableFuture.postComplete(CompletableFuture.java:193) 在 java.util.concurrent.CompletableFuture.internalComplete(CompletableFuture.java:210) 在 java.util.concurrent.CompletableFutureThenCompose.run( CompletableFuture.java:1487)

我做错了什么?

注意:当任何一个未来失败时,上述返回的未来就会失败。接受的答案也应该考虑这一点。

原文由 Jatin 发布,翻译遵循 CC BY-SA 4.0 许可协议

阅读 1.5k
2 个回答

使用 CompletableFuture.allOf(...)

 static<T> CompletableFuture<List<T>> sequence(List<CompletableFuture<T>> com) {
    return CompletableFuture.allOf(com.toArray(new CompletableFuture<?>[0]))
            .thenApply(v -> com.stream()
                .map(CompletableFuture::join)
                .collect(Collectors.toList())
            );
}

关于您的实施的一些评论:

您对 .thenComposeAsync.thenApplyAsync.thenCombineAsync 的使用可能没有达到您的预期。这些 ...Async 方法在单独的线程中运行提供给它们的函数。因此,在您的情况下,您导致将新项目添加到列表中以在提供的执行程序中运行。无需将轻量级操作塞入缓存的线程执行器中。不要在没有充分理由的情况下使用 thenXXXXAsync 方法。

此外,不应使用 reduce 累积到可变容器中。即使流是顺序的时它可能会正常工作,但如果流是并行的,它就会失败。要执行可变缩减,请改用 .collect

如果您想在第一次失败后立即异常地完成整个计算,请在您的 sequence 方法中执行以下操作:

 CompletableFuture<List<T>> result = CompletableFuture.allOf(com.toArray(new CompletableFuture<?>[0]))
        .thenApply(v -> com.stream()
                .map(CompletableFuture::join)
                .collect(Collectors.toList())
        );

com.forEach(f -> f.whenComplete((t, ex) -> {
    if (ex != null) {
        result.completeExceptionally(ex);
    }
}));

return result;

此外,如果您想在第一次失败时取消其余操作,请在 exec.shutdownNow(); result.completeExceptionally(ex); 。当然,这假设 exec 只存在于这一次计算中。如果没有,您将不得不循环并分别取消每个剩余的 Future

原文由 Misha 发布,翻译遵循 CC BY-SA 4.0 许可协议

您可以获得 Spotify 的 CompletableFutures 库并使用 allAsList 方法。我认为它的灵感来自番石榴的 Futures.allAsList 方法。

 public static <T> CompletableFuture<List<T>> allAsList(
    List<? extends CompletionStage<? extends T>> stages) {


如果您不想使用库,这里有一个简单的实现:

 public <T> CompletableFuture<List<T>> allAsList(final List<CompletableFuture<T>> futures) {
    return CompletableFuture.allOf(
        futures.toArray(new CompletableFuture[futures.size()])
    ).thenApply(ignored ->
        futures.stream().map(CompletableFuture::join).collect(Collectors.toList())
    );
}

原文由 oskansavli 发布,翻译遵循 CC BY-SA 3.0 许可协议

撰写回答
你尚未登录,登录后可以
  • 和开发者交流问题的细节
  • 关注并接收问题和回答的更新提醒
  • 参与内容的编辑和改进,让解决方法与时俱进
推荐问题