【问题标题】:Filter and foreach java 8 stream as parallel并行过滤和 foreach java 8 流
【发布时间】:2021-04-04 11:40:22
【问题描述】:

我正在尝试遍历对象列表。但是我在下面的代码中尝试了太长时间,我想通过增加线程池使其现代和可扩展。

result.getObjectSummaries().parallelStream().forEach((objectSummary) -> {
                    if (objectSummary.getLastModified().after(dateBeforeMp)
                            && objectSummary.getLastModified().before(threadStartDate)) {
                        String correlationId = UUID.randomUUID().toString() + "js" + UUID.randomUUID().toString();
                        if(1==loadFileToSnowFlake(appId, tenantId, objectSummary.getKey(), correlationId)) {
                            logger.info("File loaded for the tenantId {} and fileName: {} ", tenantId, objectSummary.getKey());
                            sendMessageToSQS(appId, tenantId, bucketName, objectSummary.getKey(), correlationId);
                        }
                    }
                    count.incrementAndGet();
                });

我没有找到合适的方法来做到这一点。到目前为止,它运作良好。我的意图是 s3Object 在单个请求中列出 1000 个文件,我想将作业分配给 for 循环内的每个线程。

【问题讨论】:

    标签: spring-boot java-8 stream aws-java-sdk


    【解决方案1】:

    Java 流 parallelStream() 使用内部的、固定大小的线程池,您无法随时修改(可能根本不会)

    您可以利用 CompletableFuture api,并使用自定义线程池来执行您的任务,例如:

    // Different flavors of thread pools available
    Executor executor = Executors.newFixedThreadPool(100);
    Stream.of(1000, 2000, 3000)
        .forEach(i -> CompletableFuture.runAsync(() -> {
            // Long running tasks
            try {
                Thread.sleep(i);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
             // Result
            System.out.println(i);
        }, executor));
    

    在上面,流将在单个线程中迭代(对于大约 100,000 个项目就足够了,对于更大的批次,您可以添加.parallelStream()),并在指定的线程中执行实际工作在执行器中。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-05-14
      • 1970-01-01
      • 1970-01-01
      • 2016-02-28
      • 1970-01-01
      相关资源
      最近更新 更多