Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,15 +19,26 @@
* The LatenessToleranceWatermarkGenerator generates watermarks with a constant lateness tolerance.
*/
public class LatenessToleranceWatermarkGenerator implements WatermarkGenerator {
private final long _latenessToleranceMs;
private long _currentWaterMarkMs = Long.MIN_VALUE;
private final long latenessToleranceMs;
private long currentWaterMarkMs = Long.MIN_VALUE;

/**
* Constructs a new LatenessToleranceWatermarkGenerator with a constant lateness tolerance.
* @param latenessToleranceMs the lateness tolerance in milliseconds.
*/
public LatenessToleranceWatermarkGenerator(long latenessToleranceMs) {
_latenessToleranceMs = latenessToleranceMs;
this.latenessToleranceMs = latenessToleranceMs;
}

/**
* Constructs a new LatenessToleranceWatermarkGenerator with a constant lateness tolerance and the current watermark
* in milliseconds.
* @param latenessToleranceMs the lateness tolerance in milliseconds.
* @param currentWaterMarkMs the current watermark in milliseconds.
*/
public LatenessToleranceWatermarkGenerator(long latenessToleranceMs, long currentWaterMarkMs) {
this.latenessToleranceMs = latenessToleranceMs;
this.currentWaterMarkMs = currentWaterMarkMs;
}

/**
Expand All @@ -37,20 +48,12 @@ public LatenessToleranceWatermarkGenerator(long latenessToleranceMs) {
*/
@Override
public long generateWatermark(long lastEventTimestampMs) {
long watermarkOption = lastEventTimestampMs - _latenessToleranceMs;
long watermarkOption = lastEventTimestampMs - latenessToleranceMs;

if (watermarkOption > _currentWaterMarkMs) {
_currentWaterMarkMs = watermarkOption;
if (watermarkOption > currentWaterMarkMs) {
currentWaterMarkMs = watermarkOption;
}

return _currentWaterMarkMs;
}

/**
* Returns the current watermark in milliseconds.
* @return the current watermark in milliseconds.
*/
public long currentWaterMarkMs() {
return _currentWaterMarkMs;
return currentWaterMarkMs;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,8 @@
* @param <T> the object type of the source collection.
*/
public class WatermarkedSessionWindowCollection<T> extends WatermarkedWindowingCollection<T> {
private long _timeoutMs;
private long _nextWindowStartTimeMs;
private long timeoutMs;
private long nextWindowStartTimeMs;

/**
* Instantiates a new SessionWindowCollection
Expand All @@ -50,41 +50,68 @@ public WatermarkedSessionWindowCollection(List<T> sourceCollection, TimestampSel
init(timeoutMs);
}

/**
* Instantiates a new SessionWindowCollection
* @param sourceCollection the underlying source collection.
* @param timestampSelector the {@link TimestampSelector} is used to pull a timestamp from an item in the source
* collection and subsequent insertions.
* @param nextWindowStartTimeMs the first time an object can be in a time window -- items before the start time will
* be evicted. The start time is also the start time of the first time window.
* @param timeoutMs the minimum amount of time between session window ranges.
* @param watermarkGenerator the {@link WatermarkGenerator} is used to generate a watermark. Entries that arrive
* before the watermark time are evicted. Windows whose inclusive end exceeds the watermark
* are closed.
* @param currentWatermarkMs the current watermark in milliseconds.
*/
public WatermarkedSessionWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long timeoutMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) {
super(sourceCollection, timestampSelector, nextWindowStartTimeMs, watermarkGenerator);
init(timeoutMs);
watermarkMs = currentWatermarkMs;
}

/**
* Retrieve the configured session window timeout in milliseconds.
* @return the timeout in milliseconds.
*/
public long getTimeoutMs() {
return timeoutMs;
}

private void init(long timeoutMs) {
_timeoutMs = timeoutMs;
this.timeoutMs = timeoutMs;
}

@Override
List<TimeWindow<T>> performEviction() {
EvictionMetadata<T> ret = Utils.performSessionWindowEviction(_sourceCollection, _timestampSelector, _watermarkMs, _timeoutMs);
EvictionMetadata<T> ret = Utils.performSessionWindowEviction(sourceCollection, timestampSelector, watermarkMs, timeoutMs);
return ret.getClosedWindows();
}

@Override
public Iterator<TimeWindow<T>> iterator() {
if(_sourceCollection == null || _sourceCollection.isEmpty()) {
if(sourceCollection == null || sourceCollection.isEmpty()) {
return Collections.emptyIterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toSessionWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _timeoutMs).iterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toSessionWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, timeoutMs).iterator();
}
}

@Override
public void forEach(Consumer<? super TimeWindow<T>> action) {
if(_sourceCollection != null && !_sourceCollection.isEmpty()) {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
Windowing.toSessionWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _timeoutMs).forEach(action);
if(sourceCollection != null && !sourceCollection.isEmpty()) {
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
Windowing.toSessionWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, timeoutMs).forEach(action);
}
}

@Override
public Spliterator<TimeWindow<T>> spliterator() {
if(_sourceCollection == null || _sourceCollection.isEmpty()) {
if(sourceCollection == null || sourceCollection.isEmpty()) {
return Spliterators.emptySpliterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toSessionWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _timeoutMs).spliterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toSessionWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, timeoutMs).spliterator();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -30,9 +30,9 @@
* @param <T> the object type of the source collection.
*/
public class WatermarkedSlidingWindowCollection<T> extends WatermarkedWindowingCollection<T> {
private long _windowDurationMs;
private long _everyMs;
private long _nextWindowStartTimeMs;
private long windowDurationMs;
private long everyMs;
private long nextWindowStartTimeMs;

/**
* Instantiates a new WatermarkedSlidingWindowCollection.
Expand All @@ -52,52 +52,96 @@ public WatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSel
init(windowDurationMs, everyMs);
}

/**
* Instantiates a new WatermarkedSlidingWindowCollection.
* @param sourceCollection the underlying source collection
* @param timestampSelector the {@link TimestampSelector} is used to pull a timestamp from an item in the source
* collection and subsequent insertions.
* @param nextWindowStartTimeMs the first time an object can be in a time window -- items before the start time will
* be evicted. The start time is also the start time of the first time window.
* @param windowDurationMs the duration of a time window
* @param everyMs the time between the starting point of each time window
* @param watermarkGenerator the {@link WatermarkGenerator} is used to generate a watermark. Entries that arrive
* before the watermark time are evicted. Windows whose inclusive end exceeds the watermark
* are closed.
* @param currentWatermarkMs the current watermark in milliseconds.
*/
public WatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, long everyMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) {
super(sourceCollection, timestampSelector, nextWindowStartTimeMs, watermarkGenerator);
init(windowDurationMs, everyMs);
this.watermarkMs = currentWatermarkMs;
}

/**
* Retrieves the configured window duration in milliseconds.
* @return the window duration in milliseconds.
*/
public long getWindowDurationMs() {
return windowDurationMs;
}

/**
* Retrieves the window created "every" in milliseconds.
* @return the window every in milliseconds.
*/
public long getEveryMs() {
return everyMs;
}

/**
* Retrieves the next window start time in milliseconds.
* @return the next window start time in milliseconds.
*/
public long getNextWindowStartTimeMs() {
return nextWindowStartTimeMs;
}

private void init(long windowDurationMs, long everyMs) {
if(windowDurationMs <=0) throw new IllegalArgumentException("window duration is <= 0 in param");
if(everyMs <= 0) throw new IllegalArgumentException("everyMs is <= 0 in param");
_windowDurationMs = windowDurationMs;
_everyMs = everyMs;
_nextWindowStartTimeMs = _startTimeMs;
this.windowDurationMs = windowDurationMs;
this.everyMs = everyMs;
nextWindowStartTimeMs = startTimeMs;
}

@Override
List<TimeWindow<T>> performEviction() {
EvictionMetadata<T> ret = Utils.performWatermarkedWindowedEviction(
_sourceCollection,
_timestampSelector,
_watermarkMs,
_windowDurationMs,
_everyMs,
_nextWindowStartTimeMs);
_nextWindowStartTimeMs = ret.getNextWindowStartTimeMs();
sourceCollection,
timestampSelector,
watermarkMs,
windowDurationMs,
everyMs,
nextWindowStartTimeMs);
nextWindowStartTimeMs = ret.getNextWindowStartTimeMs();
return ret.getClosedWindows();
}

@Override
public Iterator<TimeWindow<T>> iterator() {
if(_sourceCollection == null || _sourceCollection.isEmpty()) {
if(sourceCollection == null || sourceCollection.isEmpty()) {
return Collections.emptyIterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toSlidingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs, _everyMs).iterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toSlidingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs, everyMs).iterator();
}
}

