Skip to content

Commit 3739d67

Browse files
committed
2 parents eb50aeb + e0adc83 commit 3739d67

213 files changed

Lines changed: 1865 additions & 821 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

core-java-modules/core-java-jpms/decoupling-pattern2/consumermodule2/pom.xml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,4 +45,4 @@
4545
<providermodule.version>1.0</providermodule.version>
4646
</properties>
4747

48-
</project>
48+
</project>

core-java-modules/core-java-streams-3/pom.xml

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,17 @@
2727
<version>${lombok.version}</version>
2828
<scope>provided</scope>
2929
</dependency>
30+
<dependency>
31+
<groupId>org.openjdk.jmh</groupId>
32+
<artifactId>jmh-core</artifactId>
33+
<version>${jmh.version}</version>
34+
</dependency>
35+
<dependency>
36+
<groupId>org.openjdk.jmh</groupId>
37+
<artifactId>jmh-generator-annprocess</artifactId>
38+
<version>${jmh.version}</version>
39+
<scope>test</scope>
40+
</dependency>
3041
<!-- test scoped -->
3142
<dependency>
3243
<groupId>org.assertj</groupId>
@@ -44,11 +55,30 @@
4455
<filtering>true</filtering>
4556
</resource>
4657
</resources>
58+
<plugins>
59+
<plugin>
60+
<groupId>org.apache.maven.plugins</groupId>
61+
<artifactId>maven-compiler-plugin</artifactId>
62+
<configuration>
63+
<source>1.8</source>
64+
<target>1.8</target>
65+
<annotationProcessorPaths>
66+
<path>
67+
<groupId>org.openjdk.jmh</groupId>
68+
<artifactId>jmh-generator-annprocess</artifactId>
69+
<version>${jmh.version}</version>
70+
</path>
71+
</annotationProcessorPaths>
72+
</configuration>
73+
</plugin>
74+
</plugins>
4775
</build>
4876

4977
<properties>
78+
<lombok.version>1.18.20</lombok.version>
5079
<!-- testing -->
5180
<assertj.version>3.6.1</assertj.version>
81+
<jmh.version>1.29</jmh.version>
5282
</properties>
5383

