Skip to content

Commit 9a5bc24

Browse files
committed
2 03 HW1 concurrentMultiply3
1 parent 5d6b74c commit 9a5bc24

2 files changed

Lines changed: 33 additions & 89 deletions

File tree

src/main/java/ru/javaops/masterjava/matrix/MainMatrix.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ public static void main(String[] args) throws ExecutionException, InterruptedExc
3030
singleThreadSum += duration;
3131

3232
start = System.currentTimeMillis();
33-
final int[][] concurrentMatrixC = MatrixUtil.concurrentMultiply(matrixA, matrixB, executor);
33+
final int[][] concurrentMatrixC = MatrixUtil.concurrentMultiply2(matrixA, matrixB, executor);
3434
duration = (System.currentTimeMillis() - start) / 1000.;
3535
out("Concurrent thread time, sec: %.3f", duration);
3636
concurrentThreadSum += duration;

src/main/java/ru/javaops/masterjava/matrix/MatrixUtil.java

Lines changed: 32 additions & 88 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,12 @@
11
package ru.javaops.masterjava.matrix;
22

3-
import java.util.*;
4-
import java.util.concurrent.*;
3+
import java.util.ArrayList;
4+
import java.util.List;
5+
import java.util.Random;
6+
import java.util.concurrent.Callable;
7+
import java.util.concurrent.CountDownLatch;
8+
import java.util.concurrent.ExecutionException;
9+
import java.util.concurrent.ExecutorService;
510
import java.util.stream.Collectors;
611
import java.util.stream.IntStream;
712

@@ -11,91 +16,7 @@
1116
*/
1217
public class MatrixUtil {
1318

14-
public static int[][] concurrentMultiply(int[][] matrixA, int[][] matrixB, ExecutorService executor) throws InterruptedException, ExecutionException {
15-
final int matrixSize = matrixA.length;
16-
final int[][] matrixC = new int[matrixSize][matrixSize];
17-
18-
class ColumnMultipleResult {
19-
private final int col;
20-
private final int[] columnC;
21-
22-
private ColumnMultipleResult(int col, int[] columnC) {
23-
this.col = col;
24-
this.columnC = columnC;
25-
}
26-
}
27-
28-
final CompletionService<ColumnMultipleResult> completionService = new ExecutorCompletionService<>(executor);
29-
30-
for (int j = 0; j < matrixSize; j++) {
31-
final int col = j;
32-
final int[] columnB = new int[matrixSize];
33-
for (int k = 0; k < matrixSize; k++) {
34-
columnB[k] = matrixB[k][col];
35-
}
36-
completionService.submit(() -> {
37-
final int[] columnC = new int[matrixSize];
38-
39-
for (int row = 0; row < matrixSize; row++) {
40-
final int[] rowA = matrixA[row];
41-
int sum = 0;
42-
for (int k = 0; k < matrixSize; k++) {
43-
sum += rowA[k] * columnB[k];
44-
}
45-
columnC[row] = sum;
46-
}
47-
return new ColumnMultipleResult(col, columnC);
48-
});
49-
}
50-
51-
for (int i = 0; i < matrixSize; i++) {
52-
ColumnMultipleResult res = completionService.take().get();
53-
for (int k = 0; k < matrixSize; k++) {
54-
matrixC[k][res.col] = res.columnC[k];
55-
}
56-
}
57-
return matrixC;
58-
}
59-
60-
public static int[][] concurrentMultiplyCayman(int[][] matrixA, int[][] matrixB, ExecutorService executor) throws InterruptedException, ExecutionException {
61-
final int matrixSize = matrixA.length;
62-
final int[][] matrixResult = new int[matrixSize][matrixSize];
63-
final int threadCount = Runtime.getRuntime().availableProcessors();
64-
final int maxIndex = matrixSize * matrixSize;
65-
final int cellsInThread = maxIndex / threadCount;
66-
final int[][] matrixBFinal = new int[matrixSize][matrixSize];
67-
68-
for (int i = 0; i < matrixSize; i++) {
69-
for (int j = 0; j < matrixSize; j++) {
70-
matrixBFinal[i][j] = matrixB[j][i];
71-
}
72-
}
73-
74-
Set<Callable<Boolean>> threads = new HashSet<>();
75-
int fromIndex = 0;
76-
for (int i = 1; i <= threadCount; i++) {
77-
final int toIndex = i == threadCount ? maxIndex : fromIndex + cellsInThread;
78-
final int firstIndexFinal = fromIndex;
79-
threads.add(() -> {
80-
for (int j = firstIndexFinal; j < toIndex; j++) {
81-
final int row = j / matrixSize;
82-
final int col = j % matrixSize;
83-
84-
int sum = 0;
85-
for (int k = 0; k < matrixSize; k++) {
86-
sum += matrixA[row][k] * matrixBFinal[col][k];
87-
}
88-
matrixResult[row][col] = sum;
89-
}
90-
return true;
91-
});
92-
fromIndex = toIndex;
93-
}
94-
executor.invokeAll(threads);
95-
return matrixResult;
96-
}
97-
98-
public static int[][] concurrentMultiplyDarthVader(int[][] matrixA, int[][] matrixB, ExecutorService executor)
19+
public static int[][] concurrentMultiplyStreams(int[][] matrixA, int[][] matrixB, ExecutorService executor)
9920
throws InterruptedException, ExecutionException {
10021

10122
final int matrixSize = matrixA.length;
@@ -161,7 +82,30 @@ public static int[][] concurrentMultiply2(int[][] matrixA, int[][] matrixB, Exec
16182
return matrixC;
16283
}
16384

164-
// Optimized by https://habrahabr.ru/post/114797/
85+
public static int[][] concurrentMultiply3(int[][] matrixA, int[][] matrixB, ExecutorService executor) throws InterruptedException {
86+
final int matrixSize = matrixA.length;
87+
final int[][] matrixC = new int[matrixSize][matrixSize];
88+
final CountDownLatch latch = new CountDownLatch(matrixSize);
89+
90+
for (int row = 0; row < matrixSize; row++) {
91+
final int[] rowA = matrixA[row];
92+
final int[] rowC = matrixC[row];
93+
94+
executor.submit(() -> {
95+
for (int idx = 0; idx < matrixSize; idx++) {
96+
final int elA = rowA[idx];
97+
final int[] rowB = matrixB[idx];
98+
for (int col = 0; col < matrixSize; col++) {
99+
rowC[col] += elA * rowB[col];
100+
}
101+
}
102+
latch.countDown();
103+
});
104+
}
105+
latch.await();
106+
return matrixC;
107+
}
108+
165109
public static int[][] singleThreadMultiplyOpt(int[][] matrixA, int[][] matrixB) {
166110
final int matrixSize = matrixA.length;
167111
final int[][] matrixC = new int[matrixSize][matrixSize];

0 commit comments

Comments
 (0)