Class JmsReader
Read
Records from a Java Message Service (JMS) provider.-
Nested Class Summary
Nested classes/interfaces inherited from class com.northconcepts.datapipeline.core.DataEndpoint
DataEndpoint.State -
Field Summary
Fields inherited from class com.northconcepts.datapipeline.core.DataReader
fieldLineage, recordLineageFields inherited from class com.northconcepts.datapipeline.core.DataEndpoint
lastRecord, PRODUCT, PRODUCT_VERSION, VENDOR, XML_INPUT_FACTORY_KEYFields inherited from class com.northconcepts.datapipeline.core.Endpoint
BUFFER_SIZE, captureElapsedTime, DEFAULT_READ_BUFFER_SIZEFields inherited from class com.northconcepts.datapipeline.core.DataObject
id, log, name, TIMESTAMP_FORMAT -
Constructor Summary
ConstructorsConstructorDescriptionJmsReader(JmsSettings settings) Connects to a Java Message Service server.
Apache ActiveMQ Example:
Properties props = new Properties();
props.setProperty(Context.INITIAL_CONTEXT_FACTORY,"org.apache.activemq.jndi.ActiveMQInitialContextFactory");
props.setProperty(Context.PROVIDER_URL,"tcp://localhost:61616");
DataReader reader = new JmsReader(props, JmsDestinationType.QUEUE, "queueName");
-
Method Summary
Modifier and TypeMethodDescriptionaddExceptionProperties(DataException exception) Adds this endpoint's current state to aDataException.protected voidaddMessageHeader(Message message, Record record) Adds the message's JMS headers asjms_*fields and each message property as ajms_property_prefixed field.voidclose()Indicates that this endpoint has finished reading or writing.Returns the JMS connection or null if the reader has not been opened.Returns the timeout interval for receiving messages ornullif reads should not timeout (default tonull).protected RecordonBytesMessage(BytesMessage message) Converts a bytes message to a record whosemessagefield holds the body bytes.protected RecordonMapMessage(MapMessage message) Converts a map message to a record with one field per map entry.protected RecordonMessage(Message message) Converts a received message to a record by message type, adds the JMS header fields and acknowledges it inCLIENTmode; returns null if the message is null.protected RecordonObjectMessage(ObjectMessage message) Converts an object message to a record whosemessagefield holds the object.protected RecordonStreamMessage(StreamMessage message) Converts a stream message to a record whosemessagefield holds only the first value read from the stream.protected RecordonTextMessage(TextMessage message) Converts a text message to a record whosemessagefield holds the text.voidopen()Makes this endpoint ready for reading or writing.protected RecordreadImpl()Overridden by subclasses to read the next record from thisDataReader.setReceiveTimeout(Long timeout) Sets the timeout interval for receiving messages ornullif reads should not timeout (default tonull).Methods inherited from class com.northconcepts.datapipeline.core.DataReader
addLineage, available, getBufferSize, getNestedEndpoint, getNestedReader, getReader, getRootEndpoint, getRootReader, isExhausted, isLineageSupported, isSaveLineage, peek, pop, push, read, setSaveLineage, skipMethods inherited from class com.northconcepts.datapipeline.core.DataEndpoint
decrementRecordCount, enableJmx, getLastRecord, getRecordCount, getRecordCountAsBigInteger, getRecordCountAsString, incrementRecordCount, isRecordCountBigInteger, resetRecordCount, toStringMethods inherited from class com.northconcepts.datapipeline.core.Endpoint
addElapsedtime, assertClosed, assertNotOpened, assertOpened, finalize, getClosedOn, getDescription, getElapsedTime, getElapsedTimeAsString, getOpenedOn, getOpenElapsedTime, getOpenElapsedTimeAsString, getSelfTime, getSelfTimeAsString, getState, isCaptureElapsedTime, isClosed, isOpen, setCaptureElapsedTime, setDescription
-
Constructor Details
-
JmsReader
Connects to a Java Message Service server.
Apache ActiveMQ Example:
Properties props = new Properties();
props.setProperty(Context.INITIAL_CONTEXT_FACTORY,"org.apache.activemq.jndi.ActiveMQInitialContextFactory");
props.setProperty(Context.PROVIDER_URL,"tcp://localhost:61616");
DataReader reader = new JmsReader(props, JmsDestinationType.QUEUE, "queueName");
Read your JMS provider manual for more details
- Parameters:
properties- the properties required by the JMS servertype- the type of messaging useddestinationName- the name of the Topic or Queue- See Also:
-
-
Method Details
-
getSettings
-
getConnection
Returns the JMS connection or null if the reader has not been opened. -
getReceiveTimeout
Returns the timeout interval for receiving messages ornullif reads should not timeout (default tonull).- Returns:
- the timeout interval
-
setReceiveTimeout
Sets the timeout interval for receiving messages ornullif reads should not timeout (default tonull).- Parameters:
timeout- timeout interval for receiving messages- Returns:
- the
JmsReaderinstance
-
open
Description copied from class:DataEndpointMakes this endpoint ready for reading or writing.- Overrides:
openin classIntegrationReader- Throws:
DataException
-
close
Description copied from class:DataEndpointIndicates that this endpoint has finished reading or writing.- Overrides:
closein classDataEndpoint- Throws:
DataException
-
readImpl
Description copied from class:DataReaderOverridden by subclasses to read the next record from thisDataReader. The default implementation ofDataReader.read()now insures that this method will not be called again after it returns anull.If no record is available,
nullwill be returned.Contract for subclasses (see also
docs/authoring/DataReader.md):- Return
nullexactly once to signal end-of-stream. - Do not call this method directly —
DataReader.read()is the template method that wraps exceptions, increments the record count, and tracks lineage. - Do not call
DataEndpoint.incrementRecordCount()here;DataReader.read()already does. - Throw raw exceptions;
DataReader.read()wraps them viaexception(throwable).
- Specified by:
readImplin classDataReader- Throws:
Throwable
- Return
-
onMessage
Converts a received message to a record by message type, adds the JMS header fields and acknowledges it inCLIENTmode; returns null if the message is null.- Throws:
Throwable
-
addMessageHeader
Adds the message's JMS headers asjms_*fields and each message property as ajms_property_prefixed field.- Throws:
Throwable
-
onMapMessage
Converts a map message to a record with one field per map entry.- Throws:
JMSException
-
onBytesMessage
Converts a bytes message to a record whosemessagefield holds the body bytes.- Throws:
JMSException
-
onObjectMessage
Converts an object message to a record whosemessagefield holds the object.- Throws:
JMSException
-
onStreamMessage
Converts a stream message to a record whosemessagefield holds only the first value read from the stream.- Throws:
JMSException
-
onTextMessage
Converts a text message to a record whosemessagefield holds the text.- Throws:
JMSException
-
addExceptionProperties
Description copied from class:EndpointAdds this endpoint's current state to aDataException. Since this method is called whenever an exception is thrown, subclasses should override it to add their specific information.- Overrides:
addExceptionPropertiesin classDataReader
-