5484
</project>
Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
package com.baeldung.streams.parallel;
2+
3+
public class BenchmarkRunner {
4+
5+
public static void main(String[] args) throws Exception {
6+
org.openjdk.jmh.Main.main(args);
7+
}
8+
9+
}
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import org.openjdk.jmh.annotations.Benchmark;
4+
import org.openjdk.jmh.annotations.BenchmarkMode;
5+
import org.openjdk.jmh.annotations.Mode;
6+
import org.openjdk.jmh.annotations.OutputTimeUnit;
7+
8+
import java.util.ArrayList;
9+
import java.util.LinkedList;
10+
import java.util.List;
11+
import java.util.concurrent.TimeUnit;
12+
import java.util.stream.IntStream;
13+
14+
public class DifferentSourceSplitting {
15+
16+
private static final List<Integer> arrayListOfNumbers = new ArrayList<>();
17+
private static final List<Integer> linkedListOfNumbers = new LinkedList<>();
18+
19+
static {
20+
IntStream.rangeClosed(1, 1_000_000).forEach(i -> {
21+
arrayListOfNumbers.add(i);
22+
linkedListOfNumbers.add(i);
23+
});
24+
}
25+
26+
@Benchmark
27+
@BenchmarkMode(Mode.AverageTime)
28+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
29+
public static void differentSourceArrayListSequential() {
30+
arrayListOfNumbers.stream().reduce(0, Integer::sum);
31+
}
32+
33+
@Benchmark
34+
@BenchmarkMode(Mode.AverageTime)
35+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
36+
public static void differentSourceArrayListParallel() {
37+
arrayListOfNumbers.parallelStream().reduce(0, Integer::sum);
38+
}
39+
40+
@Benchmark
41+
@BenchmarkMode(Mode.AverageTime)
42+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
43+
public static void differentSourceLinkedListSequential() {
44+
linkedListOfNumbers.stream().reduce(0, Integer::sum);
45+
}
46+
47+
@Benchmark
48+
@BenchmarkMode(Mode.AverageTime)
49+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
50+
public static void differentSourceLinkedListParallel() {
51+
linkedListOfNumbers.parallelStream().reduce(0, Integer::sum);
52+
}
53+
54+
}
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import org.openjdk.jmh.annotations.Benchmark;
4+
import org.openjdk.jmh.annotations.BenchmarkMode;
5+
import org.openjdk.jmh.annotations.Mode;
6+
import org.openjdk.jmh.annotations.OutputTimeUnit;
7+
8+
import java.util.Arrays;
9+
import java.util.concurrent.TimeUnit;
10+
import java.util.stream.IntStream;
11+
12+
public class MemoryLocalityCosts {
13+
14+
private static final int[] intArray = new int[1_000_000];
15+
private static final Integer[] integerArray = new Integer[1_000_000];
16+
17+
static {
18+
IntStream.rangeClosed(1, 1_000_000).forEach(i -> {
19+
intArray[i-1] = i;
20+
integerArray[i-1] = i;
21+
});
22+
}
23+
24+
@Benchmark
25+
@BenchmarkMode(Mode.AverageTime)
26+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
27+
public static void localityIntArraySequential() {
28+
Arrays.stream(intArray).reduce(0, Integer::sum);
29+
}
30+
31+
@Benchmark
32+
@BenchmarkMode(Mode.AverageTime)
33+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
34+
public static void localityIntArrayParallel() {
35+
Arrays.stream(intArray).parallel().reduce(0, Integer::sum);
36+
}
37+
38+
@Benchmark
39+
@BenchmarkMode(Mode.AverageTime)
40+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
41+
public static void localityIntegerArraySequential() {
42+
Arrays.stream(integerArray).reduce(0, Integer::sum);
43+
}
44+
45+
@Benchmark
46+
@BenchmarkMode(Mode.AverageTime)
47+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
48+
public static void localityIntegerArrayParallel() {
49+
Arrays.stream(integerArray).parallel().reduce(0, Integer::sum);
50+
}
51+
52+
}
Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import org.openjdk.jmh.annotations.Benchmark;
4+
import org.openjdk.jmh.annotations.BenchmarkMode;
5+
import org.openjdk.jmh.annotations.Mode;
6+
import org.openjdk.jmh.annotations.OutputTimeUnit;
7+
8+
import java.util.ArrayList;
9+
import java.util.List;
10+
import java.util.concurrent.TimeUnit;
11+
import java.util.stream.Collectors;
12+
import java.util.stream.IntStream;
13+
14+
public class MergingCosts {
15+
16+
private static final List<Integer> arrayListOfNumbers = new ArrayList<>();
17+
18+
static {
19+
IntStream.rangeClosed(1, 1_000_000).forEach(i -> {
20+
arrayListOfNumbers.add(i);
21+
});
22+
}
23+
24+
@Benchmark
25+
@BenchmarkMode(Mode.AverageTime)
26+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
27+
public static void mergingCostsSumSequential() {
28+
arrayListOfNumbers.stream().reduce(0, Integer::sum);
29+
}
30+
31+
@Benchmark
32+
@BenchmarkMode(Mode.AverageTime)
33+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
34+
public static void mergingCostsSumParallel() {
35+
arrayListOfNumbers.stream().parallel().reduce(0, Integer::sum);
36+
}
37+
38+
@Benchmark
39+
@BenchmarkMode(Mode.AverageTime)
40+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
41+
public static void mergingCostsGroupingSequential() {
42+
arrayListOfNumbers.stream().collect(Collectors.toSet());
43+
}
44+
45+
@Benchmark
46+
@BenchmarkMode(Mode.AverageTime)
47+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
48+
public static void mergingCostsGroupingParallel() {
49+
arrayListOfNumbers.stream().parallel().collect(Collectors.toSet());
50+
}
51+
52+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import java.util.Arrays;
4+
import java.util.List;
5+
6+
public class ParallelStream {
7+
8+
public static void main(String[] args) {
9+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
10+
listOfNumbers.parallelStream().forEach(number ->
11+
System.out.println(number + " " + Thread.currentThread().getName())
12+
);
13+
}
14+
15+
}
Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,15 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import java.util.Arrays;
4+
import java.util.List;
5+
6+
public class SequentialStream {
7+
8+
public static void main(String[] args) {
9+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
10+
listOfNumbers.stream().forEach(number ->
11+
System.out.println(number + " " + Thread.currentThread().getName())
12+
);
13+
}
14+
15+
}
Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import org.openjdk.jmh.annotations.Benchmark;
4+
import org.openjdk.jmh.annotations.BenchmarkMode;
5+
import org.openjdk.jmh.annotations.Mode;
6+
import org.openjdk.jmh.annotations.OutputTimeUnit;
7+
8+
import java.util.concurrent.TimeUnit;
9+
import java.util.stream.IntStream;
10+
11+
public class SplittingCosts {
12+
13+
@Benchmark
14+
@BenchmarkMode(Mode.AverageTime)
15+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
16+
public static void sourceSplittingIntStreamSequential() {
17+
IntStream.rangeClosed(1, 100).reduce(0, Integer::sum);
18+
}
19+
20+
@Benchmark
21+
@BenchmarkMode(Mode.AverageTime)
22+
@OutputTimeUnit(TimeUnit.NANOSECONDS)
23+
public static void sourceSplittingIntStreamParallel() {
24+
IntStream.rangeClosed(1, 100).parallel().reduce(0, Integer::sum);
25+
}
26+
27+
}
Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,46 @@
1+
package com.baeldung.streams.parallel;
2+
3+
import org.junit.jupiter.api.Test;
4+
5+
import java.util.Arrays;
6+
import java.util.List;
7+
import java.util.concurrent.ExecutionException;
8+
import java.util.concurrent.ForkJoinPool;
9+
10+
import static org.assertj.core.api.Assertions.assertThat;
11+
12+
class ForkJoinUnitTest {
13+
14+
@Test
15+
void givenSequentialStreamOfNumbers_whenReducingSumWithIdentityFive_thenResultIsCorrect() {
16+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
17+
int sum = listOfNumbers.stream().reduce(5, Integer::sum);
18+
assertThat(sum).isEqualTo(15);
19+
}
20+
21+
@Test
22+
void givenParallelStreamOfNumbers_whenReducingSumWithIdentityFive_thenResultIsNotCorrect() {
23+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
24+
int sum = listOfNumbers.parallelStream().reduce(5, Integer::sum);
25+
assertThat(sum).isNotEqualTo(15);
26+
}
27+
28+
@Test
29+
void givenParallelStreamOfNumbers_whenReducingSumWithIdentityZero_thenResultIsCorrect() {
30+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
31+
int sum = listOfNumbers.parallelStream().reduce(0, Integer::sum) + 5;
32+
assertThat(sum).isEqualTo(15);
33+
}
34+
35+
@Test
36+
public void givenParallelStreamOfNumbers_whenUsingCustomThreadPool_thenResultIsCorrect()
37+
throws InterruptedException, ExecutionException {
38+
List<Integer> listOfNumbers = Arrays.asList(1, 2, 3, 4);
39+
ForkJoinPool customThreadPool = new ForkJoinPool(4);
40+
int sum = customThreadPool.submit(
41+
() -> listOfNumbers.parallelStream().reduce(0, Integer::sum)).get();
42+
customThreadPool.shutdown();
43+
assertThat(sum).isEqualTo(10);
44+
}
45+
46+
}

0 commit comments

Comments
 (0)