@Override
public void forEach(Consumer<? super TimeWindow<T>> action) {
if(_sourceCollection != null && !_sourceCollection.isEmpty()) {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
Windowing.toSlidingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs, _everyMs).forEach(action);
if(sourceCollection != null && !sourceCollection.isEmpty()) {
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
Windowing.toSlidingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs, everyMs).forEach(action);
}
}

@Override
public Spliterator<TimeWindow<T>> spliterator() {
if (_sourceCollection == null || _sourceCollection.isEmpty()) {
if (sourceCollection == null || sourceCollection.isEmpty()) {
return Spliterators.emptySpliterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toSlidingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs, _everyMs).spliterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toSlidingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs, everyMs).spliterator();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,8 +29,8 @@
* @param <T> the object type of the source collection.
*/
public class WatermarkedTumblingWindowCollection<T> extends WatermarkedWindowingCollection<T> {
private long _windowDurationMs;
private long _nextWindowStartTimeMs;
private long windowDurationMs;
private long nextWindowStartTimeMs;

/**
* Instantiates a new WatermarkedTumblingWindowCollection
Expand All @@ -49,50 +49,85 @@ public WatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSe
init(windowDurationMs);
}

/**
* Instantiates a new WatermarkedTumblingWindowCollection
* @param sourceCollection the underlying source collection.
* @param timestampSelector the {@link TimestampSelector} is used to pull a timestamp from an item in the source
* collection and subsequent insertions.
* @param nextWindowStartTimeMs the first time an object can be in a time window -- items before the start time will
* be evicted. The start time is also the start time of the first time window.
* @param windowDurationMs the window duration in milliseconds for each window.
* @param watermarkGenerator the {@link WatermarkGenerator} is used to generate a watermark. Entries that arrive
* before the watermark time are evicted. Windows whose inclusive end exceeds the watermark
* are closed.
* @param currentWatermarkMs the current watermark in milliseconds.
*/
public WatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) {
super(sourceCollection, timestampSelector, nextWindowStartTimeMs, watermarkGenerator);
init(windowDurationMs);
watermarkMs = currentWatermarkMs;
}

/**
* Retrieves the windows duration in milliseconds.
* @return the windows duration in milliseconds.
*/
public long getWindowDurationMs() {
return windowDurationMs;
}

/**
* Retrieves the next window start time in milliseconds.
* @return the next window start time in milliseconds.
*/
public long getNextWindowStartTimeMs() {
return nextWindowStartTimeMs;
}

private void init(long windowDurationMs) {
if(windowDurationMs <= 0) throw new IllegalArgumentException("window duration is <= 0");
_windowDurationMs = windowDurationMs;
_nextWindowStartTimeMs = _startTimeMs;
this.windowDurationMs = windowDurationMs;
nextWindowStartTimeMs = startTimeMs;
}

@Override
List<TimeWindow<T>> performEviction() {
EvictionMetadata<T> ret = Utils.performWatermarkedWindowedEviction(
_sourceCollection,
_timestampSelector,
_watermarkMs,
_windowDurationMs,
_windowDurationMs,
_nextWindowStartTimeMs);
_nextWindowStartTimeMs = ret.getNextWindowStartTimeMs();
sourceCollection,
timestampSelector,
watermarkMs,
windowDurationMs,
windowDurationMs,
nextWindowStartTimeMs);
nextWindowStartTimeMs = ret.getNextWindowStartTimeMs();
return ret.getClosedWindows();
}

