public class RMQConnection
extends java.lang.Object
implements javax.jms.Connection, javax.jms.QueueConnection, javax.jms.TopicConnection
Connection, QueueConnection and TopicConnection interfaces.
A RMQConnection object holds a list of RMQSession objects as well as the actual
{link com.rabbitmq.client.Connection} object that represents the TCP connection to the RabbitMQ broker.
This implementation also holds a reference to the executor service that is used by the connection so that we can pause incoming messages.
| Modifier and Type | Field and Description |
|---|---|
static int |
NO_CHANNEL_QOS |
| Constructor and Description |
|---|
RMQConnection(com.rabbitmq.client.Connection rabbitConnection)
Creates an RMQConnection object, with default termination timeout of 15 seconds, a 2 seconds timeout for onMessage, and unlimited reads from QueueBrowsers.
|
RMQConnection(com.rabbitmq.client.Connection rabbitConnection,
long terminationTimeout,
int queueBrowserReadMax,
int onMessageTimeoutMs)
Creates an RMQConnection object.
|
RMQConnection(ConnectionParams connectionParams)
Creates an RMQConnection object.
|
| Modifier and Type | Method and Description |
|---|---|
void |
close()
From the JMS Spec:
|
javax.jms.ConnectionConsumer |
createConnectionConsumer(javax.jms.Destination destination,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.ConnectionConsumer |
createConnectionConsumer(javax.jms.Queue queue,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.ConnectionConsumer |
createConnectionConsumer(javax.jms.Topic topic,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.ConnectionConsumer |
createDurableConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.QueueSession |
createQueueSession(boolean transacted,
int acknowledgeMode) |
javax.jms.Session |
createSession() |
javax.jms.Session |
createSession(boolean transacted,
int acknowledgeMode) |
javax.jms.Session |
createSession(int sessionMode) |
javax.jms.ConnectionConsumer |
createSharedConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.ConnectionConsumer |
createSharedDurableConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages) |
javax.jms.TopicSession |
createTopicSession(boolean transacted,
int acknowledgeMode) |
java.lang.String |
getClientID() |
javax.jms.ExceptionListener |
getExceptionListener() |
javax.jms.ConnectionMetaData |
getMetaData() |
ReplyToStrategy |
getReplyToStrategy()
Gets the reply to strategy that should be followed if as reply to is
found on a received message.
|
java.util.List<java.lang.String> |
getTrustedPackages() |
boolean |
isStopped() |
void |
setClientID(java.lang.String clientID) |
void |
setExceptionListener(javax.jms.ExceptionListener listener) |
void |
start() |
void |
stop() |
java.lang.String |
toString() |
public static final int NO_CHANNEL_QOS
public RMQConnection(ConnectionParams connectionParams)
connectionParams - parameters for this connectionpublic RMQConnection(com.rabbitmq.client.Connection rabbitConnection,
long terminationTimeout,
int queueBrowserReadMax,
int onMessageTimeoutMs)
rabbitConnection - the TCP connection wrapper to the RabbitMQ brokerterminationTimeout - timeout for close in millisecondsqueueBrowserReadMax - maximum number of messages to read from a QueueBrowser (before filtering)onMessageTimeoutMs - how long to wait for onMessage to return, in millisecondspublic RMQConnection(com.rabbitmq.client.Connection rabbitConnection)
rabbitConnection - the TCP connection wrapper to the RabbitMQ brokerpublic javax.jms.Session createSession(boolean transacted,
int acknowledgeMode)
throws javax.jms.JMSException
createSession in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic java.lang.String getClientID()
throws javax.jms.JMSException
getClientID in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic void setClientID(java.lang.String clientID)
throws javax.jms.JMSException
setClientID in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic java.util.List<java.lang.String> getTrustedPackages()
public javax.jms.ConnectionMetaData getMetaData()
throws javax.jms.JMSException
getMetaData in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic javax.jms.ExceptionListener getExceptionListener()
throws javax.jms.JMSException
getExceptionListener in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic void setExceptionListener(javax.jms.ExceptionListener listener)
throws javax.jms.JMSException
setExceptionListener in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic void start()
throws javax.jms.JMSException
start in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic void stop()
throws javax.jms.JMSException
stop in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic boolean isStopped()
true if this connection is in a stopped statepublic void close()
throws javax.jms.JMSException
This call blocks until a receive or message listener in progress has completed. A blocked message consumer receive call returns null when this message consumer is closed.
close in interface java.lang.AutoCloseableclose in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic javax.jms.TopicSession createTopicSession(boolean transacted,
int acknowledgeMode)
throws javax.jms.JMSException
createTopicSession in interface javax.jms.TopicConnectionjavax.jms.JMSExceptionpublic javax.jms.ConnectionConsumer createConnectionConsumer(javax.jms.Topic topic,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
throws javax.jms.JMSException
createConnectionConsumer in interface javax.jms.TopicConnectionjava.lang.UnsupportedOperationException - - optional method not implementedjavax.jms.JMSExceptionpublic javax.jms.QueueSession createQueueSession(boolean transacted,
int acknowledgeMode)
throws javax.jms.JMSException
createQueueSession in interface javax.jms.QueueConnectionjavax.jms.JMSExceptionpublic javax.jms.ConnectionConsumer createConnectionConsumer(javax.jms.Queue queue,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
throws javax.jms.JMSException
createConnectionConsumer in interface javax.jms.QueueConnectionjava.lang.UnsupportedOperationException - - optional method not implementedjavax.jms.JMSExceptionpublic javax.jms.ConnectionConsumer createConnectionConsumer(javax.jms.Destination destination,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
createConnectionConsumer in interface javax.jms.Connectionjava.lang.UnsupportedOperationException - - optional method not implementedpublic javax.jms.ConnectionConsumer createDurableConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
createDurableConnectionConsumer in interface javax.jms.ConnectioncreateDurableConnectionConsumer in interface javax.jms.TopicConnectionjava.lang.UnsupportedOperationException - - optional method not implementedpublic java.lang.String toString()
toString in class java.lang.Objectpublic javax.jms.Session createSession(int sessionMode)
throws javax.jms.JMSException
createSession in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic javax.jms.Session createSession()
throws javax.jms.JMSException
createSession in interface javax.jms.Connectionjavax.jms.JMSExceptionpublic javax.jms.ConnectionConsumer createSharedConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
createSharedConnectionConsumer in interface javax.jms.Connectionjava.lang.UnsupportedOperationException - - optional method not implementedpublic javax.jms.ConnectionConsumer createSharedDurableConnectionConsumer(javax.jms.Topic topic,
java.lang.String subscriptionName,
java.lang.String messageSelector,
javax.jms.ServerSessionPool sessionPool,
int maxMessages)
createSharedDurableConnectionConsumer in interface javax.jms.Connectionjava.lang.UnsupportedOperationException - - optional method not implementedpublic ReplyToStrategy getReplyToStrategy()
Copyright © 2023. All rights reserved.