Class WatermarkedTumblingWindowCollection<T>
java.lang.Object
com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection<T>
com.scaleoutsoftware.collections.timewindowing.WatermarkedTumblingWindowCollection<T>
- Type Parameters:
T- the object type of the source collection.
- All Implemented Interfaces:
Iterable<TimeWindow<T>>
The TumblingWindowCollection transforms a collection into an iterable collection of sequential 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
TumblingWindowCollection and WatermarkedTumblingWindowCollection is that
windows in the WatermarkedTumblingWindowCollection 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
ConstructorsConstructorDescriptionWatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long startTimeMs, long windowDurationMs, WatermarkGenerator watermarkGenerator) Instantiates a new WatermarkedTumblingWindowCollectionWatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) Instantiates a new WatermarkedTumblingWindowCollection -
Method Summary
Modifier and TypeMethodDescriptionvoidforEach(Consumer<? super TimeWindow<T>> action) longRetrieves the next window start time in milliseconds.longRetrieves the windows duration in milliseconds.iterator()Methods inherited from class com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection
add, getSourceCollection, getStartTimeMs, getTimestampSelector, getWatermarkGenerator, getWatermarkMs
-
Constructor Details
-
WatermarkedTumblingWindowCollection
public WatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long startTimeMs, long windowDurationMs, WatermarkGenerator watermarkGenerator) Instantiates a new WatermarkedTumblingWindowCollection- Parameters:
sourceCollection- the underlying source collection.timestampSelector- 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 window duration in milliseconds for each window.watermarkGenerator- theWatermarkGeneratoris used to generate a watermark. Entries that arrive before the watermark time are evicted. Windows whose inclusive end exceeds the watermark are closed.
-
WatermarkedTumblingWindowCollection
public WatermarkedTumblingWindowCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long nextWindowStartTimeMs, long windowDurationMs, WatermarkGenerator watermarkGenerator, long currentWatermarkMs) Instantiates a new WatermarkedTumblingWindowCollection- Parameters:
sourceCollection- the underlying source collection.timestampSelector- 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 window duration in milliseconds for each window.watermarkGenerator- 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 windows duration in milliseconds.- Returns:
- the windows duration 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
-