@Override
public Iterator<TimeWindow<T>> iterator() {
if(_sourceCollection == null || _sourceCollection.isEmpty()) {
if(sourceCollection == null || sourceCollection.isEmpty()) {
return Collections.emptyIterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toTumblingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs).iterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toTumblingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs).iterator();
}
}

@Override
public void forEach(Consumer<? super TimeWindow<T>> action) {
if(_sourceCollection != null && !_sourceCollection.isEmpty()) {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
Windowing.toTumblingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs).forEach(action);
if(sourceCollection != null && !sourceCollection.isEmpty()) {
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
Windowing.toTumblingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs).forEach(action);
}
}

@Override
public Spliterator<TimeWindow<T>> spliterator() {
if(_sourceCollection == null || _sourceCollection.isEmpty()) {
if(sourceCollection == null || sourceCollection.isEmpty()) {
return Spliterators.emptySpliterator();
} else {
long end = _timestampSelector.select(_sourceCollection.get(_sourceCollection.size()-1)) + 1;
return Windowing.toTumblingWindows(_sourceCollection, _timestampSelector, _nextWindowStartTimeMs, end, _windowDurationMs).spliterator();
long end = timestampSelector.select(sourceCollection.get(sourceCollection.size()-1)) + 1;
return Windowing.toTumblingWindows(sourceCollection, timestampSelector, nextWindowStartTimeMs, end, windowDurationMs).spliterator();
}
}
}
Expand Down
Loading