Class 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:
  • Constructor Details

    • AbstractPipeline

      public AbstractPipeline()
  • Method Details

    • generateJavaCode

      public void generateJavaCode(JavaCodeBuilder code)
      Description copied from interface: JavaCodeGenerator
      Appends Java source code representing this object to the given builder.
      Specified by:
      generateJavaCode in interface JavaCodeGenerator
    • generateJavaCodePreProcess

      protected void generateJavaCodePreProcess(JavaCodeBuilder code)
      Writes the input's Java code; called before generateJavaCodeImpl(JavaCodeBuilder).
    • generateJavaCodeImpl

      protected abstract void generateJavaCodeImpl(JavaCodeBuilder code)
      Writes the Java code for this pipeline's own processing, between the input and output code.
    • generateJavaCodePostProcess

      protected void generateJavaCodePostProcess(JavaCodeBuilder code)
      Writes the output's Java code and, when both an input and an output are set, the code that runs the job.
    • getJavaCode

      public JavaCodeBuilder getJavaCode()
      Returns a Java program, a DataPipelineDesktop class whose main method holds this pipeline's generated code.
    • getName

      public String getName()
    • setName

      public AbstractPipeline setName(String name)
    • getDescription

      public String getDescription()
    • setDescription

      public AbstractPipeline setDescription(String description)
    • getInput

      public PipelineInput getInput()
    • setInput

      public AbstractPipeline setInput(PipelineInput input)
    • setInputAsDataReaderFactory

      public AbstractPipeline setInputAsDataReaderFactory(DataReaderFactory factory)
      Sets the input to a DataReaderPipelineInput that creates its readers with the given factory.
    • setInputAsDataReader

      public AbstractPipeline setInputAsDataReader(DataReader reader)
      Sets the input to a DataReaderPipelineInput that always returns the given reader instance.
    • getSourceEntity

      public EntityDef getSourceEntity()
      Returns the schema entity used to convert and validate records read from the input, before any processing, or null if none.
    • setSourceEntity

      public AbstractPipeline setSourceEntity(EntityDef sourceEntity)
      Sets the schema entity used to convert and validate records read from the input, before any processing (null for none).
    • getTargetEntity

      public EntityDef getTargetEntity()
      Returns the schema entity used to convert and validate records after processing, before they are written, or null if none.
    • setTargetEntity

      public AbstractPipeline setTargetEntity(EntityDef targetEntity)
      Sets the schema entity used to convert and validate records after processing, before they are written (null for none).
    • getOutput

      public PipelineOutput getOutput()
    • setOutput

      public AbstractPipeline setOutput(PipelineOutput output)
    • setOutputAsDataWriterFactory

      public AbstractPipeline setOutputAsDataWriterFactory(DataWriterFactory factory)
      Sets the output to a DataWriterPipelineOutput that creates its writers with the given factory.
    • setOutputAsDataWriter

      public AbstractPipeline setOutputAsDataWriter(DataWriter writer)
      Sets the output to a DataWriterPipelineOutput that always returns the given writer instance.
    • isMultithreaded

      public boolean isMultithreaded()
      Indicates if each processing stage runs in its own thread using an AsyncReader (defaults to true).
    • setMultithreaded

      public AbstractPipeline setMultithreaded(boolean multithreaded)
      Indicates if each processing stage runs in its own thread using an AsyncReader (defaults to true).
    • getDiscardWriter

      public DataWriter getDiscardWriter()
      Returns the writer that receives records failing source or target entity validation, instead of stopping the pipeline, or null if none.
    • setDiscardWriter

      public AbstractPipeline setDiscardWriter(DataWriter discardWriter)
      Sets the writer that receives records failing source or target entity validation, instead of stopping the pipeline.
    • getDiscardReasonFieldName

      public String getDiscardReasonFieldName()
      Returns the name of the field added to discarded records to hold the failure message, or null if none.
    • setDiscardReasonFieldName

      public AbstractPipeline setDiscardReasonFieldName(String discardReasonFieldName)
      Sets the name of the field added to discarded records to hold the failure message; only allowed with a discard writer.
    • getDateTimePatternDetector

      public DateTimePatternDetector getDateTimePatternDetector()
      Returns the detector a Dataset built on this pipeline uses to infer date/time values in text columns.
    • setDateTimePatternDetector

      public AbstractPipeline setDateTimePatternDetector(DateTimePatternDetector dateTimePatternDetector)
      Sets the detector a Dataset built on this pipeline uses to infer date/time values in text columns; null restores the default detector.
    • getNumberDetector

      public NumberDetector getNumberDetector()
      Returns the detector a Dataset built on this pipeline uses to infer numeric values in text columns.
    • setNumberDetector

      public AbstractPipeline setNumberDetector(NumberDetector numberDetector)
      Sets the detector a Dataset built on this pipeline uses to infer numeric values in text columns; null restores the default detector.
    • applySourceEntityTransformer

      protected DataReader applySourceEntityTransformer(DataReader reader)
      Wraps the reader to convert and validate records against the source entity, if one is set, in its own thread when multithreaded.
    • applyTargetEntityTransformer

      protected DataReader applyTargetEntityTransformer(DataReader reader)
      Wraps the reader to convert and validate records against the target entity, if one is set, in its own thread when multithreaded.
    • createDataReader

      public DataReader createDataReader()
      Specified by:
      createDataReader in interface DataReaderFactory
    • applyDataProcessing

      public final DataReader applyDataProcessing(DataReader reader)
      Wraps the reader with the source entity transformer, this pipeline's own processing and the target entity transformer, in that order.
    • applyDataProcessingImpl

      protected abstract DataReader applyDataProcessingImpl(DataReader reader)
      Wraps the reader with this pipeline's own processing; called between the source and target entity transformers.
    • createDataWriter

      public DataWriter createDataWriter()
      Specified by:
      createDataWriter in interface DataWriterFactory
    • createJob

      public Job createJob()
      Creates, but does not run, a job that reads the processed input and writes it to the output.
    • run

      public Job run()
      Runs this pipeline in the current thread and returns the finished job.
    • runAsync

      public Job runAsync()
      Starts this pipeline in a new thread and returns its job without waiting for it to finish.
    • toRecord

      public Record toRecord()
      Description copied from interface: RecordSerializable
      Converts this object's state to a record that RecordSerializable.fromRecord(Record) can load.
      Specified by:
      toRecord in interface RecordSerializable
      Overrides:
      toRecord in class Bean
    • fromRecord

      public AbstractPipeline fromRecord(Record source)
      Description copied from interface: RecordSerializable
      Loads this instance's state from a record and returns this (for fluid API call chaining). For fluid API call chaining, the overridden method should change the declared return type to its class.
      Specified by:
      fromRecord in interface RecordSerializable
      Overrides:
      fromRecord in class Bean
      Parameters:
      source -
      Returns:
      this instance.
    • fromJson

      public AbstractPipeline fromJson(String jsonString)
      Description copied from interface: JsonSerializable
      Loads this object's state from the given JSON text.
      Specified by:
      fromJson in interface JsonSerializable
      Specified by:
      fromJson in interface RecordSerializable
    • toXmlElement

      public Element toXmlElement(Document document)
      Description copied from interface: XmlSerializable
      Returns an element, created with document, describing this object; the default implementation throws a DataException.
      Specified by:
      toXmlElement in interface XmlSerializable
    • fromXmlElement

      public AbstractPipeline fromXmlElement(Element pipelineElement)
      Specified by:
      fromXmlElement in interface XmlSerializable