gpt4 book ai didi

java - CompletableFuture 类中 join 方法的使用与 get 方法的使用

转载 作者:行者123 更新时间:2023-11-30 05:46:56 27 4
gpt4 key购买 nike

我想实现一个功能,将大文件分解为 block 并且可以并行处理。

我使用 CompletableFuture 并行运行任务。不幸的是,除非我使用 join,否​​则它不起作用。我对这种情况的发生感到惊讶,因为根据文档, get 也是返回结果的类中的阻塞方法。有人可以帮我弄清楚我做错了什么吗?

//cf.join(); if i uncommnet this everything works

如果我在 processChunk 方法中取消注释上述行,一切都会正常。我的值(value)观和一切都被打印出来。但是,如果我删除它,什么也不会发生。我收到的只是 future 已完成的通知,但内容未打印。

这是我的输出

i cmpleteddone
i cmpleteddone
i cmpleteddone
i cmpleteddone
i cmpleteddone

我的文本文件是一个非常小的文件(目前)

1212451,London,25000,Blocked 
1212452,London,215000,Open
1212453,London,125000,CreditBlocked
1212454,London,251000,DebitBlocked
1212455,London,2500,Open
1212456,London,4000,Closed
1212457,London,25100,Dormant
1212458,London,25010,Open
1212459,London,27000,Open
12124510,London,225000,Open
12124511,London,325000,Open
12124512,London,425000,Open
12124513,London,265000,Open
12124514,London,2577000,Open
12124515,London,2504400,Open


package com.org.java_trial.thread.executors;

import java.io.BufferedReader;
import java.io.FileReader;
import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;




public class ProcessReallyBigFile {

private static final ExecutorService ex = Executors.newFixedThreadPool(5);

private static CompletableFuture<String> processChunk(List<String> lines) {

CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> {

//just for purposes of testing, will be replaced with major function later
lines.stream().forEach(System.out::println);
return "done";
}, ex);

//cf.join(); if i uncommnet this everything works
return cf;
}

private static void readInChunks(String filepath, Integer chunksize) {

List<CompletableFuture<String>> completable = new ArrayList<>();
try (BufferedReader reader = Files.newBufferedReader(Paths.get(filepath))) {

String line = null;
List<String> collection = new ArrayList<String>();
int count = 0;

while ((line = reader.readLine()) != null) {

if (count % chunksize == chunksize - 1) {

collection.add(line);
completable.add(processChunk(collection));

collection.clear();

} else {

collection.add(line);
}
count++;
}

// any leftovers
if (collection.size() > 0)
completable.add(processChunk(collection));

} catch (IOException e) {
e.printStackTrace();
}


for (CompletableFuture c : completable) {
c.join();
if (c.isDone() || c.isCompletedExceptionally()) {
try {

System.out.println("i cmpleted" + c.get());
} catch (InterruptedException | ExecutionException e) {
// TODO Auto-generated catch block
e.printStackTrace();
}
}
}

ex.shutdown();

}

public static void main(String[] args) {

String filepath = "C:\\somak\\eclipse-workspace\\java_thingies\\java_trial\\account_1.csv";

readInChunks(filepath, 3);
}
}

最佳答案

原因是这样的:

collection.clear();

您的控件返回到不带 .join() 的调用方法,并且您的任务引用的集合已被清除。 幸运的是,您没有因并发访问而抛出异常。对共享资源的并发访问应始终同步。我宁愿这样做:

synchronized(collection) { 
collection.clear();
}

synchronized(collection) {
lines.stream().forEach(System.out::println);
}

这将确保访问 collection 对象时的线程安全,因为线程需要在实例 collection 上执行任何更新之前持有监视器。

此外,正如@Holger 所指出的,请执行以下操作:

synchronized(collection) {
collection.add(line);
}

关于java - CompletableFuture 类中 join 方法的使用与 get 方法的使用,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/54662251/

27 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com