java.lang.Object
com.northconcepts.datapipeline.multiplex.DeMux
All Implemented Interfaces:
Runnable

public class DeMux extends Object implements Runnable
Converts a single source DataReader into many using one of the provided strategies.
  • Nested Class Summary

    Nested Classes
    Modifier and Type
    Class
    Description
    static interface 
    Strategy for distributing the source records among the readers created by createReader().
    static enum 
    The built-in strategies: send every record to all readers or each record to one reader in turn.
  • Field Summary

    Fields
    Modifier and Type
    Field
    Description
    protected static final Record
     
    static final Logger
     
    static final int
     
  • Constructor Summary

    Constructors
    Constructor
    Description
    DeMux(DataReader source, DeMux.IStrategy strategy)
     
  • Method Summary

    Modifier and Type
    Method
    Description
    Creates a reader that receives source records according to the strategy; create all readers before running the feeder so none miss records.
    List<com.northconcepts.datapipeline.multiplex.DeMuxReader>
    Returns the readers created by createReader(), which receive the source's records.
     
     
    Returns the thread running the feeder or null if it is not running.
    boolean
     
    void
    Waits for the feeder process to finish reading records from the source reader and sending them to downstream readers.
    void
    join(long millis)
    Waits at most millis milliseconds for the feeder process to finish reading records from the source reader and sending them to downstream readers.
    void
    join(long millis, int nanos)
    Waits at most millis milliseconds plus nanos nanoseconds for the feeder process to finish reading records from the source reader and sending them to downstream readers.
    protected void
    Drops closed readers from getSink(); the round-robin strategy calls it when it reaches a closed reader.
    void
    run()
    Runs the feeder process in the current thread.
    void
    run(boolean async)
    Runs the feeder process in the current thread (false) or in a new thread (true).
    void
    Runs the feeder process in a new thread.
    void
    Makes the feeder stop before reading the next source record; the feeder thread is not interrupted.

    Methods inherited from class java.lang.Object

    clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
  • Field Details

    • log

      public static final Logger log
    • MAX_QUEUE_SIZE

      public static final int MAX_QUEUE_SIZE
      See Also:
    • EOF

      protected static final Record EOF
  • Constructor Details

  • Method Details

    • getStrategy

      public DeMux.IStrategy getStrategy()
    • getSource

      public DataReader getSource()
    • getSink

      public List<com.northconcepts.datapipeline.multiplex.DeMuxReader> getSink()
      Returns the readers created by createReader(), which receive the source's records.
    • getThread

      public Thread getThread()
      Returns the thread running the feeder or null if it is not running.
    • createReader

      public DataReader createReader()
      Creates a reader that receives source records according to the strategy; create all readers before running the feeder so none miss records.
    • run

      public void run(boolean async) throws DataException
      Runs the feeder process in the current thread (false) or in a new thread (true). If true, this method blocks until all data is read in from the source reader and sent to the downstream readers.
      Throws:
      DataException
    • runAsync

      public void runAsync() throws DataException
      Runs the feeder process in a new thread. This method does not block.
      Throws:
      DataException
    • run

      public void run() throws DataException
      Runs the feeder process in the current thread. This method blocks until all data is read in from the source reader and sent to the downstream readers.
      Specified by:
      run in interface Runnable
      Throws:
      DataException
    • join

      public void join(long millis) throws InterruptedException
      Waits at most millis milliseconds for the feeder process to finish reading records from the source reader and sending them to downstream readers. A timeout of 0 means to wait forever.
      Throws:
      InterruptedException
    • join

      public void join(long millis, int nanos) throws InterruptedException
      Waits at most millis milliseconds plus nanos nanoseconds for the feeder process to finish reading records from the source reader and sending them to downstream readers. A timeout of 0 means to wait forever.
      Throws:
      InterruptedException
    • join

      public void join() throws InterruptedException
      Waits for the feeder process to finish reading records from the source reader and sending them to downstream readers. A timeout of 0 means to wait forever.
      Throws:
      InterruptedException
    • isFinished

      public boolean isFinished()
    • stop

      public void stop()
      Makes the feeder stop before reading the next source record; the feeder thread is not interrupted.
    • removeClosedSinks

      protected void removeClosedSinks()
      Drops closed readers from getSink(); the round-robin strategy calls it when it reaches a closed reader.