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 e939582866..d88b55f555 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,20 +416,17 @@ 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 50e93f207e..91de46dbe3 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-2012 the original author or authors.
+ * Copyright 2002-2011 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,9 +16,6 @@
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;
@@ -26,7 +23,6 @@ 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;
@@ -35,7 +31,6 @@ 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;
@@ -45,7 +40,6 @@ 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;
@@ -53,8 +47,6 @@ 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.
@@ -66,16 +58,6 @@ import org.springframework.util.StringUtils;
*/
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;
@@ -177,7 +159,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
this.replyPubSubDomain = true;
}
this.replyDestination = replyDestination;
- this.replyDestinationExplicitlySet = true;
}
/**
@@ -186,7 +167,6 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
*/
public void setReplyDestinationName(String replyDestinationName) {
this.replyDestinationName = replyDestinationName;
- this.replyDestinationExplicitlySet = true;
}
/**
@@ -260,21 +240,15 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
/**
- * 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.
+ * 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.
*/
public void setCorrelationKey(String correlationKey) {
this.correlationKey = correlationKey;
@@ -343,7 +317,48 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
return "jms:outbound-gateway";
}
- @SuppressWarnings("unchecked")
+ 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();
+ }
+
@Override
public final void onInit() {
synchronized (this.initializationMonitor) {
@@ -361,24 +376,6 @@ 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");
- }
}
}
@@ -416,93 +413,12 @@ 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;
@@ -515,240 +431,111 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
headerMapper.fromHeaders(requestMessage.getHeaders(), jmsRequest);
// TODO: support a JmsReplyTo header in the SI Message?
- replyTo = this.determineReplyDestination(session, sessionId);
-
+ replyTo = this.getReplyDestination(session);
+ jmsRequest.setJMSReplyTo(replyTo);
connection.start();
Integer priority = requestMessage.getHeaders().getPriority();
if (priority == null) {
priority = this.priority;
}
-
- Destination requestDestinationToUse = this.determineRequestDestination(requestMessage, session);
-
- javax.jms.Message replyMessage = this.doSendAndReceive(requestDestinationToUse, jmsRequest, replyTo, session, priority, sessionId);
-
+ 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);
+ }
return replyMessage;
}
finally {
JmsUtils.closeSession(session);
-
- if (!this.replyDestinationExplicitlySet &&
- this.isTemporaryDestination(replyTo) &&
- !this.isCachedSession(session)) {
-
- this.deleteDestinationIfTemporary(replyTo);
- this.tempQueuePerSessionMap.remove(sessionId);
- }
-
+ this.deleteDestinationIfTemporary(replyTo);
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 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);
+ 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;
try {
- if (StringUtils.hasText(this.correlationKey)){
- return this.exchangeWithCorrelationKey(messageProducer, jmsRequest, session, sessionId, replyTo, priority);
- }
- else {
- return this.exchangeWithoutCorrelationKey(messageProducer, jmsRequest, session, sessionId, replyTo, priority);
- }
+ 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);
}
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);
@@ -758,6 +545,10 @@ 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
@@ -778,44 +569,17 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
}
/**
- * Will clear reply destination *only* if such reply destination is temporary
- * and was created internally
+ * Create a new JMS Connection for this JMS gateway.
*/
- 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;
- }
+ protected Connection createConnection() throws JMSException {
+ return this.connectionFactory.createConnection();
}
/**
- * Checks if this destination (topic or queue) is a temporary destination
+ * Create a new JMS Session using the provided Connection.
*/
- private boolean isTemporaryDestination(Destination destination){
- return destination instanceof TemporaryQueue || destination instanceof TemporaryTopic;
+ protected Session createSession(Connection connection) throws JMSException {
+ return connection.createSession(false, Session.AUTO_ACKNOWLEDGE);
}
- /**
- * 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
deleted file mode 100644
index 9f3ff846a2..0000000000
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/TemporaryConsumerUtils.java
+++ /dev/null
@@ -1,114 +0,0 @@
-/*
- * 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