Class GroupByReader


public class GroupByReader extends ProxyReader
A proxy that divides records into groups and applies summary operations to each group; similar to "group by" in SQL, but applied to streaming data.
  • Constructor Details

    • GroupByReader

      public GroupByReader(DataReader reader, String... groupByFields)
      Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up to 100 records.
    • GroupByReader

      public GroupByReader(DataReader reader, int queueSize, String... groupByFields)
      Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up to queueSize records.
    • GroupByReader

      public GroupByReader(DataReader reader, FieldList groupByFields)
      Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up to 100 records.
    • GroupByReader

      public GroupByReader(DataReader reader, int queueSize, FieldList groupByFields)
      Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up to queueSize records.
  • Method Details

    • getGroupByFields

      public FieldList getGroupByFields()
    • isExcludeNulls

      public boolean isExcludeNulls()
      Indicates if null groups should be excluded from results (defaults to false).
    • setExcludeNulls

      public GroupByReader setExcludeNulls(boolean excludeNulls)
      Indicates if null groups should be excluded from results (defaults to false). For example, when grouping on country and city, any cases where both country and city are null will be skipped if set to true.
    • isDebug

      public boolean isDebug()
      Indicates if window activity is logged at debug level (defaults to false).
    • setDebug

      public GroupByReader setDebug(boolean debug)
      Indicates if window activity is logged at debug level (defaults to false).
    • getOperations

      public List<GroupOperation<?>> getOperations()
    • add

      public GroupByReader add(GroupOperation<?> operation)
    • count

      public GroupByReader count(FieldPath targetFieldName)
      Adds a count of the records in each group, written to the target field.
    • count

      public GroupByReader count(String targetFieldName)
      Adds a count of the records in each group, written to the target field.
    • count

      public GroupByReader count(FieldPath targetFieldName, boolean excludeNulls)
    • count

      public GroupByReader count(String targetFieldName, boolean excludeNulls)
    • sum

      public GroupByReader sum(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the sum of the source field's non-null values as a BigDecimal, written to the target field.
    • sum

      public GroupByReader sum(FieldPath sourceFieldName)
      Adds the sum of the field's non-null values as a BigDecimal, written to the same field path.
    • sum

      public GroupByReader sum(String sourceFieldName, String targetFieldName)
      Adds the sum of the source field's non-null values as a BigDecimal, written to the target field.
    • sum

      public GroupByReader sum(String sourceFieldName)
      Adds the sum of the field's non-null values as a BigDecimal, written to a field of the same name.
    • avg

      public GroupByReader avg(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the average of the source field as a BigDecimal with 5 decimal places (HALF_UP), written to the target field; null values count as zero.
    • avg

      public GroupByReader avg(String sourceFieldName, String targetFieldName)
      Adds the average of the source field as a BigDecimal with 5 decimal places (HALF_UP), written to the target field; null values count as zero.
    • avg

      public GroupByReader avg(FieldPath sourceFieldName, FieldPath targetFieldName, int scale, RoundingMode roundingMode)
      Adds the average of the source field rounded to scale decimal places, written to the target field; null values count as zero.
    • avg

      public GroupByReader avg(String sourceFieldName, String targetFieldName, int scale, RoundingMode roundingMode)
      Adds the average of the source field rounded to scale decimal places, written to the target field; null values count as zero.
    • first

      public GroupByReader first(FieldPath sourceFieldName, FieldPath targetFieldName, boolean excludeNulls)
    • first

      public GroupByReader first(String sourceFieldName, String targetFieldName, boolean excludeNulls)
    • first

      public GroupByReader first(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the first non-null value of the source field in each group, written to the target field.
    • first

      public GroupByReader first(String sourceFieldName, String targetFieldName)
      Adds the first non-null value of the source field in each group, written to the target field.
    • first

      public GroupByReader first(FieldPath sourceFieldName, boolean excludeNulls)
      Adds the first value of the field in each group, written to the same field path.
    • first

      public GroupByReader first(String sourceFieldName, boolean excludeNulls)
      Adds the first value of the field in each group, written to a field of the same name.
    • last

      public GroupByReader last(FieldPath sourceFieldName, FieldPath targetFieldName, boolean excludeNulls)
    • last

      public GroupByReader last(String sourceFieldName, String targetFieldName, boolean excludeNulls)
    • last

      public GroupByReader last(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the last non-null value of the source field in each group, written to the target field.
    • last

      public GroupByReader last(String sourceFieldName, String targetFieldName)
      Adds the last non-null value of the source field in each group, written to the target field.
    • last

      public GroupByReader last(FieldPath sourceFieldName, boolean excludeNulls)
      Adds the last value of the field in each group, written to the same field path.
    • last

      public GroupByReader last(String sourceFieldName, boolean excludeNulls)
      Adds the last value of the field in each group, written to a field of the same name.
    • max

      public GroupByReader max(FieldPath sourceFieldName, FieldPath targetFieldName, boolean excludeNulls)
    • max

      public GroupByReader max(String sourceFieldName, String targetFieldName, boolean excludeNulls)
    • max

      public GroupByReader max(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the largest non-null value of the source field in each group, written to the target field.
    • max

      public GroupByReader max(String sourceFieldName, String targetFieldName)
      Adds the largest non-null value of the source field in each group, written to the target field.
    • max

      public GroupByReader max(FieldPath sourceFieldName, boolean excludeNulls)
      Adds the largest value of the field in each group, written to the same field path.
    • max

      public GroupByReader max(String sourceFieldName, boolean excludeNulls)
      Adds the largest value of the field in each group, written to a field of the same name.
    • max

      public GroupByReader max(FieldPath sourceFieldName)
      Adds the largest non-null value of the field in each group, written to the same field path.
    • max

      public GroupByReader max(String sourceFieldName)
      Adds the largest non-null value of the field in each group, written to a field of the same name.
    • min

      public GroupByReader min(FieldPath sourceFieldName, FieldPath targetFieldName, boolean excludeNulls)
    • min

      public GroupByReader min(String sourceFieldName, String targetFieldName, boolean excludeNulls)
    • min

      public GroupByReader min(FieldPath sourceFieldName, FieldPath targetFieldName)
      Adds the smallest non-null value of the source field in each group, written to the target field.
    • min

      public GroupByReader min(String sourceFieldName, String targetFieldName)
      Adds the smallest non-null value of the source field in each group, written to the target field.
    • min

      public GroupByReader min(FieldPath sourceFieldName, boolean excludeNulls)
    • min

      public GroupByReader min(String sourceFieldName, boolean excludeNulls)
    • min

      public GroupByReader min(FieldPath sourceFieldName)
      Adds the smallest non-null value of the field in each group, written to the same field path.
    • min

      public GroupByReader min(String sourceFieldName)
      Adds the smallest non-null value of the field in each group, written to a field of the same name.
    • collect

      public GroupByReader collect(FieldPath sourceFieldName, FieldPath targetFieldName, boolean distinctValues, boolean flattenArrayValues, boolean excludeNulls)
      Collects the individual values into an array.
      Parameters:
      sourceFieldName - the path of the field to collect values from
      targetFieldName - the path of the field where the collected array will be placed
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      excludeNulls - true to exclude null values from the collection
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • collect

      public GroupByReader collect(String sourceFieldName, String targetFieldName, boolean distinctValues, boolean flattenArrayValues, boolean excludeNulls)
      Collects the individual values into an array.
      Parameters:
      sourceFieldName - the name of the field to collect values from
      targetFieldName - the name of the field where the collected array will be placed
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      excludeNulls - true to exclude null values from the collection
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • collect

      public GroupByReader collect(FieldPath sourceFieldName, FieldPath targetFieldName, boolean distinctValues, boolean flattenArrayValues)
      Collects the individual values into an array, excluding nulls by default.
      Parameters:
      sourceFieldName - the path of the field to collect values from
      targetFieldName - the path of the field where the collected array will be placed
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • collect

      public GroupByReader collect(String sourceFieldName, String targetFieldName, boolean distinctValues, boolean flattenArrayValues)
      Collects the individual values into an array, excluding nulls by default.
      Parameters:
      sourceFieldName - the name of the field to collect values from
      targetFieldName - the name of the field where the collected array will be placed
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • collect

      public GroupByReader collect(FieldPath sourceFieldName, boolean excludeNulls, boolean distinctValues, boolean flattenArrayValues)
      Collects the individual values into an array, using the source field name as the target field name.
      Parameters:
      sourceFieldName - the path of the field to collect values from and place the result in
      excludeNulls - true to exclude null values from the collection
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • collect

      public GroupByReader collect(String sourceFieldName, boolean excludeNulls, boolean distinctValues, boolean flattenArrayValues)
      Collects the individual values into an array, using the source field name as the target field name.
      Parameters:
      sourceFieldName - the name of the field to collect values from and place the result in
      excludeNulls - true to exclude null values from the collection
      distinctValues - true to collect only distinct values, false to allow duplicates
      flattenArrayValues - true to flatten array values (add elements individually instead of as nested arrays)
      Returns:
      this GroupByReader instance for method chaining
      See Also:
    • getCloseWindowStrategy

      public CloseWindowStrategy getCloseWindowStrategy()
      Returns the strategy that decides when open windows are closed (defaults to CloseWindowStrategy.never()).
    • setCloseWindowStrategy

      public GroupByReader setCloseWindowStrategy(CloseWindowStrategy closeWindowStrategy)
      Sets the strategy that decides when open windows are closed; null restores the default CloseWindowStrategy.never(). Must be called before open().
    • getCreateWindowStrategy

      public CreateWindowStrategy getCreateWindowStrategy()
      Returns the strategy that decides when new windows are opened (defaults to limitOpened(1)).
    • setCreateWindowStrategy

      public GroupByReader setCreateWindowStrategy(CreateWindowStrategy createWindowStrategy)
      Sets the strategy that decides when new windows are opened; null restores the default limitOpened(1). Must be called before open().
    • getOpenedWindows

      public List<Window> getOpenedWindows()
      Returns the windows that are currently open and collecting records.
    • getClosedWindows

      public List<Window> getClosedWindows()
      Returns the closed windows whose summary records have not all been read yet.
    • getWindowsCreated

      public long getWindowsCreated()
    • getWindowsClosed

      public long getWindowsClosed()
    • open

      public void open()
      Description copied from class: DataEndpoint
      Makes this endpoint ready for reading or writing.
      Overrides:
      open in class ProxyReader
    • close

      public void close() throws DataException
      Description copied from class: DataEndpoint
      Indicates that this endpoint has finished reading or writing.
      Overrides:
      close in class ProxyReader
      Throws:
      DataException
    • readImpl

      protected Record readImpl() throws Throwable
      Description copied from class: DataReader
      Overridden by subclasses to read the next record from this DataReader. The default implementation of DataReader.read() now insures that this method will not be called again after it returns a null.

      If no record is available, null will be returned.

      Contract for subclasses (see also docs/authoring/DataReader.md):

      Overrides:
      readImpl in class ProxyReader
      Throws:
      Throwable
    • addExceptionProperties

      public DataException addExceptionProperties(DataException exception)
      Description copied from class: Endpoint
      Adds this endpoint's current state to a DataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.
      Overrides:
      addExceptionProperties in class ProxyReader
    • toString

      public String toString()
      Overrides:
      toString in class DataEndpoint