Class Pipeline
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
com.northconcepts.datapipeline.foundations.pipeline.Pipeline
- All Implemented Interfaces:
DataExceptionContributor,JsonSerializable,RecordSerializable,XmlSerializable,DataReaderFactory,DataWriterFactory,JavaCodeGenerator,Serializable
A declarative pipeline that applies an ordered list of
PipelineActions to its input records, with undo and redo of
the most recent actions.- 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 TypeMethodDescriptionaddAction(PipelineAction action) Appends an action, attaches it to this pipeline and clears the redo history.addAction(Integer index, PipelineAction action) Inserts an action at the given 0-based index (or appends it if the index is null), attaches it to this pipeline and clears the redo history.protected DataReaderapplyDataProcessingImpl(DataReader reader) Wraps the reader with this pipeline's own processing; called between the source and target entity transformers.fromJson(InputStream inputStream) Loads this object's state from the UTF-8 encoded JSON in the stream.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).fromXml(InputStream inputStream) Parses the XML in the stream and loads this object's state from its root element.fromXmlElement(Element pipelineElement) voidWrites the Java code for this pipeline's own processing, between the input and output code.intReturns a read-only view of the actions, in the order they are applied.intReturns a read-only view of the actions removed byundoLastAction()that can still be redone.Re-appends the most recently undone action; returns it, or null if there is nothing to redo.setAction(int index, PipelineAction action) Replaces the action at the given 0-based index.setActions(List<PipelineAction> actions) Replaces all the actions with the given list and clears the redo history.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).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.Removes the last action and keeps it forredoLastAction(); returns it, or null if there are no actions.Methods inherited from class com.northconcepts.datapipeline.foundations.pipeline.AbstractPipeline
applyDataProcessing, applySourceEntityTransformer, applyTargetEntityTransformer, createDataReader, createDataWriter, createJob, generateJavaCode, generateJavaCodePostProcess, generateJavaCodePreProcess, getDateTimePatternDetector, getDescription, getDiscardReasonFieldName, getDiscardWriter, getInput, getJavaCode, getName, getNumberDetector, getOutput, getSourceEntity, getTargetEntity, isMultithreaded, run, runAsync, setDateTimePatternDetector, setNumberDetectorMethods 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
toJson, toJson, toJson
-
Constructor Details
-
Pipeline
public Pipeline()
-
-
Method Details
-
getActions
Returns a read-only view of the actions, in the order they are applied. -
getUndoneActions
Returns a read-only view of the actions removed byundoLastAction()that can still be redone. -
getActionCount
public int getActionCount() -
getUndoneActionCount
public int getUndoneActionCount() -
addAction
Appends an action, attaches it to this pipeline and clears the redo history. -
addAction
Inserts an action at the given 0-based index (or appends it if the index is null), attaches it to this pipeline and clears the redo history. -
setAction
Replaces the action at the given 0-based index. -
setActions
Replaces all the actions with the given list and clears the redo history. -
undoLastAction
Removes the last action and keeps it forredoLastAction(); returns it, or null if there are no actions. -
redoLastAction
Re-appends the most recently undone action; returns it, or null if there is nothing to redo. -
generateJavaCodeImpl
Description copied from class:AbstractPipelineWrites the Java code for this pipeline's own processing, between the input and output code.- Specified by:
generateJavaCodeImplin classAbstractPipeline
-
applyDataProcessingImpl
Description copied from class:AbstractPipelineWraps the reader with this pipeline's own processing; called between the source and target entity transformers.- Specified by:
applyDataProcessingImplin classAbstractPipeline
-
toRecord
Description copied from interface:RecordSerializableConverts this object's state to a record thatRecordSerializable.fromRecord(Record)can load.- Specified by:
toRecordin interfaceRecordSerializable- Overrides:
toRecordin classAbstractPipeline
-
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 classAbstractPipeline- Parameters:
source-- Returns:
- this instance.
-
toXmlElement
Description copied from interface:XmlSerializableReturns an element, created withdocument, describing this object; the default implementation throws aDataException.- Specified by:
toXmlElementin interfaceXmlSerializable- Overrides:
toXmlElementin classAbstractPipeline
-
fromXmlElement
- Specified by:
fromXmlElementin interfaceXmlSerializable- Overrides:
fromXmlElementin classAbstractPipeline
-
fromXml
Description copied from interface:XmlSerializableParses the XML in the stream and loads this object's state from its root element. -
fromXml
-
fromJson
Description copied from interface:JsonSerializableLoads this object's state from the UTF-8 encoded JSON in the stream. -
fromJson
Description copied from interface:JsonSerializableLoads this object's state from the given JSON text.- Specified by:
fromJsonin interfaceJsonSerializable- Specified by:
fromJsonin interfaceRecordSerializable- Overrides:
fromJsonin classAbstractPipeline
-
setName
- Overrides:
setNamein classAbstractPipeline
-
setDescription
- Overrides:
setDescriptionin classAbstractPipeline
-
setInput
- Overrides:
setInputin classAbstractPipeline
-
setInputAsDataReaderFactory
Description copied from class:AbstractPipelineSets the input to aDataReaderPipelineInputthat creates its readers with the given factory.- Overrides:
setInputAsDataReaderFactoryin classAbstractPipeline
-
setInputAsDataReader
Description copied from class:AbstractPipelineSets the input to aDataReaderPipelineInputthat always returns the given reader instance.- Overrides:
setInputAsDataReaderin classAbstractPipeline
-
setOutput
- Overrides:
setOutputin classAbstractPipeline
-
setOutputAsDataWriterFactory
Description copied from class:AbstractPipelineSets the output to aDataWriterPipelineOutputthat creates its writers with the given factory.- Overrides:
setOutputAsDataWriterFactoryin classAbstractPipeline
-
setOutputAsDataWriter
Description copied from class:AbstractPipelineSets the output to aDataWriterPipelineOutputthat always returns the given writer instance.- Overrides:
setOutputAsDataWriterin classAbstractPipeline
-
setSourceEntity
Description copied from class:AbstractPipelineSets the schema entity used to convert and validate records read from the input, before any processing (null for none).- Overrides:
setSourceEntityin classAbstractPipeline
-
setTargetEntity
Description copied from class:AbstractPipelineSets the schema entity used to convert and validate records after processing, before they are written (null for none).- Overrides:
setTargetEntityin classAbstractPipeline
-
setMultithreaded
Description copied from class:AbstractPipelineIndicates if each processing stage runs in its own thread using anAsyncReader(defaults to true).- Overrides:
setMultithreadedin classAbstractPipeline
-
setDiscardWriter
Description copied from class:AbstractPipelineSets the writer that receives records failing source or target entity validation, instead of stopping the pipeline.- Overrides:
setDiscardWriterin classAbstractPipeline
-
setDiscardReasonFieldName
Description copied from class:AbstractPipelineSets the name of the field added to discarded records to hold the failure message; only allowed with a discard writer.- Overrides:
setDiscardReasonFieldNamein classAbstractPipeline
-