Обертка Stream.parallel() службой Executor в Java ⇐ JAVA

Программисты JAVA общаются здесь
Anonymous
Обертка Stream.parallel() службой Executor в Java

Сообщение Anonymous »

Я изучаю Java Core на курсе DMDEV в Udemy и YouTube. На уровне 2 в уроках Concurrent я, вероятно, увидел ошибку в коде, допущенном репетитором. Но мой вопрос к нему остался без ответа. Так что, возможно, эксперты здесь смогут помочь.
Задача: имеется массив из 1_000_000 элементов, заполненный случайными целыми числами в диапазоне от 1 до 300. Напишите код на Java, чтобы найти максимальный элемент по используя 10 нитей.
Учитель предложил такое решение:

Код: Выделить всё

public class TaskFromDmdevForSOF {
public static void main(String[] args) throws ExecutionException, InterruptedException {
int[] values = new int[1_000_000];
Random random = new Random();
for (int i = 0; i < values.length; i++) {
values[i] = random.nextInt(300) + 1;
}

ExecutorService threadPool = Executors.newFixedThreadPool(10);
int max = findMaxParallel(values, threadPool);
System.out.println(max);
threadPool.shutdown();
threadPool.awaitTermination(1, TimeUnit.MINUTES);
}
private static int findMaxParallel(int[] values, ExecutorService executorService) throws ExecutionException, InterruptedException {
return executorService.submit(() -> IntStream.of(values)
.parallel()
.max()
.orElse(Integer.MIN_VALUE)).get();
}
}
Но мне стало любопытно, реально ли под этой задачей работают 10 потоков?
Я добавил .peek() после .parallel() и увидел, что есть общий ForkJoinPool с Под капотом работают 4 потока.
На мой взгляд, нам нужно использовать новый ForkJoinPool(10) вместо Executors.newFixedThreadPool(10), чтобы гарантировать, что работают настоящие 10 потоков.
Или лучше нам нужно реализовать наш собственный класс от расширения RecursiveTask для правильного решения этой задачи.
Я прав? Большое спасибо за ваши ответы.
UPD
Это мое решение этой задачи с помощью RecursiveTask. Верно ли это?

Код: Выделить всё

public class TaskFromDmdevForSOF2 {
private static int[] values;

public static void main(String[] args) throws InterruptedException {
values = new int[1_000_000];
Random random = new Random();
for (int i = 0; i < values.length; i++) {
values[i] = random.nextInt(300) + 1;
}

ForkJoinPool forkJoinPool = new ForkJoinPool(10);
MyRecursiveTaskFJP myRecursiveTask = new MyRecursiveTaskFJP(0, values.length);
Integer max = forkJoinPool.invoke(myRecursiveTask);
System.out.println(max);
forkJoinPool.shutdown();
forkJoinPool.awaitTermination(1, TimeUnit.MINUTES);

}

public static class MyRecursiveTaskFJP extends RecursiveTask {
int from;
int to;

public MyRecursiveTaskFJP(int from, int to) {
this.from = from;
this.to = to;
}
@Override
protected Integer compute() {
if ((to - from)  max) {
max = values[i];
}
}
return max;
}
int middle = ((to - from) / 2) + from;
MyRecursiveTaskFJP task1 = new MyRecursiveTaskFJP(from, middle);
MyRecursiveTaskFJP task2 = new MyRecursiveTaskFJP(middle, to);
task2.fork();
task1.fork();
return Integer.max(task1.join(), task2.join());
}
}
}

Подробнее здесь: https://stackoverflow.com/questions/791 ... ce-in-java

Вернуться в «JAVA»