diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java
index d88b55f555..e939582866 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/ChannelPublishingJmsMessageListener.java
@@ -47,18 +47,18 @@ import org.springframework.util.Assert;
* Message and sends that Message to a channel. If the 'expectReply' value is
* true, it will also wait for a Spring Integration reply Message
* and convert that into a JMS reply.
- *
+ *
* @author Mark Fisher
* @author Juergen Hoeller
* @author Oleg Zhurakousky
*/
-public class ChannelPublishingJmsMessageListener
+public class ChannelPublishingJmsMessageListener
implements SessionAwareMessageListener, InitializingBean, TrackableComponent {
-
+
protected final Log logger = LogFactory.getLog(getClass());
-
+
private volatile boolean expectReply;
-
+
private volatile MessageConverter messageConverter = new SimpleMessageConverter();
private volatile boolean extractRequestPayload = true;
@@ -80,7 +80,7 @@ public class ChannelPublishingJmsMessageListener
private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver();
private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
-
+
private final GatewayDelegate gatewayDelegate = new GatewayDelegate();
/**
@@ -93,27 +93,27 @@ public class ChannelPublishingJmsMessageListener
public void setComponentName(String componentName){
this.gatewayDelegate.setComponentName(componentName);
}
-
+
public void setRequestChannel(MessageChannel requestChannel){
this.gatewayDelegate.setRequestChannel(requestChannel);
}
-
+
public void setReplyChannel(MessageChannel replyChannel){
this.gatewayDelegate.setReplyChannel(replyChannel);
}
-
+
public void setErrorChannel(MessageChannel errorChannel){
this.gatewayDelegate.setErrorChannel(errorChannel);
}
-
+
public void setRequestTimeout(long requestTimeout){
this.gatewayDelegate.setRequestTimeout(requestTimeout);
}
-
+
public void setReplyTimeout(long replyTimeout){
this.gatewayDelegate.setReplyTimeout(replyTimeout);
}
-
+
public void setShouldTrack(boolean shouldTrack) {
this.gatewayDelegate.setShouldTrack(shouldTrack);
}
@@ -125,7 +125,7 @@ public class ChannelPublishingJmsMessageListener
public String getComponentType() {
return this.gatewayDelegate.getComponentType();
}
-
+
/**
* Set the default reply destination to send reply messages to. This will
* be applied in case of a request message that does not carry a
@@ -189,7 +189,7 @@ public class ChannelPublishingJmsMessageListener
* JMSMessageID from the request will be copied into the JMSCorrelationID of the reply
* unless there is already a value in the JMSCorrelationID property of the newly created
* reply Message in which case nothing will be copied. If the JMSCorrelationID of the
- * request Message should be copied into the JMSCorrelationID of the reply Message
+ * request Message should be copied into the JMSCorrelationID of the reply Message
* instead, then this value should be set to "JMSCorrelationID".
* Any other value will be treated as a JMS String Property to be copied as-is
* from the request Message into the reply Message with the same property name.
@@ -200,7 +200,7 @@ public class ChannelPublishingJmsMessageListener
/**
* Specify whether explicit QoS should be enabled for replies
- * (for timeToLive, priority, and deliveryMode settings).
+ * (for timeToLive, priority, and deliveryMode settings).
*/
public void setExplicitQosEnabledForReplies(boolean explicitQosEnabledForReplies) {
this.explicitQosEnabledForReplies = explicitQosEnabledForReplies;
@@ -224,7 +224,7 @@ public class ChannelPublishingJmsMessageListener
* converting between JMS Messages and Spring Integration Messages.
* If none is provided, a {@link SimpleMessageConverter} will
* be used.
- *
+ *
* @param messageConverter
*/
public void setMessageConverter(MessageConverter messageConverter) {
@@ -268,10 +268,10 @@ public class ChannelPublishingJmsMessageListener
logger.debug("converted JMS Message [" + jmsMessage + "] to integration Message payload [" + result + "]");
}
}
-
+
Map headers = headerMapper.toHeaders(jmsMessage);
Message> requestMessage = (result instanceof Message>) ?
- MessageBuilder.fromMessage((Message>) result).copyHeaders(headers).build() :
+ MessageBuilder.fromMessage((Message>) result).copyHeaders(headers).build() :
MessageBuilder.withPayload(result).copyHeaders(headers).build();
if (!this.expectReply) {
this.gatewayDelegate.send(requestMessage);
@@ -292,7 +292,7 @@ public class ChannelPublishingJmsMessageListener
headerMapper.fromHeaders(replyMessage.getHeaders(), jmsReply);
this.copyCorrelationIdFromRequestToReply(jmsMessage, jmsReply);
this.sendReply(jmsReply, destination, session);
- }
+ }
catch (RuntimeException e) {
logger.error("Failed to generate JMS Reply Message from: " + replyResult, e);
throw e;
@@ -308,15 +308,15 @@ public class ChannelPublishingJmsMessageListener
public void afterPropertiesSet() {
this.gatewayDelegate.afterPropertiesSet();
}
-
+
protected void start(){
this.gatewayDelegate.start();
}
-
+
protected void stop(){
this.gatewayDelegate.stop();
}
-
+
private void copyCorrelationIdFromRequestToReply(javax.jms.Message requestMessage, javax.jms.Message replyMessage) throws JMSException {
if (this.correlationKey != null) {
if (this.correlationKey.equals("JMSCorrelationID")) {
@@ -416,17 +416,20 @@ public class ChannelPublishingJmsMessageListener
this.isTopic = isTopic;
}
}
-
+
private class GatewayDelegate extends MessagingGatewaySupport {
+ @Override
protected void send(Object request) {
super.send(request);
}
-
+
+ @Override
protected Message> sendAndReceiveMessage(Object request) {
return super.sendAndReceiveMessage(request);
}
+ @Override
public String getComponentType() {
if (expectReply) {
return "jms:inbound-gateway";
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java
index 91de46dbe3..50e93f207e 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsOutboundGateway.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2011 the original author or authors.
+ * Copyright 2002-2012 the original author or authors.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -16,6 +16,9 @@
package org.springframework.integration.jms;
+import java.util.HashMap;
+import java.util.LinkedList;
+import java.util.List;
import java.util.Map;
import java.util.UUID;
@@ -23,6 +26,7 @@ import javax.jms.Connection;
import javax.jms.ConnectionFactory;
import javax.jms.DeliveryMode;
import javax.jms.Destination;
+import javax.jms.InvalidDestinationException;
import javax.jms.JMSException;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
@@ -31,6 +35,7 @@ import javax.jms.TemporaryQueue;
import javax.jms.TemporaryTopic;
import javax.jms.Topic;
+import org.springframework.beans.DirectFieldAccessor;
import org.springframework.expression.Expression;
import org.springframework.integration.Message;
import org.springframework.integration.MessageChannel;
@@ -40,6 +45,7 @@ import org.springframework.integration.MessageTimeoutException;
import org.springframework.integration.handler.AbstractReplyProducingMessageHandler;
import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor;
import org.springframework.integration.support.MessageBuilder;
+import org.springframework.jms.connection.CachingConnectionFactory;
import org.springframework.jms.connection.ConnectionFactoryUtils;
import org.springframework.jms.support.JmsUtils;
import org.springframework.jms.support.converter.MessageConverter;
@@ -47,6 +53,8 @@ import org.springframework.jms.support.converter.SimpleMessageConverter;
import org.springframework.jms.support.destination.DestinationResolver;
import org.springframework.jms.support.destination.DynamicDestinationResolver;
import org.springframework.util.Assert;
+import org.springframework.util.ObjectUtils;
+import org.springframework.util.StringUtils;
/**
* An outbound Messaging Gateway for request/reply JMS.
@@ -58,6 +66,16 @@ import org.springframework.util.Assert;
*/
public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
+ private final String gatewayId = UUID.randomUUID().toString();
+
+ private volatile boolean cachedConsumers;
+
+ private volatile Map> cachedSessionView;
+
+ private volatile boolean replyDestinationExplicitlySet;
+
+ private final Map tempQueuePerSessionMap = new HashMap(10);
+
private volatile Destination requestDestination;
private volatile String requestDestinationName;
@@ -159,6 +177,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
this.replyPubSubDomain = true;
}
this.replyDestination = replyDestination;
+ this.replyDestinationExplicitlySet = true;
}
/**
@@ -167,6 +186,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
*/
public void setReplyDestinationName(String replyDestinationName) {
this.replyDestinationName = replyDestinationName;
+ this.replyDestinationExplicitlySet = true;
}
/**
@@ -240,15 +260,21 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
/**
- * Provide the name of a JMS property that should hold a generated UUID that
- * the receiver of the JMS Message would expect to represent the CorrelationID.
- * When waiting for the reply Message, a MessageSelector will be configured
- * to match this property name and the UUID value that was sent in the request.
- * If this value is NULL (the default) then the reply consumer's MessageSelector
- * will be expecting the JMSCorrelationID to equal the Message ID of the request.
- * If you want to store the outbound correlation UUID value in the actual
- * JMSCorrelationID property, then set this value to "JMSCorrelationID".
- * However, any other value will be treated as a JMS String Property.
+ * Provide the name of a JMS property that should hold a generated JMSCorrelationID.
+ * If NO value is provided for this attribute, then the reply consumer's
+ * MessageSelector will be expecting the JMSCorrelationID to equal the Message ID
+ * of the request. However this would also mean that the MessageSelector will be different
+ * for each new message, thus requiring a new reply consumer to be created for each Message sent.
+ * Although this is a very common JMS exchange pattern it is not optimal for high throughput request/reply
+ * applications. Most MOMs (including Spring Integration's Inbound Gateway) also support propagation
+ * of the 'JMSCorrelationID' of the consumed request message into the 'JMSCorrelationID' of the newly produced
+ * reply message. This also allows us to use an optimized value generation algorithm which results in caching
+ * and reuse of reply consumers resulting in significant performance improvement of the request/reply
+ * scenarios. So if you have a Spring JMS consumer (e.g., Spring Integration's Inbound Gateway) or
+ * you know that your MOM supports propagation of the 'CorrelationID' it is highly recommended to set this
+ * value to 'JMSCorrelationID'. If you set this value to something else the receiving end must now how
+ * to propagate it. For example: if the receiving end is an inbound-gateway its correlation-key attribute
+ * must have the same value.
*/
public void setCorrelationKey(String correlationKey) {
this.correlationKey = correlationKey;
@@ -317,48 +343,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
return "jms:outbound-gateway";
}
- private Destination getRequestDestination(Message> message, Session session) throws JMSException {
- if (this.requestDestination != null) {
- return this.requestDestination;
- }
- if (this.requestDestinationName != null) {
- return this.resolveRequestDestination(this.requestDestinationName, session);
- }
- if (this.requestDestinationExpressionProcessor != null) {
- Object result = this.requestDestinationExpressionProcessor.processMessage(message);
- if (result instanceof Destination) {
- return (Destination) result;
- }
- if (result instanceof String) {
- return this.resolveRequestDestination((String) result, session);
- }
- throw new MessageDeliveryException(message,
- "Evaluation of requestDestinationExpression failed to produce a Destination or destination name. Result was: " + result);
- }
- throw new MessageDeliveryException(message,
- "No requestDestination, requestDestinationName, or requestDestinationExpression has been configured.");
- }
-
- private Destination resolveRequestDestination(String requestDestinationName, Session session) throws JMSException {
- Assert.notNull(this.destinationResolver,
- "DestinationResolver is required when relying upon the 'requestDestinationName' property.");
- return this.destinationResolver.resolveDestinationName(
- session, requestDestinationName, this.requestPubSubDomain);
- }
-
- private Destination getReplyDestination(Session session) throws JMSException {
- if (this.replyDestination != null) {
- return this.replyDestination;
- }
- if (this.replyDestinationName != null) {
- Assert.notNull(this.destinationResolver,
- "DestinationResolver is required when relying upon the 'replyDestinationName' property.");
- return this.destinationResolver.resolveDestinationName(
- session, this.replyDestinationName, this.replyPubSubDomain);
- }
- return session.createTemporaryQueue();
- }
-
+ @SuppressWarnings("unchecked")
@Override
public final void onInit() {
synchronized (this.initializationMonitor) {
@@ -376,6 +361,24 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
this.requestDestinationExpressionProcessor.setConversionService(getConversionService());
}
this.initialized = true;
+ if (this.connectionFactory instanceof CachingConnectionFactory){
+ this.cachedConsumers = ((CachingConnectionFactory)this.connectionFactory).isCacheConsumers();
+ DirectFieldAccessor cfDfa = new DirectFieldAccessor(this.connectionFactory);
+ this.cachedSessionView = (Map>) cfDfa.getPropertyValue("cachedSessions");
+ }
+
+ if (this.cachedConsumers && !StringUtils.hasText(this.correlationKey)
+ && (this.replyDestination != null || StringUtils.hasText(this.replyDestinationName))){
+ logger.warn("Caching consumers without custom correlation-key can lead to " +
+ "significant performance degradation and OutOfMemoryError. Either do not use consumer caching " +
+ "by setting 'cacheConsumers' attribute to FALSE in the CachingConnectionFactory or set 'correlation-key'" +
+ "to a value (e.g., correlation-key=\"JMSCorrelationID\")");
+ }
+ else if (!this.cachedConsumers && StringUtils.hasText(this.correlationKey)){
+ logger.warn("Using custom 'correlationKey' without caching consumers can lead to " +
+ "significant performance degradation. Consider using the CachingConnectionFactory " +
+ "with 'cacheConsumers' attribute set to TRUE");
+ }
}
}
@@ -413,12 +416,93 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
}
+ /**
+ * Create a new JMS Connection for this JMS gateway.
+ */
+ protected Connection createConnection() throws JMSException {
+ return this.connectionFactory.createConnection();
+ }
+
+ /**
+ * Create a new JMS Session using the provided Connection.
+ */
+ protected Session createSession(Connection connection) throws JMSException {
+ return connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
+ }
+
+ private Destination determineRequestDestination(Message> message, Session session) throws JMSException {
+ if (this.requestDestination != null) {
+ return this.requestDestination;
+ }
+
+ if (this.requestDestinationName != null) {
+ return this.resolveRequestDestination(this.requestDestinationName, session);
+ }
+
+ if (this.requestDestinationExpressionProcessor != null) {
+ Object result = this.requestDestinationExpressionProcessor.processMessage(message);
+ if (result instanceof Destination) {
+ return (Destination) result;
+ }
+ if (result instanceof String) {
+ return this.resolveRequestDestination((String) result, session);
+ }
+ throw new MessageDeliveryException(message,
+ "Evaluation of requestDestinationExpression failed to produce a Destination or destination name. Result was: " + result);
+ }
+ throw new MessageDeliveryException(message,
+ "No requestDestination, requestDestinationName, or requestDestinationExpression has been configured.");
+ }
+
+ private Destination resolveRequestDestination(String requestDestinationName, Session session) throws JMSException {
+ Assert.notNull(this.destinationResolver,
+ "DestinationResolver is required when relying upon the 'requestDestinationName' property.");
+ return this.destinationResolver.resolveDestinationName(
+ session, requestDestinationName, this.requestPubSubDomain);
+ }
+
+ private Destination determineReplyDestination(Session session, long sessionId) throws JMSException {
+
+ Destination replyDestinationToReturn = this.replyDestination;
+
+ if (replyDestinationToReturn == null){
+ if (StringUtils.hasText(this.replyDestinationName)) {
+ Assert.notNull(this.destinationResolver,
+ "DestinationResolver is required when relying upon the 'replyDestinationName' property.");
+ replyDestinationToReturn = this.destinationResolver.resolveDestinationName(
+ session, this.replyDestinationName, this.replyPubSubDomain);
+ }
+ else if (this.tempQueuePerSessionMap.containsKey(sessionId)){
+ replyDestinationToReturn = this.tempQueuePerSessionMap.get(sessionId);
+ }
+ else if (this.cachedConsumers){
+ if (StringUtils.hasText(this.correlationKey)){
+ // will rely on LIKE MessageSelector. Only need one TemporaryQueue for reply
+ // as if the replyDestinationName was specified
+ replyDestinationToReturn = session.createTemporaryQueue();
+ this.replyDestination = replyDestinationToReturn;
+ }
+ else {
+ replyDestinationToReturn = session.createTemporaryQueue();
+ this.tempQueuePerSessionMap.put(sessionId, replyDestinationToReturn);
+ }
+ }
+ else {
+ replyDestinationToReturn = session.createTemporaryQueue();
+ }
+ }
+
+ return replyDestinationToReturn;
+ }
+
private javax.jms.Message sendAndReceive(Message> requestMessage) throws JMSException {
Connection connection = this.createConnection();
Session session = null;
Destination replyTo = null;
+ long sessionId = 0;
try {
session = this.createSession(connection);
+ sessionId = System.identityHashCode(session);
// convert to JMS Message
Object objectToSend = requestMessage;
@@ -431,111 +515,240 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
headerMapper.fromHeaders(requestMessage.getHeaders(), jmsRequest);
// TODO: support a JmsReplyTo header in the SI Message?
- replyTo = this.getReplyDestination(session);
- jmsRequest.setJMSReplyTo(replyTo);
+ replyTo = this.determineReplyDestination(session, sessionId);
+
connection.start();
Integer priority = requestMessage.getHeaders().getPriority();
if (priority == null) {
priority = this.priority;
}
- javax.jms.Message replyMessage = null;
- Destination requestDestination = this.getRequestDestination(requestMessage, session);
- if (this.correlationKey != null) {
- replyMessage = this.doSendAndReceiveWithGeneratedCorrelationId(requestDestination, jmsRequest, replyTo, session, priority);
- }
- else if (replyTo instanceof TemporaryQueue || replyTo instanceof TemporaryTopic) {
- replyMessage = this.doSendAndReceiveWithTemporaryReplyToDestination(requestDestination, jmsRequest, replyTo, session, priority);
- }
- else {
- replyMessage = this.doSendAndReceiveWithMessageIdCorrelation(requestDestination, jmsRequest, replyTo, session, priority);
- }
+
+ Destination requestDestinationToUse = this.determineRequestDestination(requestMessage, session);
+
+ javax.jms.Message replyMessage = this.doSendAndReceive(requestDestinationToUse, jmsRequest, replyTo, session, priority, sessionId);
+
return replyMessage;
}
finally {
JmsUtils.closeSession(session);
- this.deleteDestinationIfTemporary(replyTo);
+
+ if (!this.replyDestinationExplicitlySet &&
+ this.isTemporaryDestination(replyTo) &&
+ !this.isCachedSession(session)) {
+
+ this.deleteDestinationIfTemporary(replyTo);
+ this.tempQueuePerSessionMap.remove(sessionId);
+ }
+
ConnectionFactoryUtils.releaseConnection(connection, this.connectionFactory, true);
}
}
- /**
- * Creates the MessageConsumer before sending the request Message since we are generating our own correlationId value for the MessageSelector.
- */
- private javax.jms.Message doSendAndReceiveWithGeneratedCorrelationId(Destination requestDestination,
- javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
- MessageProducer messageProducer = null;
- MessageConsumer messageConsumer = null;
- try {
- messageProducer = session.createProducer(requestDestination);
- String correlationId = UUID.randomUUID().toString().replaceAll("'", "''");
- Assert.state(this.correlationKey != null, "correlationKey must not be null");
- String messageSelector = null;
- if (this.correlationKey.equals("JMSCorrelationID")) {
- jmsRequest.setJMSCorrelationID(correlationId);
- messageSelector = "JMSCorrelationID = '" + correlationId + "'";
- }
- else {
- jmsRequest.setStringProperty(this.correlationKey, correlationId);
- messageSelector = this.correlationKey + " = '" + correlationId + "'";
- }
- messageConsumer = session.createConsumer(replyTo, messageSelector);
- this.sendRequestMessage(jmsRequest, messageProducer, priority);
- return this.receiveReplyMessage(messageConsumer);
- }
- finally {
- JmsUtils.closeMessageProducer(messageProducer);
- JmsUtils.closeMessageConsumer(messageConsumer);
- }
- }
-
- /**
- * Creates the MessageConsumer before sending the request Message since we do not need any correlation.
- */
- private javax.jms.Message doSendAndReceiveWithTemporaryReplyToDestination(Destination requestDestination,
- javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
- MessageProducer messageProducer = null;
- MessageConsumer messageConsumer = null;
- try {
- messageProducer = session.createProducer(requestDestination);
- messageConsumer = session.createConsumer(replyTo);
- this.sendRequestMessage(jmsRequest, messageProducer, priority);
- return this.receiveReplyMessage(messageConsumer);
- }
- finally {
- JmsUtils.closeMessageProducer(messageProducer);
- JmsUtils.closeMessageConsumer(messageConsumer);
- }
- }
-
/**
* Creates the MessageConsumer after sending the request Message since we need the MessageID for correlation with a MessageSelector.
*/
- private javax.jms.Message doSendAndReceiveWithMessageIdCorrelation(Destination requestDestination,
- javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority) throws JMSException {
- if (replyTo instanceof Topic && logger.isWarnEnabled()) {
- logger.warn("Relying on the MessageID for correlation is not recommended when using a Topic as the replyTo Destination " +
- "because that ID can only be provided to a MessageSelector after the reuqest Message has been sent thereby " +
- "creating a race condition where a fast response might be sent before the MessageConsumer has been created. " +
- "Consider providing a value to the 'correlationKey' property of this gateway instead. Then the MessageConsumer " +
- "will be created before the request Message is sent.");
- }
- MessageProducer messageProducer = null;
- MessageConsumer messageConsumer = null;
+ private javax.jms.Message doSendAndReceive(Destination requestDestinationToUse,
+ javax.jms.Message jmsRequest, Destination replyTo, Session session, int priority, long sessionId) throws JMSException {
+
+ MessageProducer messageProducer = session.createProducer(requestDestinationToUse);
+ jmsRequest.setJMSReplyTo(replyTo);
try {
- messageProducer = session.createProducer(requestDestination);
- this.sendRequestMessage(jmsRequest, messageProducer, priority);
- String messageId = jmsRequest.getJMSMessageID().replaceAll("'", "''");
- String messageSelector = "JMSCorrelationID = '" + messageId + "'";
- messageConsumer = session.createConsumer(replyTo, messageSelector);
- return this.receiveReplyMessage(messageConsumer);
+ if (StringUtils.hasText(this.correlationKey)){
+ return this.exchangeWithCorrelationKey(messageProducer, jmsRequest, session, sessionId, replyTo, priority);
+ }
+ else {
+ return this.exchangeWithoutCorrelationKey(messageProducer, jmsRequest, session, sessionId, replyTo, priority);
+ }
}
finally {
JmsUtils.closeMessageProducer(messageProducer);
+ }
+ }
+
+
+ private javax.jms.Message exchangeWithCorrelationKey(MessageProducer messageProducer, javax.jms.Message jmsRequest,
+ Session session, long sessionId, Destination replyTo, int priority) throws JMSException {
+
+ String messageCorrelationId = ObjectUtils.getIdentityHexString(jmsRequest);
+
+ String consumerCorrelationId = gatewayId + "_" + sessionId;
+
+ if (this.correlationKey.equals("JMSCorrelationID")) {
+ jmsRequest.setJMSCorrelationID(consumerCorrelationId + "$" + messageCorrelationId);
+ }
+ else {
+ jmsRequest.setStringProperty(correlationKey, consumerCorrelationId + "$" + messageCorrelationId);
+ }
+
+ String messageSelector = correlationKey + " LIKE '" + consumerCorrelationId + "%'";
+
+ MessageConsumer messageConsumer = null;
+
+ try {
+ messageConsumer = this.createMessageConsumer(session, replyTo, messageSelector, sessionId);
+
+ this.sendRequestMessage(jmsRequest, messageProducer, priority);
+
+ return this.doReceive(messageConsumer, this.correlationKey, messageCorrelationId, sessionId);
+ }
+ finally {
JmsUtils.closeMessageConsumer(messageConsumer);
}
}
+ private javax.jms.Message exchangeWithoutCorrelationKey(MessageProducer messageProducer, javax.jms.Message jmsRequest,
+ Session session, long sessionId, Destination replyTo, int priority) throws JMSException {
+
+ MessageConsumer messageConsumer = null;
+ try {
+ String messageCorrelationId = null;
+
+ /*
+ * When using the message id for correlation, we must null any correlation id because we
+ * may have an upstream gateway waiting on this correlation id; consider:
+ * ob->ib->ob->ib->ob->ib
+ * If we propagate the correlation id, the third ob will attempt to correlate on
+ * the same selector.
+ */
+ if (jmsRequest.getJMSCorrelationID() != null) {
+ jmsRequest.setJMSCorrelationID(null);
+ }
+
+ if (this.replyDestinationExplicitlySet) {
+ if (replyTo instanceof Topic && logger.isWarnEnabled()) {
+ logger.warn("Relying on the MessageID for correlation is not recommended when using a Topic as the replyTo Destination " +
+ "because that ID can only be provided to a MessageSelector after the request Message has been sent thereby " +
+ "creating a race condition where a fast response might be sent before the MessageConsumer has been created. " +
+ "Consider providing a value to the 'correlationKey' property of this gateway instead. Then the MessageConsumer " +
+ "will be created before the request Message is sent.");
+ }
+ this.sendRequestMessage(jmsRequest, messageProducer, priority);
+
+ messageCorrelationId = jmsRequest.getJMSMessageID().replaceAll("'", "''");
+ String messageSelector = "JMSCorrelationID = '" + messageCorrelationId + "'";
+ messageConsumer = this.createMessageConsumer(session, replyTo, messageSelector, sessionId);
+ }
+ else {
+ messageConsumer = this.createMessageConsumer(session, replyTo, null, sessionId);
+ this.sendRequestMessage(jmsRequest, messageProducer, priority);
+ messageCorrelationId = jmsRequest.getJMSMessageID().replaceAll("'", "''");
+ }
+
+ return this.doReceive(messageConsumer, null, messageCorrelationId, sessionId);
+ }
+ finally {
+ JmsUtils.closeMessageConsumer(messageConsumer);
+ }
+ }
+
+ /**
+ * Will perform clean up if failures detected during message receive
+ */
+ private javax.jms.Message doReceive(MessageConsumer messageConsumer, String correlationKey, String messageCorrelationIdToMatch,
+ long sessionId) throws JMSException{
+ try {
+ return this.receiveCorrelatedReplyMessage(messageConsumer, correlationKey, messageCorrelationIdToMatch);
+ }
+ catch (javax.jms.IllegalStateException e) {
+ this.clearDestination(sessionId);
+ throw e;
+ }
+ }
+
+ private javax.jms.Message receiveCorrelatedReplyMessage(MessageConsumer messageConsumer,
+ String correlationKey, String messageCorrelationIdToMatch) throws JMSException {
+
+ long timeLeft = this.receiveTimeout < 0 ? Long.MAX_VALUE : this.receiveTimeout;
+
+ javax.jms.Message replyMessage = null;
+ boolean replyFound = false;
+
+ while (!replyFound && timeLeft > 0){
+ long startTime = System.currentTimeMillis();
+
+ replyMessage = (this.receiveTimeout >= 0) ? messageConsumer.receive(timeLeft) : messageConsumer.receive();
+
+ if (replyMessage != null) {
+ replyFound = checkReplyMessage(correlationKey, messageCorrelationIdToMatch, replyMessage);
+ if (!replyFound) {
+ timeLeft -= (System.currentTimeMillis() - startTime);
+ }
+ else {
+ replyMessage.setJMSCorrelationID(null);
+ }
+ }
+ else {
+ timeLeft = 0;
+ }
+ }
+
+ return replyFound ? replyMessage : null;
+ }
+
+ /**
+ * Examine the reply to determine if it is the one we are expecting
+ * @return true if the reply is correlated.
+ */
+ private boolean checkReplyMessage(String correlationKey, String messageCorrelationIdToMatch,
+ javax.jms.Message replyMessage) throws JMSException {
+ String jmsCorrelationId = null;
+ boolean replyFound;
+ if (correlationKey != null && !correlationKey.equals("JMSCorrelationID")){
+ jmsCorrelationId = replyMessage.getStringProperty(correlationKey);
+ }
+ else {
+ jmsCorrelationId = replyMessage.getJMSCorrelationID();
+ }
+
+ if (StringUtils.hasText(jmsCorrelationId) && StringUtils.hasText(messageCorrelationIdToMatch)){
+
+ if (jmsCorrelationId.endsWith(messageCorrelationIdToMatch)){
+ replyFound = true;
+ }
+ else {
+ // Essentially we are discarding the uncorrelated message here and moving on to
+ // the next one since we can no longer communicate with its originating producer
+ // since it has already timed out waiting for the reply. We are also honoring the original timeout
+ if (this.logger.isDebugEnabled()){
+ this.logger.debug("Discarded late arriving reply: " + replyMessage);
+ }
+ replyFound = false;
+ }
+ }
+ else {
+ replyFound = true;
+ }
+ return replyFound;
+ }
+
+ private MessageConsumer createMessageConsumer(Session session, Destination replyTo, String messageSelector, long sessionId) throws JMSException{
+ MessageConsumer consumer;
+ try {
+ if (this.cachedConsumers && this.isTemporaryDestination(replyTo)){
+ if (StringUtils.hasText(messageSelector)){
+ consumer = TemporaryConsumerUtils.createConsumer(session, replyTo, messageSelector);
+ }
+ else {
+ consumer = TemporaryConsumerUtils.createConsumer(session, replyTo, null);
+ }
+ }
+ else {
+ if (StringUtils.hasText(messageSelector)){
+ consumer = session.createConsumer(replyTo, messageSelector);
+ }
+ else {
+ consumer = session.createConsumer(replyTo);
+ }
+ }
+ }
+ catch (InvalidDestinationException e) {
+ this.clearDestination(sessionId);
+ throw e;
+ }
+ return consumer;
+ }
+
private void sendRequestMessage(javax.jms.Message jmsRequest, MessageProducer messageProducer, int priority) throws JMSException {
if (this.explicitQosEnabled) {
messageProducer.send(jmsRequest, this.deliveryMode, priority, this.timeToLive);
@@ -545,10 +758,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
}
- private javax.jms.Message receiveReplyMessage(MessageConsumer messageConsumer) throws JMSException {
- return (this.receiveTimeout >= 0) ? messageConsumer.receive(receiveTimeout) : messageConsumer.receive();
- }
-
/**
* Deletes either a {@link TemporaryQueue} or {@link TemporaryTopic}.
* Ignores any other {@link Destination} type and also ignores any
@@ -569,17 +778,44 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
/**
- * Create a new JMS Connection for this JMS gateway.
+ * Will clear reply destination *only* if such reply destination is temporary
+ * and was created internally
*/
- protected Connection createConnection() throws JMSException {
- return this.connectionFactory.createConnection();
+ private void clearDestination(long sessionId){
+ if (this.tempQueuePerSessionMap.containsKey(sessionId)){
+ this.tempQueuePerSessionMap.remove(sessionId);
+ }
+ else if (this.isTemporaryDestination(this.replyDestination) && !this.replyDestinationExplicitlySet){
+ this.replyDestination = null;
+ }
}
/**
- * Create a new JMS Session using the provided Connection.
+ * Checks if this destination (topic or queue) is a temporary destination
*/
- protected Session createSession(Connection connection) throws JMSException {
- return connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
+ private boolean isTemporaryDestination(Destination destination){
+ return destination instanceof TemporaryQueue || destination instanceof TemporaryTopic;
}
+ /**
+ * Validates if this session is cached by the CachingConnectionFactory
+ * Ignores any exception that may be thrown
+ */
+ private boolean isCachedSession(Session session){
+ try {
+ // this.cacheConsumers means we are using CCF
+ if (this.cachedConsumers && this.cachedSessionView != null &&
+ this.cachedSessionView.containsKey(Session.AUTO_ACKNOWLEDGE)){
+ List cachedSessions = this.cachedSessionView.get(Session.AUTO_ACKNOWLEDGE);
+ if (cachedSessions != null){
+ return cachedSessions.contains(session);
+ }
+ }
+ } catch (Exception e) {
+ if (logger.isTraceEnabled()) {
+ logger.trace("Exception while detecting cached session", e);
+ }
+ }
+ return false;
+ }
}
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/TemporaryConsumerUtils.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/TemporaryConsumerUtils.java
new file mode 100644
index 0000000000..9f3ff846a2
--- /dev/null
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/TemporaryConsumerUtils.java
@@ -0,0 +1,114 @@
+/*
+ * Copyright 2002-2012 the original author or authors.
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.springframework.integration.jms;
+
+import java.lang.reflect.Constructor;
+import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.Proxy;
+import java.util.Map;
+
+import javax.jms.Destination;
+import javax.jms.JMSException;
+import javax.jms.MessageConsumer;
+import javax.jms.Session;
+
+import org.apache.commons.logging.Log;
+import org.apache.commons.logging.LogFactory;
+
+import org.springframework.beans.DirectFieldAccessor;
+/**
+ * NOT A PUBLIC API
+ *
+ * Temporary class, will be removed as soon as
+ * https://jira.springsource.org/browse/SPR-9648 is addressed
+ *
+ * @author Oleg Zhurakousky
+ * @since 2.2
+ */
+class TemporaryConsumerUtils {
+
+ private static volatile Constructor> consumerCacheKeyConstructor;
+
+ private static volatile Constructor> cachedMessageConsumerConstructor;
+
+ private final static Log logger = LogFactory.getLog(TemporaryConsumerUtils.class);
+
+ public static MessageConsumer createConsumer(Session session, Destination destination, String messageSelector)
+ throws JMSException {
+ if (Proxy.isProxyClass(session.getClass())){
+ InvocationHandler invocationHandler = Proxy.getInvocationHandler(session);
+ DirectFieldAccessor ihDa = new DirectFieldAccessor(invocationHandler);
+ @SuppressWarnings("unchecked")
+ Map