diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java
index 4026c71b32..4c33f0b18a 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/AbstractJmsTemplateBasedAdapter.java
@@ -31,8 +31,11 @@ import org.springframework.util.Assert;
* Base class for adapters that delegate to a {@link JmsTemplate}.
*
* @author Mark Fisher
+ * @author Oleg Zhurakousky
*/
public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBean {
+
+ private volatile boolean extractPayload = true;
private volatile ConnectionFactory connectionFactory;
@@ -56,7 +59,7 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
private volatile MessageConverter messageConverter;
- private volatile JmsHeaderMapper headerMapper;
+ private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
private volatile boolean initialized;
@@ -84,7 +87,18 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
public AbstractJmsTemplateBasedAdapter() {
}
-
+ /**
+ * Specify whether the payload should be extracted from each received JMS
+ * Message to be used as the Spring Integration Message payload.
+ *
+ *
The default value is true. To force creation of Spring
+ * Integration Messages whose payload is the actual JMS Message, set this
+ * to false.
+ */
+ public void setExtractPayload(boolean extractPayload) {
+ this.extractPayload = extractPayload;
+ }
+
public void setConnectionFactory(ConnectionFactory connectionFactory) {
this.connectionFactory = connectionFactory;
}
@@ -105,13 +119,13 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
* Provide a {@link MessageConverter} strategy to use for converting
* between Spring Integration Messages and JMS Messages.
*
- * The default is a {@link HeaderMappingMessageConverter} that delegates to
+ * The default is a {@link DefaultMessageConverter} that delegates to
* a {@link SimpleMessageConverter}.
*/
public void setMessageConverter(MessageConverter messageConverter) {
this.messageConverter = messageConverter;
}
-
+
public void setDestinationResolver(DestinationResolver destinationResolver) {
this.destinationResolver = destinationResolver;
}
@@ -119,6 +133,10 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
public void setHeaderMapper(JmsHeaderMapper headerMapper) {
this.headerMapper = headerMapper;
}
+
+ JmsHeaderMapper getHeaderMapper(){
+ return this.headerMapper;
+ }
/**
* @see JmsTemplate#setExplicitQosEnabled(boolean)
@@ -182,7 +200,7 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
if (this.messageConverter != null) {
this.jmsTemplate.setMessageConverter(this.messageConverter);
}
- this.configureMessageConverter(this.jmsTemplate, this.headerMapper);
+ this.configureMessageConverter(this.jmsTemplate);
this.initialized = true;
}
}
@@ -203,6 +221,12 @@ public abstract class AbstractJmsTemplateBasedAdapter implements InitializingBea
return jmsTemplate;
}
- protected abstract void configureMessageConverter(JmsTemplate jmsTemplate, JmsHeaderMapper headerMapper);
-
+ protected void configureMessageConverter(JmsTemplate jmsTemplate) {
+ MessageConverter converter = jmsTemplate.getMessageConverter();
+ if (converter == null || !(converter instanceof DefaultMessageConverter)) {
+ DefaultMessageConverter dmc = new DefaultMessageConverter(converter);
+ dmc.setExtractIntegrationMessagePayload(this.extractPayload);
+ jmsTemplate.setMessageConverter(dmc);
+ }
+ }
}
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 2e3fcb431c..c1087753a1 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
@@ -66,7 +66,7 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa
private volatile DestinationResolver destinationResolver = new DynamicDestinationResolver();
- private volatile JmsHeaderMapper headerMapper;
+ private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
/**
* Specify whether a JMS reply Message is expected.
@@ -155,7 +155,7 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa
/**
* Provide a {@link MessageConverter} implementation to use when
* converting between JMS Messages and Spring Integration Messages.
- * If none is provided, a {@link HeaderMappingMessageConverter} will
+ * If none is provided, a {@link DefaultMessageConverter} will
* be used and the {@link JmsHeaderMapper} instance provided to the
* {@link #setHeaderMapper(JmsHeaderMapper)} method will be included
* in the conversion process.
@@ -172,7 +172,7 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa
*
This property will be ignored if a {@link MessageConverter} is
* provided to the {@link #setMessageConverter(MessageConverter)} method.
* However, you may provide your own implementation of the delegating
- * {@link HeaderMappingMessageConverter} implementation.
+ * {@link DefaultMessageConverter} implementation.
*/
public void setHeaderMapper(JmsHeaderMapper headerMapper) {
this.headerMapper = headerMapper;
@@ -199,8 +199,8 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa
}
public final void onInit() throws Exception {
- if (!(this.messageConverter instanceof HeaderMappingMessageConverter)) {
- HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(this.messageConverter, this.headerMapper);
+ if (!(this.messageConverter instanceof DefaultMessageConverter)) {
+ DefaultMessageConverter hmmc = new DefaultMessageConverter(this.messageConverter);
hmmc.setExtractJmsMessageBody(this.extractRequestPayload);
hmmc.setExtractIntegrationMessagePayload(this.extractReplyPayload);
this.messageConverter = hmmc;
@@ -220,7 +220,11 @@ public class ChannelPublishingJmsMessageListener extends AbstractMessagingGatewa
if (replyMessage != null) {
Destination destination = this.getReplyDestination(jmsMessage, session);
if (destination != null){
+ // convert SI Message to JMS Message
javax.jms.Message jmsReply = this.messageConverter.toMessage(replyMessage, session);
+ // map SI Message Headers to JMS Message Properties/Headers
+ headerMapper.fromHeaders(replyMessage.getHeaders(), jmsReply);
+
if (jmsReply.getJMSCorrelationID() == null) {
jmsReply.setJMSCorrelationID(jmsMessage.getJMSMessageID());
}
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/HeaderMappingMessageConverter.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/DefaultMessageConverter.java
similarity index 78%
rename from spring-integration-jms/src/main/java/org/springframework/integration/jms/HeaderMappingMessageConverter.java
rename to spring-integration-jms/src/main/java/org/springframework/integration/jms/DefaultMessageConverter.java
index 1d7c450fd3..8beb42b87a 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/HeaderMappingMessageConverter.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/DefaultMessageConverter.java
@@ -16,16 +16,12 @@
package org.springframework.integration.jms;
-import java.util.Map;
-
import javax.jms.JMSException;
import javax.jms.Session;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
-
import org.springframework.integration.core.Message;
-import org.springframework.integration.core.MessageHeaders;
import org.springframework.integration.message.MessageBuilder;
import org.springframework.jms.support.converter.MessageConversionException;
import org.springframework.jms.support.converter.MessageConverter;
@@ -33,9 +29,8 @@ import org.springframework.jms.support.converter.SimpleMessageConverter;
/**
* A {@link MessageConverter} implementation that is capable of delegating to
- * an existing converter instance and an existing {@link JmsHeaderMapper}. The
- * default MessageConverter implementation is {@link SimpleMessageConverter},
- * and the default header mapper implementation is {@link DefaultJmsHeaderMapper}.
+ * an existing converter instance. The default MessageConverter implementation
+ * is {@link SimpleMessageConverter}.
*
*
If 'extractJmsMessageBody' is true (the default), the body
* of each received JMS Message will become the payload of a Spring Integration
@@ -51,15 +46,14 @@ import org.springframework.jms.support.converter.SimpleMessageConverter;
* specified for Message extraction.
*
* @author Mark Fisher
+ * @author Oleg Zhurakousky
*/
-public class HeaderMappingMessageConverter implements MessageConverter {
+public class DefaultMessageConverter implements MessageConverter {
private final Log logger = LogFactory.getLog(this.getClass());
private final MessageConverter converter;
- private final JmsHeaderMapper headerMapper;
-
private volatile boolean extractJmsMessageBody = true;
private volatile boolean extractIntegrationMessagePayload = true;
@@ -69,8 +63,8 @@ public class HeaderMappingMessageConverter implements MessageConverter {
* Create a HeaderMappingMessageConverter instance that will rely on the
* default {@link SimpleMessageConverter} and {@link DefaultJmsHeaderMapper}.
*/
- public HeaderMappingMessageConverter() {
- this(null, null);
+ public DefaultMessageConverter() {
+ this(null);
}
/**
@@ -78,26 +72,8 @@ public class HeaderMappingMessageConverter implements MessageConverter {
* the provided {@link MessageConverter} instance and will use the default
* implementation of the {@link JmsHeaderMapper} strategy.
*/
- public HeaderMappingMessageConverter(MessageConverter converter) {
- this(converter, null);
- }
-
- /**
- * Create a HeaderMappingMessageConverter instance that will delegate to
- * the provided {@link JmsHeaderMapper} instance and will use the default
- * {@link SimpleMessageConverter} implementation.
- */
- public HeaderMappingMessageConverter(JmsHeaderMapper headerMapper) {
- this(null, headerMapper);
- }
-
- /**
- * Create a HeaderMappingMessageConverter instance that will delegate to
- * the provided {@link MessageConverter} and {@link JmsHeaderMapper}.
- */
- public HeaderMappingMessageConverter(MessageConverter converter, JmsHeaderMapper headerMapper) {
+ public DefaultMessageConverter(MessageConverter converter) {
this.converter = (converter != null ? converter : new SimpleMessageConverter());
- this.headerMapper = (headerMapper != null ? headerMapper : new DefaultJmsHeaderMapper());
}
/**
@@ -152,8 +128,7 @@ public class HeaderMappingMessageConverter implements MessageConverter {
else {
builder = MessageBuilder.withPayload(jmsMessage);
}
- Map headers = this.headerMapper.toHeaders(jmsMessage);
- Message> message = builder.copyHeadersIfAbsent(headers).build();
+ Message> message = builder.build();
if (logger.isDebugEnabled()) {
logger.debug("converted JMS Message [" + jmsMessage + "] to integration Message [" + message + "]");
}
@@ -164,18 +139,14 @@ public class HeaderMappingMessageConverter implements MessageConverter {
* Converts from an integration Message to a JMS {@link javax.jms.Message}.
*/
public javax.jms.Message toMessage(Object object, Session session) throws JMSException, MessageConversionException {
- MessageHeaders headers = null;
javax.jms.Message jmsMessage = null;
if (object instanceof Message>) {
- headers = ((Message>) object).getHeaders();
if (this.extractIntegrationMessagePayload) {
object = ((Message>) object).getPayload();
}
}
jmsMessage = this.converter.toMessage(object, session);
- if (headers != null) {
- this.headerMapper.fromHeaders(headers, jmsMessage);
- }
+
if (logger.isDebugEnabled()) {
logger.debug("converted [" + object + "] to JMS Message [" + jmsMessage + "]");
}
diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java
index 242795e416..3fcc23c9e9 100644
--- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java
+++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/JmsDestinationPollingSource.java
@@ -1,5 +1,5 @@
/*
- * Copyright 2002-2009 the original author or authors.
+ * Copyright 2002-2010 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,11 +16,14 @@
package org.springframework.integration.jms;
+import java.util.Map;
+
import javax.jms.ConnectionFactory;
import javax.jms.Destination;
import org.springframework.integration.core.Message;
-import org.springframework.integration.message.GenericMessage;
+import org.springframework.integration.core.MessagingException;
+import org.springframework.integration.message.MessageBuilder;
import org.springframework.integration.message.MessageSource;
import org.springframework.jms.core.JmsTemplate;
import org.springframework.jms.support.converter.MessageConverter;
@@ -32,14 +35,12 @@ import org.springframework.jms.support.converter.MessageConverter;
* support is a better option.
*
* @author Mark Fisher
+ * @author Oleg Zhurakousky
*/
public class JmsDestinationPollingSource extends AbstractJmsTemplateBasedAdapter implements MessageSource