Package com.northconcepts.datapipeline.core
package com.northconcepts.datapipeline.core
Core DataPipeline classes for reading, writing, and processing data.
-
ClassDescriptionAbstract super-class, with some common logic, for reading records.Abstract super-class, with some common logic, for writing records.ArrayValue holds an ordered collection (also known as a sequence) of persistent data.AsyncMultiReader reads from one or more source DataReaders asynchronously using a separate thread for each one.A proxy that reads data asynchronously using a separate thread.A proxy that uses multiple threads to process incoming data, where the work applied to the incoming data is specified by the a
DataReaderDecorator.A proxy that writes data asynchronously using a separate thread.Abstract super-class for writing records to a binary stream.A fixed-size list of values used as a composite key (for example in lookups, grouping and duplicate removal); a single-value instance hashes like, and equals, that value itself.ARecordListthat guards many of its methods with a read-write lock; others, such assize()andclear(), are unguarded, as are the iterators, streams and lists it hands out.Abstract super-class for reading and writing records.Lifecycle states of an endpoint:NEWuntil opened, thenOPENED, thenCLOSED.Helper class--makes it easy to work with multiple readers and writers as a single unit.Indicates an error has occurred during the course of execution.An interface for classes capable of adding context/information to exceptions usually to aid in identifying and troubleshooting issues.Abstract super-class for classes that work with data and require a consistent error handling mechanism.Abstract super-class for reading records.Abstract super-class for writing records.A proxy that prints records passing through to a stream in a human-readable formatA proxy that prints records passing through to a stream in a human-readable formatBase class for anything that is opened and closed, such as readers, writers and file systems; tracks its state, open and close times, and elapsed time.Field holds persistent key-value data as part of a record.SeeFieldPathfor supported field name expressions.Declares the Java type of named fields for expressions such asFilterExpressionandSetCalculatedField; declared types take precedence over those inferred from each record.SeeFieldPathfor supported field name expressions.Represents the location of a field within a data pipeline context.A functional interface for predicates that testFieldLocationinstances.An abstract representation for the location of a field within a record.FieldType lists of all field data types used in records.Functions is the entry point for adding new method aliases to the dynamic expression language.Receives callbacks whileNodeVisitor.visit(Node, INodeVisitor)walks a tree of records, fields, arrays and values depth-first.Character-level parser with lookahead, used by the text readers to peek at, match and consume their input.The interface implemented by classes to participate in serialization to and from JSON.A proxy that can skip a number of upstream records, limit the number of records sent downstream, or both.Abstract super-class for writing records to a text stream.AnIParserthat reads its input one line at a time; lookahead and matching only see the current line untilLineParser.cacheNextLine()moves to the next one.A diagnostic message with a severity level, an optional exception or stack trace, and the time and thread it was created on; collected byMessages.The severity of aMessage.Thread-local container for messages and exceptions originating in asynchronous operations.Writes records to multipleDataWriter.Writes each record to a single writer by choosing the one with the highest available capacity (DataWriter.available()If all writers have identical capacities, this strategy behaves likeMultiWriter.RoundRobinWriteStrategyand writes to each writer in turn.Writes a clone (to prevent side effects) of each record to all writers.Writes each record to all writers.Divides records evenly between all writers by cycling through the list and writing each record to a single writer in turn.Strategy for distributing each record among aMultiWriter's target writers.Node is the base class for all persistent data in Data Pipeline.DuplicateNodeAction lists all actions that can be taken automatically when a target field already exists during a record copy or lookup.NodeType lists all concrete data node types used in Data Pipeline.BaseINodeVisitorwhose callbacks do nothing; override only the ones you need and walk a tree withNodeVisitor.visit(Node, INodeVisitor).A data source that produces no records.Discards records.Abstract super-class for obtaining records from a text stream using a parser.Reads the records written to a connectedPipedWriter, usually on another thread; reads block until a record arrives and the stream ends when the writer is closed.Writes records to a connectedPipedReader, blocking while its queue is full; closing this writer ends the reader's stream.Abstract super-class for obtaining records from anotherDataReader, possibly transforming them along the way.Abstract super-class for writing records to anotherDataWriter, possibly transforming them along the way.Record holds persistent data in key-value fields as it flows through the pipeline.Lists the lifecycle states of a record as it flows through the pipeline.Orders records by a list of fields, each ascending or descending with an optionalCollator; when no fields are added, records are compared on all their fields by name.A general interface for objects that are associated with aRecord.An in-memory list of records that can be sorted, searched with aFilter, and converted to binary, arrays or a reader.The interface implemented by classes to participate in serialization to and from Record.A proxy that removes duplicate records.Combines one or more DataReaders into a single stream by reading from each until empty then moving to the next.Writes to sequence of data writers created by a factory in turn and rolled based on a ISequenceStrategy policy.Creates the nested writer for each new sequence of aSequenceWriter.Ends a sequence once it has been running for a given time; checked only when the next record arrives.Decides when aSequenceWriterends the current sequence and starts a new one.Ends the current sequence when any of the given strategies does.Ends each sequence once it has received the given number of records.Ends the current sequence at each time produced by aScheduler; checked only when the next record arrives.One segment of aSequenceWriter's output, with its 0-based index, start time and record counts.Session is the base type for objects that can contain temporary, non-persistent data in DataPipeline.Stores session properties in a lazily created map, keyed by property name, class name, or both.SingleValue is an immutableValueNodethat holds a single scalar value (ornull).A proxy that sorts records.Writes records to a stream in a human-readable format.How aStreamWriterprints each record (after its 0-based number): plain text, JSON, XML, or text plus session properties.AnIParserover an in-memory string.A proxy reader that also writes every record passing through it to a DataWriter.Abstract super-class for obtaining records from a text stream.Abstract super-class for writing records to a text stream.Abstract super-class for writing records to a text stream.Reads an underlying DataReader for a maximum period of time or until the source DataReader is finished.ValueNode<T>ValueNode is the base class for all persistent values held in fields and arrays.A comparator forValueNodeinstances that supports sorting and comparing records, arrays, and single values.A general interface for nodes that allow values to be added to itself.The interface implemented by classes to participate in serialization to and from XML.