【问题标题】:Using CompletableFuture within Filter Function在过滤器函数中使用 CompletableFuture
【发布时间】:2018-12-07 07:59:44
【问题描述】:

我有一个用例,我想根据对元素执行的网络调用过滤掉列表中的几个元素。为此,我使用了流、过滤器和 Completable Future。目标是进行异步执行,以使操作变得高效。下面提到了这个的伪代码。

public List<Integer> afterFilteringList(List<Integer> initialList){
   List<Integer> afterFilteringList =initialList.stream().filter(element -> {
        boolean valid = true;
        try{
            valid = makeNetworkCallAndCheck().get();
        } catch (Exception e) {

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

    return afterFilteringList;
}
public CompletableFuture<Boolean> makeNetworkCallAndCheck(Integer value){
   return CompletableFuture.completedFuture(resultOfNetWorkCall(value);
 }

我在这里遇到的问题是,我是否以异步方式本身执行此操作?(当我在过滤器中使用“get”函数时,它会阻止执行并使其仅按顺序执行)或者是否有使用 Java 8 中的 Completable Future 和 Filters 以异步方式执行此操作的更好方法。

【问题讨论】:

  • 您忘记将element 传递给makeNetworkCallAndCheck。此外,例外情况被认为是“有效的”,这看起来很奇怪。
  • @Holger 标题与我认为的实际问题不匹配。 我自己是否以异步方式执行此操作? ...恕我直言,调用不是异步的,gets会阻塞...想法?
  • @nullpointer 好吧,调用get() 会立即破坏异步执行的好处,毫无疑问,但我不知道,鉴于此代码,该建议作为解决方案。例如。输入List&lt;Integer&gt; 在流操作期间神奇地变为List&lt;Long&gt;,最终以List&lt;Integer&gt; 返回。我想,它应该一直是相同的 Integer 对象,但我不想基于假设编写代码……

标签: java lambda java-8 completable-future


【解决方案1】:

当您立即调用get 时,您确实在破坏异步执行的好处。解决方案是先收集所有异步作业,然后再加入。

public List<Integer> afterFilteringList(List<Integer> initialList){
    Map<Integer,CompletableFuture<Boolean>> jobs = initialList.stream()
        .collect(Collectors.toMap(Function.identity(), this::makeNetworkCallAndCheck));
    return initialList.stream()
        .filter(element -> jobs.get(element).join())
        .collect(Collectors.toList());
}
public CompletableFuture<Boolean> makeNetworkCallAndCheck(Integer value){
   return CompletableFuture.supplyAsync(() -> resultOfNetWorkCall(value));
}

当然,makeNetworkCallAndCheck 方法也必须启动真正的异步操作。同步调用方法并返回 completedFuture 是不够的。我在这里提供了一个简单的示例异步操作,但对于 I/O 操作,您可能希望提供自己的 Executor,根据您希望允许的同时连接数进行定制。

【讨论】:

    【解决方案2】:

    如果使用get(),则不会是异步的

    get():如有必要,等待此未来完成,然后返回其结果。

    如果您想以异步方式处理所有请求。你可以使用CompletetableFuture.allOf()

    public List<Integer> filterList(List<Integer> initialList){
        List<Integer> filteredList = Collections.synchronizedList(new ArrayList());
        AtomicInteger atomicInteger = new AtomicInteger(0);
        CompletableFuture[] completableFutures = new CompletableFuture[initialList.size()];
        initialList.forEach(x->{
            completableFutures[atomicInteger.getAndIncrement()] = CompletableFuture
                .runAsync(()->{
                    if(makeNetworkCallAndCheck(x)){
                        filteredList.add(x);
                    }
            });
        });
    
        CompletableFuture.allOf(completableFutures).join();
        return filteredList;
    }
    
    private Boolean makeNetworkCallAndCheck(Integer value){
        // TODO: write the logic;
        return true;
    }
    

    【讨论】:

    • 此代码无法保证结果列表的顺序正确。
    • 这里不需要AtomicInteger,因为它没有与其他线程共享。
    【解决方案3】:

    Collection.parallelStream() 是一种为集合执行异步操作的简单方法。你可以修改你的代码如下:

    public List<Integer> afterFilteringList(List<Integer> initialList){
        List<Integer> afterFilteringList =initialList
                .parallelStream()
                .filter(this::makeNetworkCallAndCheck)
                .collect(Collectors.toList());
    
        return afterFilteringList;
    }
    public Boolean makeNetworkCallAndCheck(Integer value){
        return resultOfNetWorkCall(value);
    }
    

    您可以通过this way 自定义自己的执行器。并且按照this保证结果顺序。

    我已经编写了以下代码来验证我所说的。

    public class  DemoApplication {
        public static void main(String[] args) throws ExecutionException, InterruptedException {
            ForkJoinPool forkJoinPool = new ForkJoinPool(50);
            final List<Integer> integers = new ArrayList<>();
            for (int i = 0; i < 50; i++) {
                integers.add(i);
            }
            long before = System.currentTimeMillis();
            List<Integer> items = forkJoinPool.submit(() ->
                    integers
                            .parallelStream()
                            .filter(it -> {
                                try {
                                    Thread.sleep(10000);
                                } catch (InterruptedException e) {
                                    e.printStackTrace();
                                }
                                return true;
                            })
                            .collect(Collectors.toList()))
                    .get();
            long after = System.currentTimeMillis();
            System.out.println(after - before);
        }
    }
    

    我创建了自己的ForkJoinPool,我需要 10019 毫秒才能并行完成 50 个作业,尽管每个作业花费 10000 毫秒。

    【讨论】:

    • 并行流的缺点是缺乏对执行器的控制,对于 I/O 操作,你通常不想要默认的池并行度,它是为 CPU 核心数量身定制的。
    • @Holger 我之前这么认为,甚至在我的帖子中提到过,直到我看到这篇帖子,stackoverflow.com/a/22269778/5053214
    • 嗯,这是一个没有官方支持的无证副作用。它仅适用于 Fork/Join 池,不适用于任意的 Executor 实现。
    猜你喜欢
    • 2021-11-30
    • 1970-01-01
    • 1970-01-01
    • 2013-02-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-18
    相关资源
    最近更新 更多