-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathMain.java
More file actions
131 lines (118 loc) · 4.66 KB
/
Copy pathMain.java
File metadata and controls
131 lines (118 loc) · 4.66 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
import java.io.File;
import java.io.FileOutputStream;
import java.io.OutputStream;
import java.util.*;
import java.util.concurrent.*;
public class Main {
public static void main(String[] args) {
int[] lineCounts = {1, 10, 100, 1000};
int[] groupCounts = {1000, 100, 10, 1};
int[] bufferCapacitys = {10, 100, 1000};
int[] arrayThreshholds = {5, 10, 20};
boolean[] randomInt = {false, true};
System.out.println("Test File\t\t\tExecution Time(ms)");
for (int i = 0; i < lineCounts.length; i++) {
for (int j = 0; j < groupCounts.length; j++) {
for (boolean random : randomInt) {
String path = "l" + lineCounts[i] + "_g" + groupCounts[j] + "_r" + random + ".csv";
makeFile(path.toString(), lineCounts[i] * 10000, groupCounts[j], random);
long start = System.currentTimeMillis();
hashAggregate(path.toString(), bufferCapacitys[2], arrayThreshholds[1]);
long end = System.currentTimeMillis();
System.out.println(path.toString() + "\t\t\t" + (end - start));
}
}
}
return;
}
public static void makeFile(String path, int rowCount, int groupCount, boolean random) {
OutputStream out = null;
try {
out = new FileOutputStream(new File(path));
Random r = new Random();
int a,b;
for (int i = 0; i < rowCount; i++) {
if(!random) {
a = i % groupCount;
b = i;
}
else {
a = r.nextInt() / groupCount;
b = r.nextInt();
}
out.write((Integer.valueOf(a).toString() + "," + Integer.valueOf(b).toString() + "\r\n").getBytes());
}
out.close();
} catch (Exception e) {
e.printStackTrace();
} finally {
try {
out.close();
} catch (Exception e) {
e.printStackTrace();
}
}
}
public static void hashAggregate(String path, int bufferCapacity, int arrayThreshhold) {
int consumerCount = Runtime.getRuntime().availableProcessors();
BlockingQueue<DataRow> queue = new LinkedBlockingDeque<>(bufferCapacity);
Map result = Collections.synchronizedMap(new HashMap<Integer, Object>());
ExecutorService service = Executors.newCachedThreadPool();
DataReader reader = new DataReader(queue, path);
Future readerFuture = service.submit(reader);
HashSet<DataConsumer> builders = new HashSet<>();
HashSet<Future> builderFutures = new HashSet<>();
for (int count = 0; count < consumerCount; count++) {
DataConsumer builder = new DataConsumer(queue, result, arrayThreshhold);
builders.add(builder);
builderFutures.add(service.submit(builder));
}
try {
Object f = readerFuture.get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
}
if (readerFuture.isDone()) {
for (DataConsumer builder : builders) {
builder.stop();
}
}
for (Future builderFuture : builderFutures) {
try {
Object f = builderFuture.get();
} catch (InterruptedException e) {
e.printStackTrace();
} catch (ExecutionException e) {
e.printStackTrace();
}
}
service.shutdown();
processResult(result);
}
public static void processResult(Map result) {
/* sum(b) and avg(b) can be optimized by SIMD */
//System.out.println("a\tcdb\tadb");
for (Object aKey : result.keySet()) {
int bSum = 0;
Object bKeys = result.get(aKey);
if(bKeys instanceof HashSet) {
for (Integer bKey : (HashSet<Integer>)bKeys) {
bSum += bKey;
}
int bCount = ((HashSet<Integer>)bKeys).size();
double bAvg = (double) bSum / bCount;
//System.out.println(aKey + "\t" + bCount + "\t" + bAvg);
}
else {
for (Integer bKey : (ArrayList<Integer>)bKeys) {
bSum += bKey;
}
int bCount = ((ArrayList<Integer>)bKeys).size();
double bAvg = (double) bSum / bCount;
//System.out.println(aKey + "\t" + bCount + "\t" + bAvg);
}
}
}
}