diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java index 3dfd9ddba3..fdefae56b2 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractEndpoint.java @@ -34,6 +34,7 @@ import org.springframework.integration.message.MessageSource; import org.springframework.integration.message.MessageTarget; import org.springframework.integration.message.selector.MessageSelector; import org.springframework.integration.scheduling.Schedule; +import org.springframework.util.Assert; /** * Base class for {@link MessageEndpoint} implementations. @@ -86,9 +87,18 @@ public abstract class AbstractEndpoint implements MessageEndpoint, ChannelRegist } public void setMessageExchangeTemplate(MessageExchangeTemplate messageExchangeTemplate) { + Assert.notNull(messageExchangeTemplate, "messageExchangeTemplate must not be null"); this.messageExchangeTemplate = messageExchangeTemplate; } + public MessageExchangeTemplate getMessageExchangeTemplate() { + if (this.messageExchangeTemplate == null) { + this.messageExchangeTemplate = new MessageExchangeTemplate(); + this.messageExchangeTemplate.afterPropertiesSet(); + } + return this.messageExchangeTemplate; + } + public void setInputChannelName(String inputChannelName) { this.inputChannelName = inputChannelName; } @@ -134,7 +144,7 @@ public abstract class AbstractEndpoint implements MessageEndpoint, ChannelRegist } public void setSendTimeout(long sendTimeout) { - this.messageExchangeTemplate.setSendTimeout(sendTimeout); + this.getMessageExchangeTemplate().setSendTimeout(sendTimeout); } public MessageTarget getTarget() { @@ -181,10 +191,6 @@ public abstract class AbstractEndpoint implements MessageEndpoint, ChannelRegist } public void afterPropertiesSet() { - if (this.messageExchangeTemplate == null) { - this.messageExchangeTemplate = new MessageExchangeTemplate(); - this.messageExchangeTemplate.afterPropertiesSet(); - } if (this.target == null) { this.target = this.getTarget(); } @@ -235,15 +241,12 @@ public abstract class AbstractEndpoint implements MessageEndpoint, ChannelRegist } private boolean doSend(Message message) { - if (this.messageExchangeTemplate == null) { - this.afterPropertiesSet(); - } if (!this.supports(message)) { throw new MessageRejectedException(message, "unsupported message"); } Message result = this.handleMessage(message); if (result != null) { - return this.messageExchangeTemplate.send(message, this.target); + return this.getMessageExchangeTemplate().send(message, this.target); } return true; } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointPoller.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointPoller.java index cb01cbf823..83822471f4 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointPoller.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/EndpointPoller.java @@ -17,7 +17,6 @@ package org.springframework.integration.endpoint; import org.springframework.integration.ConfigurationException; -import org.springframework.integration.message.MessageExchangeTemplate; import org.springframework.integration.message.MessageSource; import org.springframework.integration.message.PollableSource; @@ -26,14 +25,6 @@ import org.springframework.integration.message.PollableSource; */ public class EndpointPoller implements EndpointVisitor { - private final MessageExchangeTemplate template; - - - public EndpointPoller() { - this.template = new MessageExchangeTemplate(); - this.template.setSendTimeout(0); - } - public void visitEndpoint(MessageEndpoint endpoint) { MessageSource source = endpoint.getSource(); if (source == null) { @@ -44,7 +35,7 @@ public class EndpointPoller implements EndpointVisitor { throw new ConfigurationException("unable to poll for endpoint '" + endpoint + ", source is not a PollableSource"); } - this.template.receiveAndForward((PollableSource) source, endpoint); + endpoint.getMessageExchangeTemplate().receiveAndForward((PollableSource) source, endpoint); } } diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java index 0c3911d6f6..e44f61f092 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/MessageEndpoint.java @@ -16,6 +16,7 @@ package org.springframework.integration.endpoint; +import org.springframework.integration.message.MessageExchangeTemplate; import org.springframework.integration.message.MessageSource; import org.springframework.integration.message.MessageTarget; import org.springframework.integration.scheduling.Schedule; @@ -43,4 +44,6 @@ public interface MessageEndpoint extends MessageTarget { String getOutputChannelName(); + MessageExchangeTemplate getMessageExchangeTemplate(); + }