Class DeMux
java.lang.Object
com.northconcepts.datapipeline.multiplex.DeMux
- All Implemented Interfaces:
Runnable
Converts a single source DataReader into many using one of the provided strategies.
-
Nested Class Summary
Nested ClassesModifier and TypeClassDescriptionstatic interfaceStrategy for distributing the source records among the readers created bycreateReader().static enumThe built-in strategies: send every record to all readers or each record to one reader in turn. -
Field Summary
FieldsModifier and TypeFieldDescriptionprotected static final Recordstatic final Loggerstatic final int -
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionCreates 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> getSink()Returns the readers created bycreateReader(), which receive the source's records.Returns the thread running the feeder ornullif it is not running.booleanvoidjoin()Waits for the feeder process to finish reading records from the source reader and sending them to downstream readers.voidjoin(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.voidjoin(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 voidDrops closed readers fromgetSink(); the round-robin strategy calls it when it reaches a closed reader.voidrun()Runs the feeder process in the current thread.voidrun(boolean async) Runs the feeder process in the current thread (false) or in a new thread (true).voidrunAsync()Runs the feeder process in a new thread.voidstop()Makes the feeder stop before reading the next source record; the feeder thread is not interrupted.
-
Field Details
-
log
public static final Logger log -
MAX_QUEUE_SIZE
public static final int MAX_QUEUE_SIZE- See Also:
-
EOF
-
-
Constructor Details
-
DeMux
-
-
Method Details
-
getStrategy
-
getSource
-
getSink
Returns the readers created bycreateReader(), which receive the source's records. -
getThread
Returns the thread running the feeder ornullif it is not running. -
createReader
Creates a reader that receives source records according to the strategy; create all readers before running the feeder so none miss records. -
run
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
Runs the feeder process in a new thread. This method does not block.- Throws:
DataException
-
run
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:
runin interfaceRunnable- Throws:
DataException
-
join
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
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
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 fromgetSink(); the round-robin strategy calls it when it reaches a closed reader.
-