diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/AbstractEndpointParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/AbstractEndpointParser.java index d190327066..84e4fb1245 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/config/AbstractEndpointParser.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/AbstractEndpointParser.java @@ -94,6 +94,7 @@ public abstract class AbstractEndpointParser extends AbstractSingleBeanDefinitio if (txElement != null) { IntegrationNamespaceUtils.configureTransactionAttributes(txElement, builder); } + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, pollerElement, "task-executor"); } IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, INPUT_CHANNEL_ATTRIBUTE); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, OUTPUT_CHANNEL_ATTRIBUTE); diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractMessageConsumingEndpoint.java b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractMessageConsumingEndpoint.java index 61f58064a3..523d187a4d 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractMessageConsumingEndpoint.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/endpoint/AbstractMessageConsumingEndpoint.java @@ -17,6 +17,7 @@ package org.springframework.integration.endpoint; import org.springframework.context.Lifecycle; +import org.springframework.core.task.TaskExecutor; import org.springframework.integration.ConfigurationException; import org.springframework.integration.channel.MessageChannel; import org.springframework.integration.channel.PollableChannel; @@ -41,6 +42,8 @@ public abstract class AbstractMessageConsumingEndpoint extends AbstractEndpoint private volatile ChannelPoller poller; + private volatile TaskExecutor taskExecutor; + private volatile int maxMessagesPerPoll = -1; private volatile boolean initialized; @@ -58,6 +61,10 @@ public abstract class AbstractMessageConsumingEndpoint extends AbstractEndpoint this.schedule = schedule; } + public void setTaskExecutor(TaskExecutor taskExecutor) { + this.taskExecutor = taskExecutor; + } + public void setMaxMessagesPerPoll(int maxMessagesPerPoll) { this.maxMessagesPerPoll = maxMessagesPerPoll; if (this.poller != null) { @@ -76,6 +83,9 @@ public abstract class AbstractMessageConsumingEndpoint extends AbstractEndpoint this.poller = new ChannelPoller((PollableChannel) this.inputChannel, this.schedule); this.poller.setMaxMessagesPerPoll(this.maxMessagesPerPoll); this.configureTransactionSettingsForPoller(this.poller); + if (this.taskExecutor != null) { + this.poller.setTaskExecutor(this.taskExecutor); + } this.poller.subscribe(this); } this.initialized = true;