public abstract class RMQMessage
extends java.lang.Object
implements javax.jms.Message, java.lang.Cloneable
| Modifier and Type | Field and Description |
|---|---|
protected static int |
DEFAULT_MESSAGE_BODY_SIZE
When we create a message that has a byte[] as the underlying
structure, BytesMessage and StreamMessage, this is the default size.
|
protected org.slf4j.Logger |
logger
Logger shared with derived classes
|
protected static java.lang.String |
MSG_EOF
Error message when we get an EOF exception
|
protected static java.lang.String |
NOT_READABLE
Error message when the message is not readable
|
protected static java.lang.String |
NOT_WRITEABLE
Error message when the message is not writeable
|
protected static java.lang.String |
UNABLE_TO_CAST
Error message used when throwing
MessageFormatException |
| Constructor and Description |
|---|
RMQMessage()
Constructor for auto de-serialization
|
| Modifier and Type | Method and Description |
|---|---|
void |
acknowledge()
Although the JavaDoc implies otherwise, this call acknowledges all unacknowledged messages on the message's
session.
|
void |
clearBody() |
protected abstract void |
clearBodyInternal() |
void |
clearProperties() |
java.lang.Object |
clone() |
protected static void |
copyAttributes(RMQMessage rmqMessage,
javax.jms.Message message)
Assign generic attributes.
|
protected abstract <T> T |
doGetBody(java.lang.Class<T> c) |
static void |
doNotDeclareReplyToDestination(javax.jms.Message message)
Indicates to not declare a reply-to
RMQDestination. |
boolean |
equals(java.lang.Object obj) |
<T> T |
getBody(java.lang.Class<T> c) |
boolean |
getBooleanProperty(java.lang.String name) |
byte |
getByteProperty(java.lang.String name) |
double |
getDoubleProperty(java.lang.String name) |
float |
getFloatProperty(java.lang.String name) |
java.lang.String |
getInternalID()
Returns the unique ID for this message.
|
int |
getIntProperty(java.lang.String name) |
java.lang.String |
getJMSCorrelationID() |
byte[] |
getJMSCorrelationIDAsBytes() |
int |
getJMSDeliveryMode() |
long |
getJMSDeliveryTime() |
javax.jms.Destination |
getJMSDestination() |
long |
getJMSExpiration() |
java.lang.String |
getJMSMessageID() |
int |
getJMSPriority() |
boolean |
getJMSRedelivered() |
javax.jms.Destination |
getJMSReplyTo() |
long |
getJMSTimestamp() |
java.lang.String |
getJMSType() |
long |
getLongProperty(java.lang.String name) |
java.lang.Object |
getObjectProperty(java.lang.String name) |
java.util.Enumeration<?> |
getPropertyNames() |
long |
getRabbitDeliveryTag()
returns the delivery tag for this message
|
RMQSession |
getSession()
Returns the session this object belongs to
|
short |
getShortProperty(java.lang.String name) |
java.lang.String |
getStringProperty(java.lang.String name) |
int |
hashCode() |
protected boolean |
isReadonlyBody()
Returns true if this message body is read only
This means that message has been received and can not
be modified
|
protected boolean |
isReadOnlyProperties() |
protected void |
loggerDebugByteArray(java.lang.String format,
byte[] buffer,
java.lang.Object arg) |
boolean |
propertyExists(java.lang.String name) |
protected abstract void |
readAmqpBody(byte[] barr)
Invoked when an AMQP message is being transformed into a RMQMessage
The implementing class should only read its body by this method
|
protected abstract void |
readBody(java.io.ObjectInput inputStream,
java.io.ByteArrayInputStream bin)
Invoked when a message is being deserialized to read and decode the message body.
|
protected static java.lang.Object |
readPrimitive(java.io.ObjectInput in)
Utility method to read objects from a stream.
|
void |
setBooleanProperty(java.lang.String name,
boolean value) |
void |
setByteProperty(java.lang.String name,
byte value) |
void |
setDoubleProperty(java.lang.String name,
double value) |
void |
setFloatProperty(java.lang.String name,
float value) |
void |
setIntProperty(java.lang.String name,
int value) |
void |
setJMSCorrelationID(java.lang.String correlationID) |
void |
setJMSCorrelationIDAsBytes(byte[] correlationID) |
void |
setJMSDeliveryMode(int deliveryMode) |
void |
setJMSDeliveryTime(long deliveryTime) |
void |
setJMSDestination(javax.jms.Destination destination) |
void |
setJMSExpiration(long expiration) |
void |
setJMSMessageID(java.lang.String id) |
void |
setJMSPriority(int priority) |
void |
setJMSRedelivered(boolean redelivered) |
void |
setJMSReplyTo(javax.jms.Destination replyTo) |
void |
setJMSTimestamp(long timestamp) |
void |
setJMSType(java.lang.String type) |
void |
setLongProperty(java.lang.String name,
long value) |
void |
setObjectProperty(java.lang.String name,
java.lang.Object value) |
protected void |
setRabbitDeliveryTag(long rabbitDeliveryTag)
Sets the RabbitMQ delivery tag for this message.
|
protected void |
setReadonly(boolean readonly)
Sets the read only flag on this message
|
protected void |
setReadOnlyBody(boolean readonly) |
protected void |
setReadOnlyProperties(boolean readonly) |
protected void |
setSession(RMQSession session)
Sets the session this object was received by
|
void |
setShortProperty(java.lang.String name,
short value) |
void |
setStringProperty(java.lang.String name,
java.lang.String value) |
protected abstract void |
writeAmqpBody(java.io.ByteArrayOutputStream out)
Invoked when
toAmqpByteArray() is called to create
a byte[] from a message. |
protected abstract void |
writeBody(java.io.ObjectOutput out,
java.io.ByteArrayOutputStream bout)
Invoked when
toByteArray() is called to create
a byte[] from a message. |
protected static void |
writePrimitive(java.lang.Object s,
java.io.ObjectOutput out)
Utility method used to be able to write primitives and objects to a data
stream without keeping track of order and type.
|
protected static void |
writePrimitive(java.lang.Object s,
java.io.ObjectOutput out,
boolean allowSerializable) |
protected final org.slf4j.Logger logger
protected static final java.lang.String NOT_READABLE
protected static final java.lang.String NOT_WRITEABLE
protected static final java.lang.String UNABLE_TO_CAST
MessageFormatExceptionprotected static final java.lang.String MSG_EOF
protected static final int DEFAULT_MESSAGE_BODY_SIZE
protected void loggerDebugByteArray(java.lang.String format,
byte[] buffer,
java.lang.Object arg)
protected boolean isReadonlyBody()
protected boolean isReadOnlyProperties()
protected void setReadonly(boolean readonly)
readonly - read only flag valueconvertMessage(RMQSession, RMQDestination, com.rabbitmq.client.GetResponse, ReceivingContextConsumer)protected void setReadOnlyBody(boolean readonly)
protected void setReadOnlyProperties(boolean readonly)
public long getRabbitDeliveryTag()
protected void setRabbitDeliveryTag(long rabbitDeliveryTag)
rabbitDeliveryTag - RabbitMQ delivery tagconvertMessage(RMQSession, RMQDestination, com.rabbitmq.client.GetResponse, ReceivingContextConsumer)public RMQSession getSession()
protected void setSession(RMQSession session)
session - the session this object was received byconvertMessage(RMQSession, RMQDestination, com.rabbitmq.client.GetResponse, ReceivingContextConsumer),
RMQSession.acknowledgeMessage(RMQMessage)public java.lang.String getJMSMessageID()
throws javax.jms.JMSException
getJMSMessageID in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSMessageID(java.lang.String id)
throws javax.jms.JMSException
setJMSMessageID in interface javax.jms.Messagejavax.jms.JMSExceptionpublic long getJMSTimestamp()
throws javax.jms.JMSException
getJMSTimestamp in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSTimestamp(long timestamp)
throws javax.jms.JMSException
setJMSTimestamp in interface javax.jms.Messagejavax.jms.JMSExceptionpublic byte[] getJMSCorrelationIDAsBytes()
throws javax.jms.JMSException
getJMSCorrelationIDAsBytes in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSCorrelationIDAsBytes(byte[] correlationID)
throws javax.jms.JMSException
setJMSCorrelationIDAsBytes in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSCorrelationID(java.lang.String correlationID)
throws javax.jms.JMSException
setJMSCorrelationID in interface javax.jms.Messagejavax.jms.JMSExceptionpublic java.lang.String getJMSCorrelationID()
throws javax.jms.JMSException
getJMSCorrelationID in interface javax.jms.Messagejavax.jms.JMSExceptionpublic javax.jms.Destination getJMSReplyTo()
throws javax.jms.JMSException
getJMSReplyTo in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSReplyTo(javax.jms.Destination replyTo)
throws javax.jms.JMSException
setJMSReplyTo in interface javax.jms.Messagejavax.jms.JMSExceptionpublic javax.jms.Destination getJMSDestination()
throws javax.jms.JMSException
getJMSDestination in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSDestination(javax.jms.Destination destination)
throws javax.jms.JMSException
setJMSDestination in interface javax.jms.Messagejavax.jms.JMSExceptionpublic int getJMSDeliveryMode()
throws javax.jms.JMSException
getJMSDeliveryMode in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSDeliveryMode(int deliveryMode)
throws javax.jms.JMSException
setJMSDeliveryMode in interface javax.jms.Messagejavax.jms.JMSExceptionpublic boolean getJMSRedelivered()
throws javax.jms.JMSException
getJMSRedelivered in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSRedelivered(boolean redelivered)
throws javax.jms.JMSException
setJMSRedelivered in interface javax.jms.Messagejavax.jms.JMSExceptionpublic java.lang.String getJMSType()
throws javax.jms.JMSException
getJMSType in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSType(java.lang.String type)
throws javax.jms.JMSException
setJMSType in interface javax.jms.Messagejavax.jms.JMSExceptionpublic long getJMSExpiration()
throws javax.jms.JMSException
getJMSExpiration in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSExpiration(long expiration)
throws javax.jms.JMSException
setJMSExpiration in interface javax.jms.Messagejavax.jms.JMSExceptionpublic int getJMSPriority()
throws javax.jms.JMSException
getJMSPriority in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSPriority(int priority)
throws javax.jms.JMSException
setJMSPriority in interface javax.jms.Messagejavax.jms.JMSExceptionpublic final void clearProperties()
throws javax.jms.JMSException
clearProperties in interface javax.jms.Messagejavax.jms.JMSExceptionpublic boolean propertyExists(java.lang.String name)
propertyExists in interface javax.jms.Messagepublic boolean getBooleanProperty(java.lang.String name)
throws javax.jms.JMSException
getBooleanProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic byte getByteProperty(java.lang.String name)
throws javax.jms.JMSException
getByteProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic short getShortProperty(java.lang.String name)
throws javax.jms.JMSException
getShortProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic int getIntProperty(java.lang.String name)
throws javax.jms.JMSException
getIntProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic long getLongProperty(java.lang.String name)
throws javax.jms.JMSException
getLongProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic float getFloatProperty(java.lang.String name)
throws javax.jms.JMSException
getFloatProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic double getDoubleProperty(java.lang.String name)
throws javax.jms.JMSException
getDoubleProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic java.lang.String getStringProperty(java.lang.String name)
throws javax.jms.JMSException
getStringProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic java.lang.Object getObjectProperty(java.lang.String name)
throws javax.jms.JMSException
getObjectProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic java.util.Enumeration<?> getPropertyNames()
throws javax.jms.JMSException
getPropertyNames in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setBooleanProperty(java.lang.String name,
boolean value)
throws javax.jms.JMSException
setBooleanProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setByteProperty(java.lang.String name,
byte value)
throws javax.jms.JMSException
setByteProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setShortProperty(java.lang.String name,
short value)
throws javax.jms.JMSException
setShortProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setIntProperty(java.lang.String name,
int value)
throws javax.jms.JMSException
setIntProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setLongProperty(java.lang.String name,
long value)
throws javax.jms.JMSException
setLongProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setFloatProperty(java.lang.String name,
float value)
throws javax.jms.JMSException
setFloatProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setDoubleProperty(java.lang.String name,
double value)
throws javax.jms.JMSException
setDoubleProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setStringProperty(java.lang.String name,
java.lang.String value)
throws javax.jms.JMSException
setStringProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setObjectProperty(java.lang.String name,
java.lang.Object value)
throws javax.jms.JMSException
setObjectProperty in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void acknowledge()
throws javax.jms.JMSException
The JavaDoc says:
A client may individually acknowledge each message as it is consumed, or it may choose to acknowledge messages as an application-defined group (which is done by calling acknowledge on the last received message of the group, thereby acknowledging all messages consumed by the session.)
...but there is no way to do either of these things, as is explained by the specification[1].
[1] JMS 1.1 spec Section 11.3.2.2.
acknowledge in interface javax.jms.Messagejavax.jms.JMSExceptionRMQSession.acknowledgeMessage(RMQMessage)public final void clearBody()
throws javax.jms.JMSException
clearBody in interface javax.jms.Messagejavax.jms.JMSExceptionprotected abstract void clearBodyInternal()
throws javax.jms.JMSException
javax.jms.JMSExceptionprotected abstract void writeBody(java.io.ObjectOutput out,
java.io.ByteArrayOutputStream bout)
throws java.io.IOException
toByteArray() is called to create
a byte[] from a message. Each subclass must implement this, but ONLY
write its specific body. All the properties defined in Message
will be written by the parent class.out - - the output stream to which the structured part of message body (scalar types) is writtenbout - - the output stream to which the un-structured part of message body (explicit bytes) is writtenjava.io.IOException - if the body can not be writtenprotected abstract void writeAmqpBody(java.io.ByteArrayOutputStream out)
throws java.io.IOException
toAmqpByteArray() is called to create
a byte[] from a message. Each subclass must implement this, but ONLY
write its specific body.out - - the output stream to which the message body is writtenjava.io.IOException - if the body can not be writtenprotected abstract void readBody(java.io.ObjectInput inputStream,
java.io.ByteArrayInputStream bin)
throws java.io.IOException,
java.lang.ClassNotFoundException
inputStream - - the stream to read its body frombin - - the underlying byte input streamjava.io.IOException - if a read error occurs on the input streamjava.lang.ClassNotFoundException - if the object class cannot be foundprotected abstract void readAmqpBody(byte[] barr)
barr - - the byte array payload of the AMQP messagepublic static void doNotDeclareReplyToDestination(javax.jms.Message message)
throws javax.jms.JMSException
RMQDestination.
This is used for inbound messages, when the reply-to destination is supposed to be already created. This avoids trying to create a second time a temporary queue.
message - javax.jms.JMSExceptionpublic int hashCode()
hashCode in class java.lang.Objectpublic boolean equals(java.lang.Object obj)
equals in class java.lang.Objectpublic java.lang.String getInternalID()
protected static void writePrimitive(java.lang.Object s,
java.io.ObjectOutput out)
throws java.io.IOException,
javax.jms.MessageFormatException
This also allows any Object to be written.
The purpose of this method is to optimise the writing of a primitive that is represented as an object by only writing the type and the primitive value to the stream.
s - the primitive to be writtenout - the stream to write the primitive to.java.io.IOException - if an I/O error occursjavax.jms.MessageFormatException - if message cannot be parsedprotected static void writePrimitive(java.lang.Object s,
java.io.ObjectOutput out,
boolean allowSerializable)
throws java.io.IOException,
javax.jms.MessageFormatException
java.io.IOExceptionjavax.jms.MessageFormatExceptionprotected static java.lang.Object readPrimitive(java.io.ObjectInput in)
throws java.io.IOException,
java.lang.ClassNotFoundException
writePrimitive(Object, ObjectOutput) otherwise
deserialization will fail and an IOException will be thrownin - the stream to read fromjava.io.IOException - if an I/O error occursjava.lang.ClassNotFoundException - if a class of serialized object cannot be foundpublic java.lang.Object clone()
throws java.lang.CloneNotSupportedException
clone in class java.lang.Objectjava.lang.CloneNotSupportedExceptionprotected static void copyAttributes(RMQMessage rmqMessage, javax.jms.Message message) throws javax.jms.JMSException
“rmqMessage = (RMQMessage) message;”
With conversion as appropriate.
rmqMessage - message filled in with attributesmessage - message from which attributes are gained.javax.jms.JMSException - if attributes cannot be copiedpublic long getJMSDeliveryTime()
throws javax.jms.JMSException
getJMSDeliveryTime in interface javax.jms.Messagejavax.jms.JMSExceptionpublic void setJMSDeliveryTime(long deliveryTime)
throws javax.jms.JMSException
setJMSDeliveryTime in interface javax.jms.Messagejavax.jms.JMSExceptionpublic <T> T getBody(java.lang.Class<T> c)
throws javax.jms.JMSException
getBody in interface javax.jms.Messagejavax.jms.JMSExceptionprotected abstract <T> T doGetBody(java.lang.Class<T> c)
throws javax.jms.JMSException
javax.jms.JMSExceptionCopyright © 2022. All rights reserved.