Class ParquetDataReader
java.lang.Object
com.northconcepts.datapipeline.core.DataObject
com.northconcepts.datapipeline.core.Endpoint
com.northconcepts.datapipeline.core.DataEndpoint
com.northconcepts.datapipeline.core.DataReader
com.northconcepts.datapipeline.internal.lang.IntegrationReader
com.northconcepts.datapipeline.parquet.ParquetDataReader
Read records from Apache Parquet columnar files.
See Apache Parquet columnar storage.
-
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
ConstructorsConstructorDescriptionParquetDataReader(InputFile inputFile) Reads parquet data from anInputFile. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.protected RecordaddLineage(Record record) Called byDataReader.read()for each record fromDataReader.readImpl()while lineage is saved; the default copiesrecordLineageinto every field along with its original index and name.voidclose()Indicates that this endpoint has finished reading or writing.ConfigurationReturns the Parquet configuration parameters.MetadataFilterReturns the filter settings.MessageTypeIndicates the modified schema used to read the file.MessageTypeIndicates the schema used to read the file.protected Dateint96ToTimestamp(byte[] int96Bytes) Converts a 12-byte Parquet INT96 timestamp to aDatein the system default time zone.booleanisDebug()Indicates if debugging is enabled to print log statements.booleanIndicates if this reader can capture record and field lineage (false unless overridden by a reader that supports it).booleanUsed whether to check foroptionalcolumns that can be maderequiredusingduringinvalid reference
ColumnChunkMetaData#getStatistics()open().booleanWhen true, check forrequiredcolumns that can be madeoptionalusing.invalid reference
ColumnChunkMetaData#getStatistics()booleanWhen true, remove fields in the schema does not have a correspondingColumnChunkMetaDatain the file.booleanWhen true, remove fields in the schema whose value count is invalid input: '<'= 0 SeegetFieldsWithoutColumnMetadata()for details.voidopen()Makes this endpoint ready for reading or writing.protected voidreadGroupValue(Group group, int depth, int fieldIndex, ValueNodeContainer valueContainer) Adds the value of the group's field atfieldIndexto the container: primitives as values, lists as arrays, and maps and other groups as records.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.protected voidreadPrimitiveValue(Group group, int fieldIndex, ValueNodeContainer valueContainer, Type fieldType) Adds the values of the group's primitive field atfieldIndexto the container, converted using its logical type, or a typed null if it has none; times and timestamps are truncated to milliseconds.protected RecordreadRecord(Group group, int depth) Converts a Parquet group into a record with one field per group field;depthis 0 for a row and grows with nesting.setConfiguration(Configuration configuration) Sets the Parquet configuration parameters.setDebug(boolean debug) Indicates if debugging is enabled to print log statements.setFilter(MetadataFilter filter) Sets the filter settings.setMakeOptionalFieldsRequired(boolean makeOptionalFieldsRequired) Used whether to check foroptionalcolumns that can be maderequiredusingduringinvalid reference
ColumnChunkMetaData#getStatistics()open().setMakeRequiredFieldsOptional(boolean makeRequiredFieldsOptional) When true, check forrequiredcolumns that can be madeoptionalusing.invalid reference
ColumnChunkMetaData#getStatistics()setRemoveFieldsWithoutColumnMetadata(boolean removeFieldsWithoutColumnMetadata) Used whether to remove fields in the schema does not have a correspondingColumnChunkMetaDatain the file duringopen().setRemoveFieldsWithoutValues(boolean removeFieldsWithoutValues) Used whether to remove fields in the schema whose value count is invalid input: '<'= 0 duringopen().setSchema(MessageType schema) Indicates the schema used to read the file.Methods inherited from class com.northconcepts.datapipeline.core.DataReader
available, getBufferSize, getNestedEndpoint, getNestedReader, getReader, getRootEndpoint, getRootReader, isExhausted, isSaveLineage, peek, pop, push, read, setSaveLineage, skipMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, resetRecordCount, toStringMethods 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
-
ParquetDataReader
public ParquetDataReader(InputFile inputFile) Reads parquet data from anInputFile.
-
-
Method Details
-
isDebug
public boolean isDebug()Indicates if debugging is enabled to print log statements. -
setDebug
Indicates if debugging is enabled to print log statements. -
getConfiguration
public Configuration getConfiguration()Returns the Parquet configuration parameters. -
setConfiguration
Sets the Parquet configuration parameters. -
getFilter
public MetadataFilter getFilter()Returns the filter settings. -
setFilter
Sets the filter settings. -
getSchema
public MessageType getSchema()Indicates the schema used to read the file. -
getModifiedSchema
public MessageType getModifiedSchema()Indicates the modified schema used to read the file. -
setSchema
Indicates the schema used to read the file. -
isMakeRequiredFieldsOptional
public boolean isMakeRequiredFieldsOptional()When true, check forrequiredcolumns that can be madeoptionalusing. Seeinvalid reference
ColumnChunkMetaData#getStatistics()getRequiredColumnsWithNullValues()for details. -
setMakeRequiredFieldsOptional
When true, check forrequiredcolumns that can be madeoptionalusing. Seeinvalid reference
ColumnChunkMetaData#getStatistics()getRequiredColumnsWithNullValues()for details. -
isMakeOptionalFieldsRequired
public boolean isMakeOptionalFieldsRequired()Used whether to check foroptionalcolumns that can be maderequiredusingduringinvalid reference
ColumnChunkMetaData#getStatistics()open(). SeegetOptionalColumnsWithoutNullValues()()} for details. -
setMakeOptionalFieldsRequired
Used whether to check foroptionalcolumns that can be maderequiredusingduringinvalid reference
ColumnChunkMetaData#getStatistics()open(). SeegetOptionalColumnsWithoutNullValues()()} for details. -
isRemoveFieldsWithoutColumnMetadata
public boolean isRemoveFieldsWithoutColumnMetadata()When true, remove fields in the schema does not have a correspondingColumnChunkMetaDatain the file. SeegetFieldsWithoutColumnMetadata()for details. -
setRemoveFieldsWithoutColumnMetadata
public ParquetDataReader setRemoveFieldsWithoutColumnMetadata(boolean removeFieldsWithoutColumnMetadata) Used whether to remove fields in the schema does not have a correspondingColumnChunkMetaDatain the file duringopen(). SeegetFieldsWithoutColumnMetadata()for details. -
isRemoveFieldsWithoutValues
public boolean isRemoveFieldsWithoutValues()When true, remove fields in the schema whose value count is invalid input: '<'= 0 SeegetFieldsWithoutColumnMetadata()for details. -
setRemoveFieldsWithoutValues
Used whether to remove fields in the schema whose value count is invalid input: '<'= 0 duringopen(). SeegetFieldsWithoutColumnMetadata()for details. -
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classIntegrationReader- Throws:
DataException
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classDataEndpoint- 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).
- Specified by:
readImplin classDataReader- Throws:
Throwable
- Return
-
readRecord
Converts a Parquet group into a record with one field per group field;depthis 0 for a row and grows with nesting. -
readGroupValue
protected void readGroupValue(Group group, int depth, int fieldIndex, ValueNodeContainer valueContainer) Adds the value of the group's field atfieldIndexto the container: primitives as values, lists as arrays, and maps and other groups as records. -
readPrimitiveValue
protected void readPrimitiveValue(Group group, int fieldIndex, ValueNodeContainer valueContainer, Type fieldType) Adds the values of the group's primitive field atfieldIndexto the container, converted using its logical type, or a typed null if it has none; times and timestamps are truncated to milliseconds. -
int96ToTimestamp
Converts a 12-byte Parquet INT96 timestamp to aDatein the system default time zone. -
isLineageSupported
public boolean isLineageSupported()Description copied from class:DataReaderIndicates if this reader can capture record and field lineage (false unless overridden by a reader that supports it).- Overrides:
isLineageSupportedin classDataReader
-
addLineage
Description copied from class:DataReaderCalled byDataReader.read()for each record fromDataReader.readImpl()while lineage is saved; the default copiesrecordLineageinto every field along with its original index and name. Overrides set their source details onrecordLineagefirst and end withsuper.addLineage(record).- Overrides:
addLineagein classDataReader
-
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 classDataReader
-