-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathBucketTest.java
More file actions
129 lines (102 loc) · 4.19 KB
/
Copy pathBucketTest.java
File metadata and controls
129 lines (102 loc) · 4.19 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
package bucket;
import org.junit.jupiter.api.Test;
import java.time.Duration;
import java.time.Instant;
import java.time.temporal.ChronoUnit;
import java.util.Arrays;
import java.util.Comparator;
import java.util.LinkedList;
import java.util.List;
import java.util.stream.Collectors;
class BucketTest {
@Test
void sample() {
Instant now = Instant.now();
Instant earlier = now.minus(20, ChronoUnit.MINUTES);
List<Event> events = Arrays.asList(new Event(now, "now"), new Event(earlier, "earlier"));
List<Event> sortedEvents = events.stream().sorted(byTimeStamp()).collect(Collectors.toList());
Duration samplingRate = Duration.ofMinutes(1);
Duration inspectionRange = Duration.ofMinutes(10);
Instant earliest = sortedEvents.get(0).timestamp;
Instant latest = sortedEvents.get(sortedEvents.size() - 1).timestamp;
BucketAssemblyLine bucketAssemblyLine = new BucketAssemblyLine(samplingRate, inspectionRange, earlier);
}
private Comparator<Event> byTimeStamp() {
return new Comparator<Event>() {
@Override
public int compare(Event o1, Event o2) {
return o1.timestamp.compareTo(o2.timestamp);
}
};
}
public static class BucketAssemblyLine {
private final LinkedList<Bucket> allBuckets = new LinkedList<>();
private final LinkedList<Bucket> assemblyLine = new LinkedList<>();
private final Duration samplingRate;
private final Duration inspectionRange;
public BucketAssemblyLine(Duration samplingRate, Duration inspectionRange, Instant earliest) {
this.samplingRate = samplingRate;
this.inspectionRange = inspectionRange;
addBucket(new Bucket(earliest, earliest.plus(inspectionRange)));
}
public void putIntoBucket(Event event) {
assemblyLine.removeIf(bucket -> !event.timestamp.isBefore(bucket.latest));
Instant nextNewBucketStart = allBuckets.peekLast().earliest().plus(samplingRate);
if (wouldAlreadyBeAddedToTheNextBucket(event, nextNewBucketStart)) {
for (; wouldAlreadyBeAddedToTheNextBucket(event, nextNewBucketStart); nextNewBucketStart = nextNewBucketStart.plus(samplingRate)) {
Bucket newBucket = new Bucket(nextNewBucketStart, nextNewBucketStart.plus(inspectionRange));
addBucket(newBucket);
}
}
Boolean addedToAtLeastOneBucket = assemblyLine.stream().reduce(Boolean.FALSE, (aBoolean, aBucket) -> aBucket.put(event), (a, b) -> a || b);
if (!addedToAtLeastOneBucket) {
throw new IllegalStateException("There was no bucket to take the event");
}
}
private boolean wouldAlreadyBeAddedToTheNextBucket(Event event, Instant nextNewBucketStart) {
return !nextNewBucketStart.isAfter(event.timestamp);
}
public List<Bucket> allBuckets() {
return allBuckets;
}
private void addBucket(Bucket bucket) {
allBuckets.add(bucket);
assemblyLine.add(bucket);
}
}
public static class Bucket {
private final List<Event> events = new LinkedList<>();
private final Instant earliest;
private final Instant latest;
public Bucket(Instant earliest, Instant latest) {
this.latest = latest;
this.earliest = earliest;
}
public Instant earliest() {
return earliest;
}
public Instant latest() {
return latest;
}
public boolean put(Event event) {
if (!event.timestamp.isBefore(latest)) {
return false;
}
if (earliest().isAfter(event.timestamp)) {
return false;
}
return this.events.add(event);
}
public List<Event> items() {
return events;
}
}
public static class Event {
public final Instant timestamp;
public final String data;
public Event(Instant timestamp, String data) {
this.timestamp = timestamp;
this.data = data;
}
}
}