Class AbstractPipeline
java.lang.Object
com.northconcepts.datapipeline.foundations.core.Bean
com.northconcepts.datapipeline.foundations.core.FoundationObject
com.northconcepts.datapipeline.foundations.pipeline.PipelineObject
com.northconcepts.datapipeline.foundations.pipeline.AbstractPipeline
- All Implemented Interfaces:
DataExceptionContributor,JsonSerializable,RecordSerializable,XmlSerializable,DataReaderFactory,DataWriterFactory,JavaCodeGenerator,Serializable
- Direct Known Subclasses:
DataMappingPipeline,NullPipeline,Pipeline
public abstract class AbstractPipeline
extends PipelineObject
implements DataReaderFactory, DataWriterFactory
Base class for declarative pipelines that read records from a
PipelineInput, process them and write them to a
PipelineOutput.- See Also:
-
Field Summary
Fields inherited from class com.northconcepts.datapipeline.foundations.core.FoundationObject
internalId, internalName, log, TIMESTAMP_FORMATFields inherited from interface com.northconcepts.datapipeline.core.RecordSerializable
SERIALIZED_CLASS_NAME, TYPEFields inherited from interface com.northconcepts.datapipeline.core.XmlSerializable
XML_SERIALIZED_CLASS_NAME -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionfinal DataReaderapplyDataProcessing(DataReader reader) Wraps the reader with the source entity transformer, this pipeline's own processing and the target entity transformer, in that order.protected abstract DataReaderapplyDataProcessingImpl(DataReader reader) Wraps the reader with this pipeline's own processing; called between the source and target entity transformers.protected DataReaderWraps the reader to convert and validate records against the source entity, if one is set, in its own thread when multithreaded.protected DataReaderWraps the reader to convert and validate records against the target entity, if one is set, in its own thread when multithreaded.Creates, but does not run, a job that reads the processed input and writes it to the output.Loads this object's state from the given JSON text.fromRecord(Record source) Loads this instance's state from a record and returnsthis(for fluid API call chaining).fromXmlElement(Element pipelineElement) voidAppends Java source code representing this object to the given builder.protected abstract voidWrites the Java code for this pipeline's own processing, between the input and output code.protected voidWrites the output's Java code and, when both an input and an output are set, the code that runs the job.protected voidWrites the input's Java code; called beforegenerateJavaCodeImpl(JavaCodeBuilder).Returns the detector aDatasetbuilt on this pipeline uses to infer date/time values in text columns.Returns the name of the field added to discarded records to hold the failure message, or null if none.Returns the writer that receives records failing source or target entity validation, instead of stopping the pipeline, or null if none.getInput()Returns a Java program, aDataPipelineDesktopclass whosemainmethod holds this pipeline's generated code.getName()Returns the detector aDatasetbuilt on this pipeline uses to infer numeric values in text columns.Returns the schema entity used to convert and validate records read from the input, before any processing, or null if none.Returns the schema entity used to convert and validate records after processing, before they are written, or null if none.booleanIndicates if each processing stage runs in its own thread using anAsyncReader(defaults to true).run()Runs this pipeline in the current thread and returns the finished job.runAsync()Starts this pipeline in a new thread and returns its job without waiting for it to finish.setDateTimePatternDetector(DateTimePatternDetector dateTimePatternDetector) Sets the detector aDatasetbuilt on this pipeline uses to infer date/time values in text columns; null restores the default detector.setDescription(String description) setDiscardReasonFieldName(String discardReasonFieldName) Sets the name of the field added to discarded records to hold the failure message; only allowed with a discard writer.setDiscardWriter(DataWriter discardWriter) Sets the writer that receives records failing source or target entity validation, instead of stopping the pipeline.setInput(PipelineInput input) setInputAsDataReader(DataReader reader) Sets the input to aDataReaderPipelineInputthat always returns the given reader instance.Sets the input to aDataReaderPipelineInputthat creates its readers with the given factory.setMultithreaded(boolean multithreaded) Indicates if each processing stage runs in its own thread using anAsyncReader(defaults to true).setNumberDetector(NumberDetector numberDetector) Sets the detector aDatasetbuilt on this pipeline uses to infer numeric values in text columns; null restores the default detector.setOutput(PipelineOutput output) setOutputAsDataWriter(DataWriter writer) Sets the output to aDataWriterPipelineOutputthat always returns the given writer instance.Sets the output to aDataWriterPipelineOutputthat creates its writers with the given factory.setSourceEntity(EntityDef sourceEntity) Sets the schema entity used to convert and validate records read from the input, before any processing (null for none).setTargetEntity(EntityDef targetEntity) Sets the schema entity used to convert and validate records after processing, before they are written (null for none).toRecord()Converts this object's state to a record thatRecordSerializable.fromRecord(Record)can load.toXmlElement(Document document) Returns an element, created withdocument, describing this object; the default implementation throws aDataException.Methods inherited from class com.northconcepts.datapipeline.foundations.core.FoundationObject
addExceptionProperties, assertValid, assertValid, clone, exception, exception, exception, getInternalId, getInternalName, resetInternalIdMethods inherited from class java.lang.Object
equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, waitMethods inherited from interface com.northconcepts.datapipeline.core.DataExceptionContributor
addExceptionPropertiesMethods inherited from interface com.northconcepts.datapipeline.core.RecordSerializable
fromJson, toJson, toJson, toJson
-
Constructor Details
-
AbstractPipeline
public AbstractPipeline()
-
-
Method Details
-
generateJavaCode
Description copied from interface:JavaCodeGeneratorAppends Java source code representing this object to the given builder.- Specified by:
generateJavaCodein interfaceJavaCodeGenerator
-
generateJavaCodePreProcess
Writes the input's Java code; called beforegenerateJavaCodeImpl(JavaCodeBuilder). -
generateJavaCodeImpl
Writes the Java code for this pipeline's own processing, between the input and output code. -
generateJavaCodePostProcess
Writes the output's Java code and, when both an input and an output are set, the code that runs the job. -
getJavaCode
Returns a Java program, aDataPipelineDesktopclass whosemainmethod holds this pipeline's generated code. -
getName
-
setName
-
getDescription
-
setDescription
-
getInput
-
setInput
-
setInputAsDataReaderFactory
Sets the input to aDataReaderPipelineInputthat creates its readers with the given factory. -
setInputAsDataReader
Sets the input to aDataReaderPipelineInputthat always returns the given reader instance. -
getSourceEntity
Returns the schema entity used to convert and validate records read from the input, before any processing, or null if none. -
setSourceEntity
Sets the schema entity used to convert and validate records read from the input, before any processing (null for none). -
getTargetEntity
Returns the schema entity used to convert and validate records after processing, before they are written, or null if none. -
setTargetEntity
Sets the schema entity used to convert and validate records after processing, before they are written (null for none). -
getOutput
-
setOutput
-
setOutputAsDataWriterFactory
Sets the output to aDataWriterPipelineOutputthat creates its writers with the given factory. -
setOutputAsDataWriter
Sets the output to aDataWriterPipelineOutputthat always returns the given writer instance. -
isMultithreaded
public boolean isMultithreaded()Indicates if each processing stage runs in its own thread using anAsyncReader(defaults to true). -
setMultithreaded
Indicates if each processing stage runs in its own thread using anAsyncReader(defaults to true). -
getDiscardWriter
Returns the writer that receives records failing source or target entity validation, instead of stopping the pipeline, or null if none. -
setDiscardWriter
Sets the writer that receives records failing source or target entity validation, instead of stopping the pipeline. -
getDiscardReasonFieldName
Returns the name of the field added to discarded records to hold the failure message, or null if none. -
setDiscardReasonFieldName
Sets the name of the field added to discarded records to hold the failure message; only allowed with a discard writer. -
getDateTimePatternDetector
Returns the detector aDatasetbuilt on this pipeline uses to infer date/time values in text columns. -
setDateTimePatternDetector
Sets the detector aDatasetbuilt on this pipeline uses to infer date/time values in text columns; null restores the default detector. -
getNumberDetector
Returns the detector aDatasetbuilt on this pipeline uses to infer numeric values in text columns. -
setNumberDetector
Sets the detector aDatasetbuilt on this pipeline uses to infer numeric values in text columns; null restores the default detector. -
applySourceEntityTransformer
Wraps the reader to convert and validate records against the source entity, if one is set, in its own thread when multithreaded. -
applyTargetEntityTransformer
Wraps the reader to convert and validate records against the target entity, if one is set, in its own thread when multithreaded. -
createDataReader
- Specified by:
createDataReaderin interfaceDataReaderFactory
-
applyDataProcessing
Wraps the reader with the source entity transformer, this pipeline's own processing and the target entity transformer, in that order. -
applyDataProcessingImpl
Wraps the reader with this pipeline's own processing; called between the source and target entity transformers. -
createDataWriter
- Specified by:
createDataWriterin interfaceDataWriterFactory
-
createJob
Creates, but does not run, a job that reads the processed input and writes it to the output. -
run
Runs this pipeline in the current thread and returns the finished job. -
runAsync
Starts this pipeline in a new thread and returns its job without waiting for it to finish. -
toRecord
Description copied from interface:RecordSerializableConverts this object's state to a record thatRecordSerializable.fromRecord(Record)can load.- Specified by:
toRecordin interfaceRecordSerializable- Overrides:
toRecordin classBean
-
fromRecord
Description copied from interface:RecordSerializableLoads this instance's state from a record and returnsthis(for fluid API call chaining). For fluid API call chaining, the overridden method should change the declared return type to its class.- Specified by:
fromRecordin interfaceRecordSerializable- Overrides:
fromRecordin classBean- Parameters:
source-- Returns:
- this instance.
-
fromJson
Description copied from interface:JsonSerializableLoads this object's state from the given JSON text.- Specified by:
fromJsonin interfaceJsonSerializable- Specified by:
fromJsonin interfaceRecordSerializable
-
toXmlElement
Description copied from interface:XmlSerializableReturns an element, created withdocument, describing this object; the default implementation throws aDataException.- Specified by:
toXmlElementin interfaceXmlSerializable
-
fromXmlElement
- Specified by:
fromXmlElementin interfaceXmlSerializable
-