Class WatermarkedWindowingCollection<T>

java.lang.Object
com.scaleoutsoftware.collections.timewindowing.WatermarkedWindowingCollection<T>
Type Parameters:
T - the object type of the source collection.
All Implemented Interfaces:
Iterable<TimeWindow<T>>
Direct Known Subclasses:
WatermarkedSessionWindowCollection, WatermarkedSlidingWindowCollection, WatermarkedTumblingWindowCollection

public abstract class WatermarkedWindowingCollection<T> extends Object implements Iterable<TimeWindow<T>>
Used to transform a List into an iterable collection of TimeWindow. Windows in a WatermarkedWindowingCollection 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. Items that exclusively reside in closed windows are evicted from the collection.
  • Field Details

    • sourceCollection

      protected List<T> sourceCollection
      The source collection of items for the watermarked windowing collection.
    • timestampSelector

      protected TimestampSelector<T> timestampSelector
      The timestamp selector is used to select a timestamp from an element in the source collection.
    • watermarkGenerator

      protected WatermarkGenerator watermarkGenerator
      The watermark generator is used to generate a watermark for the windowing collection.
    • watermarkMs

      protected long watermarkMs
      The current watermark of the windowing collection.
    • startTimeMs

      protected long startTimeMs
      The inclusive start time of first window of a windowing collection.
  • Constructor Details

    • WatermarkedWindowingCollection

      public WatermarkedWindowingCollection(List<T> sourceCollection, TimestampSelector<T> timestampSelector, long startTimeMs, WatermarkGenerator watermarkGenerator)
      Instantiates a new SlidingWindowCollection
      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.
      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.
  • Method Details

    • add

      public List<TimeWindow<T>> add(T item)
      Uses the watermark generator to generate a watermark from the parameter item. If the items timestamp exceeds the watermark then the item is added to the source collection in time order. If the items timestamp equals or preceeds the watermark, the add is ignored. The newly created watermark may cause windows in the collection to close. Windows whose inclusive end exceeds the watermark are considered closed. Closed windows are returned. If no windows are closed, an empty list is returned.
      Parameters:
      item - the item to add.
      Returns:
      returns a list of a closed windows.
    • getSourceCollection

      public List<T> getSourceCollection()
      Retrieves the source collection.
      Returns:
      the source collection.
    • getTimestampSelector

      public TimestampSelector<T> getTimestampSelector()
      Retrieve the timestamp selector.
      Returns:
      the timestamp selector.
    • getStartTimeMs

      public long getStartTimeMs()
      Retrieve the start time in milliseconds.
      Returns:
      the start time in milliseconds.
    • getWatermarkGenerator

      public WatermarkGenerator getWatermarkGenerator()
      Retrieve the watermark generator.
      Returns:
      the watermark generator.
    • getWatermarkMs

      public long getWatermarkMs()
      Return the current watermark in milliseconds.
      Returns:
      the watermark in milliseconds.