Class WatermarkedSlidingWindowCollection<T>
java.lang.Object
com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection<T>
com.scaleoutsoftware.collections.timewindowing.WatermarkedSlidingWindowCollection<T>
- Type Parameters:
T- the object type of the source collection.
- All Implemented Interfaces:
Iterable<TimeWindow<T>>
The SlidingWindowCollection transforms a collection into an iterable collection of overlapping time windows. This
wrapper class can be used to manage the retention policy and add objects in chronological order to the underlying
source collection.
The difference between
SlidingWindowCollection and WatermarkedSlidingWindowCollection is that
windows in the WatermarkedSlidingWindowCollection can be closed if the watermark passes the inclusive end of
a window. The watermark also prevents items with timestamps that exceed the watermark from being added to the collection.-
Field Summary
Fields inherited from class com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection
sourceCollection, startTimeMs, timestampSelector, watermarkGenerator, watermarkMs -
Constructor Summary
ConstructorsConstructorDescriptionWatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long startTimeMs, long windowDurationMs, long everyMs, WatermarkGenerator watermarkGenerator) Instantiates a new WatermarkedSlidingWindowCollection.WatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, long everyMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) Instantiates a new WatermarkedSlidingWindowCollection. -
Method Summary
Modifier and TypeMethodDescriptionvoidforEach(Consumer<? super TimeWindow<T>> action) longRetrieves the window created "every" in milliseconds.longRetrieves the next window start time in milliseconds.longRetrieves the configured window duration in milliseconds.iterator()Methods inherited from class com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection
add, getSourceCollection, getStartTimeMs, getTimestampSelector, getWatermarkGenerator, getWatermarkMs
-
Constructor Details
-
WatermarkedSlidingWindowCollection
public WatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long startTimeMs, long windowDurationMs, long everyMs, WatermarkGenerator watermarkGenerator) Instantiates a new WatermarkedSlidingWindowCollection.- Parameters:
sourceCollection- the underlying source collectiontimestampSelector- theTimestampSelectoris used to pull a timestamp from an item in the source collection and subsequent insertions.startTimeMs- 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.windowDurationMs- the duration of a time windoweveryMs- the time between the starting point of each time windowwatermarkGenerator- theWatermarkGeneratoris used to generate a watermark. Entries that arrive before the watermark time are evicted. Windows whose inclusive end exceeds the watermark are closed.
-
WatermarkedSlidingWindowCollection
public WatermarkedSlidingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, long everyMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) Instantiates a new WatermarkedSlidingWindowCollection.- Parameters:
sourceCollection- the underlying source collectiontimestampSelector- theTimestampSelectoris used to pull a timestamp from an item in the source collection and subsequent insertions.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.windowDurationMs- the duration of a time windoweveryMs- the time between the starting point of each time windowwatermarkGenerator- theWatermarkGeneratoris used to generate a watermark. Entries that arrive before the watermark time are evicted. Windows whose inclusive end exceeds the watermark are closed.currentWatermarkMs- the current watermark in milliseconds.
-
-
Method Details
-
getWindowDurationMs
public long getWindowDurationMs()Retrieves the configured window duration in milliseconds.- Returns:
- the window duration in milliseconds.
-
getEveryMs
public long getEveryMs()Retrieves the window created "every" in milliseconds.- Returns:
- the window every in milliseconds.
-
getNextWindowStartTimeMs
public long getNextWindowStartTimeMs()Retrieves the next window start time in milliseconds.- Returns:
- the next window start time in milliseconds.
-
iterator
-
forEach
-
spliterator
-