INT-1242 - refactoring to decouple Message conversion from Message headers mapping
This commit is contained in:
@@ -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.
|
||||
*
|
||||
* <p>The default value is <code>true</code>. To force creation of Spring
|
||||
* Integration Messages whose payload is the actual JMS Message, set this
|
||||
* to <code>false</code>.
|
||||
*/
|
||||
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.
|
||||
* <p>
|
||||
* 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
* <p>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());
|
||||
}
|
||||
|
||||
@@ -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}.
|
||||
*
|
||||
* <p>If 'extractJmsMessageBody' is <code>true</code> (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<String, Object> 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 + "]");
|
||||
}
|
||||
@@ -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<Object> {
|
||||
|
||||
private volatile boolean extractPayload = true;
|
||||
|
||||
private volatile String messageSelector;
|
||||
|
||||
|
||||
public JmsDestinationPollingSource(JmsTemplate jmsTemplate) {
|
||||
super(jmsTemplate);
|
||||
}
|
||||
@@ -59,39 +60,35 @@ public class JmsDestinationPollingSource extends AbstractJmsTemplateBasedAdapter
|
||||
public void setMessageSelector(String messageSelector) {
|
||||
this.messageSelector = messageSelector;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether the payload should be extracted from each received JMS
|
||||
* Message to be used as the Spring Integration Message payload.
|
||||
* Will receive JMS {@link javax.jms.Message} converting and returning it as
|
||||
* Spring Integration(SI) {@link Message}.
|
||||
* This method will also use the current instance of the {@link JmsHeaderMapper} to map
|
||||
* JMS headers to SI headers
|
||||
*
|
||||
* <p>The default value is <code>true</code>. To force creation of Spring
|
||||
* Integration Messages whose payload is the actual JMS Message, set this
|
||||
* to <code>false</code>.
|
||||
* @return
|
||||
*/
|
||||
public void setExtractPayload(boolean extractPayload) {
|
||||
this.extractPayload = extractPayload;
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
public Message<Object> receive() {
|
||||
Object receivedObject = this.getJmsTemplate().receiveSelectedAndConvert(this.messageSelector);
|
||||
if (receivedObject == null) {
|
||||
Message<Object> convertedMessage = null;
|
||||
// receive JMS Message
|
||||
javax.jms.Message jmsMessage = this.getJmsTemplate().receiveSelected(this.messageSelector);
|
||||
if (jmsMessage == null) {
|
||||
return null;
|
||||
}
|
||||
if (receivedObject instanceof Message) {
|
||||
return (Message) receivedObject;
|
||||
try {
|
||||
// Map headers
|
||||
Map<String, Object> mappedHeaders = this.getHeaderMapper().toHeaders(jmsMessage);
|
||||
MessageConverter converter = this.getJmsTemplate().getMessageConverter();
|
||||
Object convertedObject = converter.fromMessage(jmsMessage);
|
||||
if (convertedObject instanceof Message) {
|
||||
convertedMessage = MessageBuilder.fromMessage((Message<Object>) convertedObject).copyHeaders(mappedHeaders).build();
|
||||
} else {
|
||||
convertedMessage = MessageBuilder.withPayload(convertedObject).build();
|
||||
}
|
||||
} catch (Exception e) {
|
||||
throw new MessagingException(e.getMessage(), e);
|
||||
}
|
||||
return new GenericMessage<Object>(receivedObject);
|
||||
return convertedMessage;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void configureMessageConverter(JmsTemplate jmsTemplate, JmsHeaderMapper headerMapper) {
|
||||
MessageConverter converter = jmsTemplate.getMessageConverter();
|
||||
if (converter == null || !(converter instanceof HeaderMappingMessageConverter)) {
|
||||
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(converter, headerMapper);
|
||||
hmmc.setExtractJmsMessageBody(this.extractPayload);
|
||||
jmsTemplate.setMessageConverter(hmmc);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2008 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,9 +16,9 @@
|
||||
|
||||
package org.springframework.integration.jms;
|
||||
|
||||
import java.util.Map;
|
||||
import javax.jms.Message;
|
||||
|
||||
import org.springframework.integration.core.MessageHeaders;
|
||||
import org.springframework.integration.message.HeaderMapper;
|
||||
|
||||
/**
|
||||
* Strategy interface for mapping integration Message headers to an outbound
|
||||
@@ -26,11 +26,7 @@ import org.springframework.integration.core.MessageHeaders;
|
||||
* header values from an inbound JMS Message.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public interface JmsHeaderMapper {
|
||||
public interface JmsHeaderMapper extends HeaderMapper<Message> {}
|
||||
|
||||
void fromHeaders(MessageHeaders headers, javax.jms.Message jmsMessage);
|
||||
|
||||
Map<String, Object> toHeaders(javax.jms.Message jmsMessage);
|
||||
|
||||
}
|
||||
|
||||
@@ -49,6 +49,7 @@ import org.springframework.util.Assert;
|
||||
* @author Mark Fisher
|
||||
* @author Arjen Poutsma
|
||||
* @author Juergen Hoeller
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
|
||||
@@ -78,7 +79,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
|
||||
private volatile MessageConverter messageConverter;
|
||||
|
||||
private volatile JmsHeaderMapper headerMapper;
|
||||
private volatile JmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
|
||||
|
||||
private volatile boolean extractRequestPayload = true;
|
||||
|
||||
@@ -196,7 +197,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
* Spring Integration request Message into a JMS Message and for converting
|
||||
* the JMS reply Messages back into Spring Integration Messages.
|
||||
* <p>
|
||||
* 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) {
|
||||
@@ -211,7 +212,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
* <p>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;
|
||||
@@ -220,11 +221,11 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
/**
|
||||
* This property will take effect if no custom {@link MessageConverter}
|
||||
* has been provided to the {@link #setMessageConverter(MessageConverter)}
|
||||
* method. In that case, a {@link HeaderMappingMessageConverter} will be
|
||||
* method. In that case, a {@link DefaultMessageConverter} will be
|
||||
* used by default, and this value will be passed along to that converter's
|
||||
* 'extractIntegrationMessagePayload' property.
|
||||
*
|
||||
* @see HeaderMappingMessageConverter#setExtractIntegrationMessagePayload(boolean)
|
||||
* @see DefaultMessageConverter#setExtractIntegrationMessagePayload(boolean)
|
||||
*/
|
||||
public void setExtractRequestPayload(boolean extractRequestPayload) {
|
||||
this.extractRequestPayload = extractRequestPayload;
|
||||
@@ -233,11 +234,11 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
/**
|
||||
* This property will take effect if no custom {@link MessageConverter}
|
||||
* has been provided to the {@link #setMessageConverter(MessageConverter)}
|
||||
* method. In that case, a {@link HeaderMappingMessageConverter} will be
|
||||
* method. In that case, a {@link DefaultMessageConverter} will be
|
||||
* used by default, and this value will be passed along to that converter's
|
||||
* 'extractJmsMessageBody' property.
|
||||
*
|
||||
* @see HeaderMappingMessageConverter#setExtractJmsMessageBody(boolean)
|
||||
* @see DefaultMessageConverter#setExtractJmsMessageBody(boolean)
|
||||
*/
|
||||
public void setExtractReplyPayload(boolean extractReplyPayload) {
|
||||
this.extractReplyPayload = extractReplyPayload;
|
||||
@@ -284,7 +285,7 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
Assert.isTrue(this.requestDestination != null || this.requestDestinationName != null,
|
||||
"Either a 'requestDestination' or 'requestDestinationName' is required.");
|
||||
if (this.messageConverter == null) {
|
||||
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(null, this.headerMapper);
|
||||
DefaultMessageConverter hmmc = new DefaultMessageConverter(null);
|
||||
hmmc.setExtractIntegrationMessagePayload(this.extractRequestPayload);
|
||||
hmmc.setExtractJmsMessageBody(this.extractReplyPayload);
|
||||
this.messageConverter = hmmc;
|
||||
@@ -320,7 +321,11 @@ public class JmsOutboundGateway extends AbstractReplyProducingMessageHandler {
|
||||
Destination replyTo = null;
|
||||
try {
|
||||
session = createSession(connection);
|
||||
// convert to JMS Message
|
||||
javax.jms.Message jmsRequest = this.messageConverter.toMessage(requestMessage, session);
|
||||
// map headers
|
||||
headerMapper.fromHeaders(requestMessage.getHeaders(), jmsRequest);
|
||||
// create JMS Producer and send
|
||||
messageProducer = session.createProducer(this.getRequestDestination(session));
|
||||
replyTo = this.getReplyDestination(session);
|
||||
jmsRequest.setJMSReplyTo(replyTo);
|
||||
|
||||
@@ -20,18 +20,16 @@ import org.springframework.core.Ordered;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.jms.core.JmsTemplate;
|
||||
import org.springframework.jms.support.converter.MessageConverter;
|
||||
|
||||
/**
|
||||
* A MessageConsumer that sends the converted Message payload within
|
||||
* a JMS Message.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
public class JmsSendingMessageHandler extends AbstractJmsTemplateBasedAdapter implements MessageHandler, Ordered {
|
||||
|
||||
private volatile boolean extractPayload = true;
|
||||
|
||||
private volatile int order = Ordered.LOWEST_PRECEDENCE;
|
||||
|
||||
|
||||
@@ -47,18 +45,6 @@ public class JmsSendingMessageHandler extends AbstractJmsTemplateBasedAdapter im
|
||||
super();
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify whether the payload should be extracted from each Spring
|
||||
* Integration Message to be converted to the body of a JMS Message.
|
||||
*
|
||||
* <p>The default value is <code>true</code>. To force creation of JMS
|
||||
* Messages whose body is the actual Spring Integration Message instance,
|
||||
* set this to <code>false</code>.
|
||||
*/
|
||||
public void setExtractPayload(boolean extractPayload) {
|
||||
this.extractPayload = extractPayload;
|
||||
}
|
||||
|
||||
public void setOrder(int order) {
|
||||
this.order = order;
|
||||
}
|
||||
@@ -73,15 +59,4 @@ public class JmsSendingMessageHandler extends AbstractJmsTemplateBasedAdapter im
|
||||
}
|
||||
this.getJmsTemplate().convertAndSend(message);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void configureMessageConverter(JmsTemplate jmsTemplate, JmsHeaderMapper headerMapper) {
|
||||
MessageConverter converter = jmsTemplate.getMessageConverter();
|
||||
if (converter == null || !(converter instanceof HeaderMappingMessageConverter)) {
|
||||
HeaderMappingMessageConverter hmmc = new HeaderMappingMessageConverter(converter, headerMapper);
|
||||
hmmc.setExtractIntegrationMessagePayload(this.extractPayload);
|
||||
jmsTemplate.setMessageConverter(hmmc);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -61,7 +61,7 @@ public class ChannelPublishingJmsMessageListenerTests {
|
||||
ChannelPublishingJmsMessageListener listener = new ChannelPublishingJmsMessageListener();
|
||||
listener.afterPropertiesSet();
|
||||
Object converter = new DirectFieldAccessor(listener).getPropertyValue("messageConverter");
|
||||
assertEquals(HeaderMappingMessageConverter.class, converter.getClass());
|
||||
assertEquals(DefaultMessageConverter.class, converter.getClass());
|
||||
Object wrappedConverter = new DirectFieldAccessor(converter).getPropertyValue("converter");
|
||||
assertEquals(SimpleMessageConverter.class, wrappedConverter.getClass());
|
||||
}
|
||||
@@ -73,7 +73,7 @@ public class ChannelPublishingJmsMessageListenerTests {
|
||||
listener.setMessageConverter(originalConverter);
|
||||
listener.afterPropertiesSet();
|
||||
Object converter = new DirectFieldAccessor(listener).getPropertyValue("messageConverter");
|
||||
assertEquals(HeaderMappingMessageConverter.class, converter.getClass());
|
||||
assertEquals(DefaultMessageConverter.class, converter.getClass());
|
||||
Object wrappedConverter = new DirectFieldAccessor(converter).getPropertyValue("converter");
|
||||
assertEquals(originalConverter, wrappedConverter);
|
||||
}
|
||||
|
||||
@@ -22,20 +22,20 @@ import javax.jms.Destination;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.channel.PollableChannel;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.jms.HeaderMappingMessageConverter;
|
||||
import org.springframework.integration.jms.DefaultJmsHeaderMapper;
|
||||
import org.springframework.integration.jms.JmsHeaders;
|
||||
import org.springframework.integration.jms.StubSession;
|
||||
import org.springframework.integration.jms.StubTextMessage;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
*/
|
||||
@ContextConfiguration
|
||||
@RunWith(SpringJUnit4ClassRunner.class)
|
||||
@@ -58,8 +58,9 @@ public class JmsHeaderEnricherTests {
|
||||
valueTestInput.send(new StringMessage("test"));
|
||||
Message<?> result = output.receive(0);
|
||||
assertEquals(testDestination, result.getHeaders().get(JmsHeaders.REPLY_TO));
|
||||
HeaderMappingMessageConverter converter = new HeaderMappingMessageConverter();
|
||||
javax.jms.Message jmsMessage = converter.toMessage(result, new StubSession("foo"));
|
||||
StubTextMessage jmsMessage = new StubTextMessage();
|
||||
DefaultJmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
|
||||
headerMapper.fromHeaders(result.getHeaders(), jmsMessage);
|
||||
assertEquals(testDestination, jmsMessage.getJMSReplyTo());
|
||||
}
|
||||
|
||||
@@ -68,8 +69,9 @@ public class JmsHeaderEnricherTests {
|
||||
valueTestInput.send(new StringMessage("test"));
|
||||
Message<?> result = output.receive(0);
|
||||
assertEquals("ABC", result.getHeaders().get(JmsHeaders.CORRELATION_ID));
|
||||
HeaderMappingMessageConverter converter = new HeaderMappingMessageConverter();
|
||||
javax.jms.Message jmsMessage = converter.toMessage(result, new StubSession("foo"));
|
||||
StubTextMessage jmsMessage = new StubTextMessage();
|
||||
DefaultJmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
|
||||
headerMapper.fromHeaders(result.getHeaders(), jmsMessage);
|
||||
assertEquals("ABC", jmsMessage.getJMSCorrelationID());
|
||||
}
|
||||
|
||||
@@ -78,8 +80,9 @@ public class JmsHeaderEnricherTests {
|
||||
expressionTestInput.send(new StringMessage("test"));
|
||||
Message<?> result = output.receive(0);
|
||||
assertEquals(123, result.getHeaders().get(JmsHeaders.CORRELATION_ID));
|
||||
HeaderMappingMessageConverter converter = new HeaderMappingMessageConverter();
|
||||
javax.jms.Message jmsMessage = converter.toMessage(result, new StubSession("foo"));
|
||||
StubTextMessage jmsMessage = new StubTextMessage();
|
||||
DefaultJmsHeaderMapper headerMapper = new DefaultJmsHeaderMapper();
|
||||
headerMapper.fromHeaders(result.getHeaders(), jmsMessage);
|
||||
assertEquals("123", jmsMessage.getJMSCorrelationID());
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user