Class GroupByReader
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.
-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
Fields inherited from class com.northconcepts.datapipeline.core.DataReader
fieldLineage, recordLineageFields inherited from class com.northconcepts.datapipeline.core.DataEndpoint
lastRecord, PRODUCT, PRODUCT_VERSION, VENDOR, XML_INPUT_FACTORY_KEYFields inherited from class com.northconcepts.datapipeline.core.Endpoint
BUFFER_SIZE, captureElapsedTime, DEFAULT_READ_BUFFER_SIZEFields inherited from class com.northconcepts.datapipeline.core.DataObject
id, log, name, TIMESTAMP_FORMAT -
Constructor Summary
ConstructorsConstructorDescriptionGroupByReader(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 toqueueSizerecords.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 toqueueSizerecords.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(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. -
Method Summary
Modifier and TypeMethodDescriptionadd(GroupOperation<?> operation) addExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.Adds the average of the source field as aBigDecimalwith 5 decimal places (HALF_UP), written to the target field; null values count as zero.avg(FieldPath sourceFieldName, FieldPath targetFieldName, int scale, RoundingMode roundingMode) Adds the average of the source field rounded toscaledecimal places, written to the target field; null values count as zero.Adds the average of the source field as aBigDecimalwith 5 decimal places (HALF_UP), written to the target field; null values count as zero.avg(String sourceFieldName, String targetFieldName, int scale, RoundingMode roundingMode) Adds the average of the source field rounded toscaledecimal places, written to the target field; null values count as zero.voidclose()Indicates that this endpoint has finished reading or writing.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.collect(FieldPath sourceFieldName, FieldPath targetFieldName, boolean distinctValues, boolean flattenArrayValues) Collects the individual values into an array, excluding nulls by default.collect(FieldPath sourceFieldName, FieldPath targetFieldName, boolean distinctValues, boolean flattenArrayValues, boolean excludeNulls) Collects the individual values into an array.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.collect(String sourceFieldName, String targetFieldName, boolean distinctValues, boolean flattenArrayValues) Collects the individual values into an array, excluding nulls by default.collect(String sourceFieldName, String targetFieldName, boolean distinctValues, boolean flattenArrayValues, boolean excludeNulls) Collects the individual values into an array.Adds a count of the records in each group, written to the target field.Adds a count of the records in each group, written to the target field.Adds the first value of the field in each group, written to the same field path.Adds the first non-null value of the source field in each group, written to the target field.Adds the first value of the field in each group, written to a field of the same name.Adds the first non-null value of the source field in each group, written to the target field.Returns the closed windows whose summary records have not all been read yet.Returns the strategy that decides when open windows are closed (defaults toCloseWindowStrategy.never()).Returns the strategy that decides when new windows are opened (defaults tolimitOpened(1)).Returns the windows that are currently open and collecting records.List<GroupOperation<?>> longlongbooleanisDebug()Indicates if window activity is logged at debug level (defaults to false).booleanIndicates ifnullgroups should be excluded from results (defaults to false).Adds the last value of the field in each group, written to the same field path.Adds the last non-null value of the source field in each group, written to the target field.Adds the last value of the field in each group, written to a field of the same name.Adds the last non-null value of the source field in each group, written to the target field.Adds the largest non-null value of the field in each group, written to the same field path.Adds the largest value of the field in each group, written to the same field path.Adds the largest non-null value of the source field in each group, written to the target field.Adds the largest non-null value of the field in each group, written to a field of the same name.Adds the largest value of the field in each group, written to a field of the same name.Adds the largest non-null value of the source field in each group, written to the target field.Adds the smallest non-null value of the field in each group, written to the same field path.Adds the smallest non-null value of the source field in each group, written to the target field.Adds the smallest non-null value of the field in each group, written to a field of the same name.Adds the smallest non-null value of the source field in each group, written to the target field.voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setCloseWindowStrategy(CloseWindowStrategy closeWindowStrategy) Sets the strategy that decides when open windows are closed; null restores the defaultCloseWindowStrategy.never().setCreateWindowStrategy(CreateWindowStrategy createWindowStrategy) Sets the strategy that decides when new windows are opened; null restores the defaultlimitOpened(1).setDebug(boolean debug) Indicates if window activity is logged at debug level (defaults to false).setExcludeNulls(boolean excludeNulls) Indicates ifnullgroups should be excluded from results (defaults to false).Adds the sum of the field's non-null values as aBigDecimal, written to the same field path.Adds the sum of the source field's non-null values as aBigDecimal, written to the target field.Adds the sum of the field's non-null values as aBigDecimal, written to a field of the same name.Adds the sum of the source field's non-null values as aBigDecimal, written to the target field.toString()Methods inherited from class com.northconcepts.datapipeline.core.ProxyReader
available, getNestedReader, interceptRecord, map, map, setNestedDataReader, setNestedDataReaderMethods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, getBufferSize, getNestedEndpoint, getReader, getRootEndpoint, getRootReader, isExhausted, isLineageSupported, isSaveLineage, peek, pop, push, read, setSaveLineage, skipMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, resetRecordCountMethods inherited from class com.northconcepts.datapipeline.core.Endpoint
addElapsedtime, assertClosed, assertNotOpened, assertOpened, finalize, getClosedOn, getDescription, getElapsedTime, getElapsedTimeAsString, getOpenedOn, getOpenElapsedTime, getOpenElapsedTimeAsString, getSelfTime, getSelfTimeAsString, getState, isCaptureElapsedTime, isClosed, isOpen, setCaptureElapsedTime, setDescription
-
Constructor Details
-
GroupByReader
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
Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up toqueueSizerecords. -
GroupByReader
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
Creates a reader that groups records by the given fields; the source is read on a background thread that buffers up toqueueSizerecords.
-
-
Method Details
-
getGroupByFields
-
isExcludeNulls
public boolean isExcludeNulls()Indicates ifnullgroups should be excluded from results (defaults to false). -
setExcludeNulls
Indicates ifnullgroups 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 totrue. -
isDebug
public boolean isDebug()Indicates if window activity is logged at debug level (defaults to false). -
setDebug
Indicates if window activity is logged at debug level (defaults to false). -
getOperations
-
add
-
count
Adds a count of the records in each group, written to the target field. -
count
Adds a count of the records in each group, written to the target field. -
count
-
count
-
sum
Adds the sum of the source field's non-null values as aBigDecimal, written to the target field. -
sum
Adds the sum of the field's non-null values as aBigDecimal, written to the same field path. -
sum
Adds the sum of the source field's non-null values as aBigDecimal, written to the target field. -
sum
Adds the sum of the field's non-null values as aBigDecimal, written to a field of the same name. -
avg
Adds the average of the source field as aBigDecimalwith 5 decimal places (HALF_UP), written to the target field; null values count as zero. -
avg
Adds the average of the source field as aBigDecimalwith 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 toscaledecimal 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 toscaledecimal places, written to the target field; null values count as zero. -
first
public GroupByReader first(FieldPath sourceFieldName, FieldPath targetFieldName, boolean excludeNulls) -
first
-
first
Adds the first non-null value of the source field in each group, written to the target field. -
first
Adds the first non-null value of the source field in each group, written to the target field. -
first
Adds the first value of the field in each group, written to the same field path. -
first
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
-
last
Adds the last non-null value of the source field in each group, written to the target field. -
last
Adds the last non-null value of the source field in each group, written to the target field. -
last
Adds the last value of the field in each group, written to the same field path. -
last
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
-
max
Adds the largest non-null value of the source field in each group, written to the target field. -
max
Adds the largest non-null value of the source field in each group, written to the target field. -
max
Adds the largest value of the field in each group, written to the same field path. -
max
Adds the largest value of the field in each group, written to a field of the same name. -
max
Adds the largest non-null value of the field in each group, written to the same field path. -
max
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
-
min
Adds the smallest non-null value of the source field in each group, written to the target field. -
min
Adds the smallest non-null value of the source field in each group, written to the target field. -
min
-
min
-
min
Adds the smallest non-null value of the field in each group, written to the same field path. -
min
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 fromtargetFieldName- the path of the field where the collected array will be placeddistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- 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 fromtargetFieldName- the name of the field where the collected array will be placeddistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- 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 fromtargetFieldName- the path of the field where the collected array will be placeddistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- 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 fromtargetFieldName- the name of the field where the collected array will be placeddistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- 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 inexcludeNulls- true to exclude null values from the collectiondistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- 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 inexcludeNulls- true to exclude null values from the collectiondistinctValues- true to collect only distinct values, false to allow duplicatesflattenArrayValues- true to flatten array values (add elements individually instead of as nested arrays)- Returns:
- this GroupByReader instance for method chaining
- See Also:
-
getCloseWindowStrategy
Returns the strategy that decides when open windows are closed (defaults toCloseWindowStrategy.never()). -
setCloseWindowStrategy
Sets the strategy that decides when open windows are closed; null restores the defaultCloseWindowStrategy.never(). Must be called beforeopen(). -
getCreateWindowStrategy
Returns the strategy that decides when new windows are opened (defaults tolimitOpened(1)). -
setCreateWindowStrategy
Sets the strategy that decides when new windows are opened; null restores the defaultlimitOpened(1). Must be called beforeopen(). -
getOpenedWindows
Returns the windows that are currently open and collecting records. -
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:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classProxyReader
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classProxyReader- Throws:
DataException
-
readImpl
Description copied from class:DataReaderOverridden by subclasses to read the next record from thisDataReader. The default implementation ofDataReader.read()now insures that this method will not be called again after it returns anull.If no record is available,
nullwill be returned.Contract for subclasses (see also
docs/authoring/DataReader.md):- Return
nullexactly once to signal end-of-stream. - Do not call this method directly —
DataReader.read()is the template method that wraps exceptions, increments the record count, and tracks lineage. - Do not call
DataEndpoint.incrementRecordCount()here;DataReader.read()already does. - Throw raw exceptions;
DataReader.read()wraps them viaexception(throwable).
- Overrides:
readImplin classProxyReader- Throws:
Throwable
- Return
-
addExceptionProperties
Description copied from class:EndpointAdds this endpoint's current state to aDataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.- Overrides:
addExceptionPropertiesin classProxyReader
-
toString
- Overrides:
toStringin classDataEndpoint
-