Class AvroReader
Read an Apache Avro file and convert the contents into
Records .-
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
ConstructorsConstructorDescriptionAvroReader(File file) Read an Avro file.AvroReader(InputStream inputStream) Read a stream of Avro data. -
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.voidclose()Indicates that this endpoint has finished reading or writing.intReturns the total number of broken Avro records this reader encountered.intIndicates the number of broken Avro records to skip before this reader should fail (default 1000).SchemaReturns the Avro schema of the data being read or null if the reader has not been opened.voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setMaxInvalidRecords(int maxInvalidRecords) Indicates the number of broken Avro records to skip before this reader should fail (default 1000).Methods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, available, getBufferSize, getNestedEndpoint, getNestedReader, 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, 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
-
AvroReader
Read an Avro file.- Parameters:
file- the file to be read
-
AvroReader
Read a stream of Avro data.- Parameters:
inputStream-
-
-
Method 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
-
getMaxInvalidRecords
public int getMaxInvalidRecords()Indicates the number of broken Avro records to skip before this reader should fail (default 1000). -
setMaxInvalidRecords
Indicates the number of broken Avro records to skip before this reader should fail (default 1000). -
getInvalidRecordCount
public int getInvalidRecordCount()Returns the total number of broken Avro records this reader encountered. -
getSchema
public Schema getSchema()Returns the Avro schema of the data being read or null if the reader has not been opened. -
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
-