【问题标题】:SupplyAsync wait for all CompletableFutures to finishSupplyAsync 等待所有 CompletableFutures 完成
【发布时间】:2020-10-12 18:14:44
【问题描述】:

我正在下面运行一些异步任务,需要等到它们全部完成。我不知道为什么,但join() 没有强制等待所有任务,代码继续执行而无需等待。连接流没有按预期工作是有原因的吗?

CompletableFutures 列表只是一个映射 supplyAsync 的流

List<Integer> items = Arrays.asList(1, 2, 3);

List<CompletableFuture<Integer>> futures = items
                .stream()
                .map(item -> CompletableFuture.supplyAsync(() ->  {

                    System.out.println("processing");
                    // do some processing here
                    return item;

                }))
                .collect(Collectors.toList());

而我等待未来之类的。

CompletableFuture.allOf(futures.toArray(new CompletableFuture[futures.size()]))
                .thenApply(ignored -> futures.stream()
                        .map(CompletableFuture::join)
                        .collect(Collectors.toList()));

我可以使用 futures.forEach(CompletableFuture::join); 等待等待,但我想知道为什么我的流方法不起作用。

【问题讨论】:

  • 您是否尝试过等待CompletableFuture.allOf().thenApply() 的完整未来?您只需创建一个新的未来,而不对其进行任何操作。
  • 这不是我的CompletableFuture.allOf().thenApply() 应该做的吗?它加入了每个应该阻塞的期货,直到所有期货都完成。

标签: java asynchronous java-8 completable-future


【解决方案1】:

这段代码:

CompletableFuture.allOf(futures.toArray(new CompletableFuture[futures.size()]))
                .thenApply(ignored -> futures.stream()
                        .map(CompletableFuture::join)
                        .collect(Collectors.toList()));

是否等待futures 中的所有期货完成。它所做的是创建一个新的未来,它将等待futures 中的所有异步执行在它自己完成之前完成(但在所有这些未来完成之前不会阻塞)。当这个 allOf 未来完成时,您的 thenApply 代码就会运行。但是allOf() 会立即返回而不会阻塞。

这意味着您的代码中的futures.stream().map(CompletableFuture::join).collect(Collectors.toList()) 仅在所有异步执行完成后运行,这与您的目的相违背。 join() 调用将立即返回。但这不是更大的问题。您的挑战是 allOf().thenApply() 不会等待异步执行完成。它只会创造另一个不会阻塞的未来。

最简单的解决方案是使用第一个管道并映射到一个整数列表:

List<Integer> results = items.stream()
    .map(item -> CompletableFuture.supplyAsync(() -> {

        System.out.println("processing " + item);
        // do some processing here
        return item;

    }))
    .collect(Collectors.toList()) //force-submit all
    .stream()
    .map(CompletableFuture::join) //wait for each
    .collect(Collectors.toList());

如果您想使用类似于原始代码的内容,则必须将您的第二个 sn-p 更改为:

List<Integer> reuslts = futures.stream()
    .map(CompletableFuture::join)
    .collect(Collectors.toList());

那是因为CompletableFuture.allOf 不等待,它只是将所有期货组合成一个新的期货,当全部完成时完成:

返回一个新的 CompletableFuture,当所有给定的 CompletableFuture 都完成时,它就完成了。

或者,您仍然可以将allOf()join() 一起使用,然后运行您当前的thenApply() 代码:

//wrapper future completes when all futures have completed
CompletableFuture.allOf(futures.toArray(new CompletableFuture[futures.size()]))
        .join(); 

//join() calls below return immediately
List<Integer> result = futures.stream()
        .map(CompletableFuture::join) 
        .collect(Collectors.toList());

最后一条语句中的 join() 调用立即返回,因为包装器 (allOf()) 上的 join() 调用将等待传递给它的所有期货完成。这就是为什么当您可以使用第一种方法时,我看不到这样做的原因。

【讨论】:

  • 我不明白,.allOf(…).thenApply(…) 方法是如何违背目的的。唯一缺少的是链接另一个join() 以获取结果列表(这意味着等待它)。
  • @Holger 在allOf 之后传递给thenApply() 的回调不会运行,直到包裹在allOf() 中的所有期货都完成。这意味着在该回调中对那些期货调用join() 就强制系统等待结果而言没有任何意义(所有join() 是否从已经完成的期货中获取结果,它不会导致等待) .我意识到,这句话可以更好地表达,但这就是我的意思。我已经建议了您在上一个 sn-p 中建议的 join() 调用,但我只是不喜欢对该方法的双重调用。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2018-09-22
  • 1970-01-01
  • 1970-01-01
  • 2015-05-27
  • 2010-09-20
  • 2021-04-14
  • 2022-01-23
相关资源
最近更新 更多