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>>

public class WatermarkedTumblingWindowCollection<T> extends WatermarkedWindowingCollection<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.
  • 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 - the TimestampSelector is 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 - the 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.
    • 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 - the TimestampSelector is 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 - the 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.
      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

      public Iterator<TimeWindow<T>> iterator()
    • forEach

      public void forEach(Consumer<? super TimeWindow<T>> action)
    • spliterator

      public Spliterator<TimeWindow<T>> spliterator()