diff --git a/libraries-stream/pom.xml b/libraries-stream/pom.xml index 6cd931c299..fea437db6b 100644 --- a/libraries-stream/pom.xml +++ b/libraries-stream/pom.xml @@ -49,11 +49,23 @@ 0.9.12 - 2.6.0 + 3.0.0 0.9.0 8.2.0 0.8.1 1.15 + + + + org.apache.maven.plugins + maven-compiler-plugin + + 21 + 21 + + + + diff --git a/libraries-stream/src/test/java/com/baeldung/parallel_collectors/ParallelCollectorsVirtualThreadsManualTest.java b/libraries-stream/src/test/java/com/baeldung/parallel_collectors/ParallelCollectorsVirtualThreadsManualTest.java new file mode 100644 index 0000000000..3038fb74f9 --- /dev/null +++ b/libraries-stream/src/test/java/com/baeldung/parallel_collectors/ParallelCollectorsVirtualThreadsManualTest.java @@ -0,0 +1,73 @@ +package com.baeldung.parallel_collectors; + +import com.pivovarit.collectors.ParallelCollectors; +import org.junit.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.time.Duration; +import java.time.Instant; +import java.util.concurrent.Executors; +import java.util.function.Supplier; +import java.util.stream.Stream; + +import static java.util.stream.Collectors.toList; + +public class ParallelCollectorsVirtualThreadsManualTest { + + private static final Logger log = LoggerFactory.getLogger(ParallelCollectorsVirtualThreadsManualTest.class); + + // increase the number of parallel processes to find the max number of threads on your machine + @Test + public void givenParallelism_whenUsingOSThreads_thenShouldRunOutOfThreads() { + int parallelProcesses = 50_000; + + var e = Executors.newFixedThreadPool(parallelProcesses); + + var result = timed(() -> Stream.iterate(0, i -> i + 1).limit(parallelProcesses) + .collect(ParallelCollectors.parallel(i -> fetchById(i), toList(), e, parallelProcesses)) + .join()); + + log.info("{}", result); + } + + @Test + public void givenParallelism_whenUsingVThreads_thenShouldProcessInParallel() { + int parallelProcesses = 1000_000; + + var result = timed(() -> Stream.iterate(0, i -> i + 1).limit(parallelProcesses) + .collect(ParallelCollectors.parallel(i -> fetchById(i), toList())) + .join()); + + log.info("{}", result); + } + + @Test + public void givenParallelismAndPCollectors2_whenUsingVThreads_thenShouldProcessInParallel() { + int parallelProcesses = 1000_000; + + var result = timed(() -> Stream.iterate(0, i -> i + 1).limit(parallelProcesses) + .collect(ParallelCollectors.parallel(i -> fetchById(i), toList(), Executors.newVirtualThreadPerTaskExecutor(), Integer.MAX_VALUE)) + .join()); + + log.info("{}", result); + } + + private static String fetchById(int id) { + try { + Thread.sleep(1000); + } catch (InterruptedException e) { + // ignore shamelessly + } + + return "user-" + id; + } + + private static T timed(Supplier supplier) { + var before = Instant.now(); + T result = supplier.get(); + var after = Instant.now(); + log.info("Execution time: {} ms", Duration.between(before, after).toMillis()); + return result; + } +} diff --git a/pom.xml b/pom.xml index 60d42c5517..54c54bf7b8 100644 --- a/pom.xml +++ b/pom.xml @@ -737,7 +737,7 @@ libraries-security libraries-server-2 libraries-server - libraries-stream + libraries-testing-2 libraries-transform libraries