diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/handler/DelayHandler.java b/org.springframework.integration/src/main/java/org/springframework/integration/handler/DelayHandler.java
index 0664c64484..59acb9d7dd 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/handler/DelayHandler.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/handler/DelayHandler.java
@@ -17,9 +17,6 @@
package org.springframework.integration.handler;
import java.util.Date;
-import java.util.concurrent.Executors;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -28,6 +25,7 @@ import org.springframework.beans.BeansException;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.beans.factory.DisposableBean;
+import org.springframework.beans.factory.InitializingBean;
import org.springframework.core.Ordered;
import org.springframework.integration.channel.BeanFactoryChannelResolver;
import org.springframework.integration.channel.ChannelResolutionException;
@@ -40,6 +38,9 @@ import org.springframework.integration.core.MessageHeaders;
import org.springframework.integration.message.ErrorMessage;
import org.springframework.integration.message.MessageHandler;
import org.springframework.integration.message.MessageHandlingException;
+import org.springframework.scheduling.TaskScheduler;
+import org.springframework.scheduling.concurrent.ExecutorConfigurationSupport;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
import org.springframework.util.Assert;
/**
@@ -69,7 +70,7 @@ import org.springframework.util.Assert;
* @author Mark Fisher
* @since 1.0.3
*/
-public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware, DisposableBean {
+public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware, InitializingBean, DisposableBean {
private final Log logger = LogFactory.getLog(this.getClass());
@@ -77,7 +78,7 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
private volatile String delayHeaderName;
- private final ScheduledExecutorService scheduler;
+ private final TaskScheduler taskScheduler;
private volatile MessageChannel outputChannel;
@@ -85,8 +86,6 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
private final MessageChannelTemplate channelTemplate = new MessageChannelTemplate();
- private volatile boolean waitForTasksToCompleteOnShutdown;
-
private volatile int order = Ordered.LOWEST_PRECEDENCE;
@@ -99,13 +98,13 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
}
/**
- * Create a DelayHandler with the given default delay. The sending of Messages after
- * the delay will be handled by the provided {@link ScheduledExecutorService}.
+ * Create a DelayHandler with the given default delay. The sending of Messages
+ * after the delay will be handled by the provided {@link TaskScheduler}.
*/
- public DelayHandler(long defaultDelay, ScheduledExecutorService scheduledExecutorService) {
+ public DelayHandler(long defaultDelay, TaskScheduler taskScheduler) {
this.defaultDelay = defaultDelay;
- this.scheduler = (scheduledExecutorService != null)
- ? scheduledExecutorService : Executors.newScheduledThreadPool(1);
+ this.taskScheduler = (taskScheduler != null)
+ ? taskScheduler : new ThreadPoolTaskScheduler();
}
@@ -147,11 +146,19 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
* Set whether to wait for scheduled tasks to complete on shutdown.
*
Default is "false". Switch this to "true" if you prefer
* fully completed tasks at the expense of a longer shutdown phase.
- * @see java.util.concurrent.ExecutorService#shutdown()
- * @see java.util.concurrent.ExecutorService#shutdownNow()
+ *
+ * This property will only have an effect for TaskScheduler implementations
+ * that extend from {@link ExecutorConfigurationSupport}.
+ * @see ExecutorConfigurationSupport#setWaitForTasksToCompleteOnShutdown(boolean)
*/
public void setWaitForTasksToCompleteOnShutdown(boolean waitForJobsToCompleteOnShutdown) {
- this.waitForTasksToCompleteOnShutdown = waitForJobsToCompleteOnShutdown;
+ if (this.taskScheduler instanceof ExecutorConfigurationSupport) {
+ ((ExecutorConfigurationSupport) this.taskScheduler).setWaitForTasksToCompleteOnShutdown(waitForJobsToCompleteOnShutdown);
+ }
+ else if (logger.isWarnEnabled()) {
+ logger.warn("The 'waitForJobsToCompleteOnShutdown' property is not supported for TaskScheduler of type [" +
+ this.taskScheduler.getClass() + "]");
+ }
}
public void setOrder(int order) {
@@ -166,6 +173,12 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
this.channelResolver = new BeanFactoryChannelResolver(beanFactory);
}
+ public void afterPropertiesSet() throws Exception {
+ if (this.taskScheduler instanceof InitializingBean) {
+ ((InitializingBean) this.taskScheduler).afterPropertiesSet();
+ }
+ }
+
public final void handleMessage(final Message> message) {
long delay = this.determineDelayForMessage(message);
if (delay > 0) {
@@ -200,7 +213,7 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
}
private void releaseMessageAfterDelay(final Message> message, long delay) {
- this.scheduler.schedule(new Runnable() {
+ this.taskScheduler.schedule(new Runnable() {
public void run() {
try {
releaseMessage(message);
@@ -220,7 +233,7 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
}
}
}
- }, delay, TimeUnit.MILLISECONDS);
+ }, new Date(System.currentTimeMillis() + delay));
}
private void releaseMessage(Message> message) {
@@ -276,12 +289,9 @@ public class DelayHandler implements MessageHandler, Ordered, BeanFactoryAware,
return channel;
}
- public void destroy() {
- if (this.waitForTasksToCompleteOnShutdown) {
- this.scheduler.shutdown();
- }
- else {
- this.scheduler.shutdownNow();
+ public void destroy() throws Exception {
+ if (this.taskScheduler instanceof DisposableBean) {
+ ((DisposableBean) this.taskScheduler).destroy();
}
}
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml
index c8169d81ef..3ca64320b5 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests-context.xml
@@ -2,8 +2,11 @@
@@ -28,8 +31,6 @@
default-delay="0"
scheduler="testScheduler"/>
-
-
-
+
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java
index 83c41fb4bd..8c44240615 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerParserTests.java
@@ -56,7 +56,8 @@ public class DelayerParserTests {
assertEquals("foo", accessor.getPropertyValue("delayHeaderName"));
assertEquals(new Long(987), new DirectFieldAccessor(
accessor.getPropertyValue("channelTemplate")).getPropertyValue("sendTimeout"));
- assertEquals(Boolean.TRUE, accessor.getPropertyValue("waitForTasksToCompleteOnShutdown"));
+ assertEquals(Boolean.TRUE, new DirectFieldAccessor(
+ accessor.getPropertyValue("taskScheduler")).getPropertyValue("waitForTasksToCompleteOnShutdown"));
}
@@ -71,7 +72,7 @@ public class DelayerParserTests {
DirectFieldAccessor accessor = new DirectFieldAccessor(delayHandler);
assertEquals(context.getBean("output"), accessor.getPropertyValue("outputChannel"));
assertEquals(new Long(0), accessor.getPropertyValue("defaultDelay"));
- assertEquals(context.getBean("testScheduler"), accessor.getPropertyValue("scheduler"));
+ assertEquals(context.getBean("testScheduler"), accessor.getPropertyValue("taskScheduler"));
}
}
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml
index e475cdf8c5..5e911655f4 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/xml/DelayerUsageTests-context.xml
@@ -2,8 +2,11 @@
@@ -35,9 +38,7 @@
send-timeout="20000"
scheduler="multiThreadScheduler"/>
-
-
-
+
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java
index 2eca3eb0b1..53ad6e2a25 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/handler/DelayHandlerTests.java
@@ -22,7 +22,6 @@ import static org.junit.Assert.assertSame;
import java.util.Date;
import java.util.concurrent.CountDownLatch;
-import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.junit.Before;
@@ -38,6 +37,7 @@ import org.springframework.integration.message.MessageDeliveryException;
import org.springframework.integration.message.MessageHandler;
import org.springframework.integration.message.MessageHandlingException;
import org.springframework.integration.message.StringMessage;
+import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
/**
* @author Mark Fisher
@@ -60,10 +60,11 @@ public class DelayHandlerTests {
@Test
- public void noDelayHeaderAndDefaultDelayIsZero() {
+ public void noDelayHeaderAndDefaultDelayIsZero() throws Exception {
DelayHandler delayHandler = new DelayHandler(0);
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test").build();
@@ -73,10 +74,11 @@ public class DelayHandlerTests {
}
@Test
- public void noDelayHeaderAndDefaultDelayIsPositive() {
+ public void noDelayHeaderAndDefaultDelayIsPositive() throws Exception {
DelayHandler delayHandler = new DelayHandler(10);
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test").build();
@@ -87,11 +89,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderAndDefaultDelayWouldTimeout() {
+ public void delayHeaderAndDefaultDelayWouldTimeout() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -103,11 +106,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsNegativeAndDefaultDelayWouldTimeout() {
+ public void delayHeaderIsNegativeAndDefaultDelayWouldTimeout() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -119,11 +123,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsInvalidFallsBackToDefaultDelay() {
+ public void delayHeaderIsInvalidFallsBackToDefaultDelay() throws Exception {
DelayHandler delayHandler = new DelayHandler(5);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -135,11 +140,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsDateInTheFutureAndDefaultDelayWouldTimeout() {
+ public void delayHeaderIsDateInTheFutureAndDefaultDelayWouldTimeout() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -151,11 +157,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsDateInThePastAndDefaultDelayWouldTimeout() {
+ public void delayHeaderIsDateInThePastAndDefaultDelayWouldTimeout() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -167,11 +174,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsNullDateAndDefaultDelayIsZero() {
+ public void delayHeaderIsNullDateAndDefaultDelayIsZero() throws Exception {
DelayHandler delayHandler = new DelayHandler(0);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Date nullDate = null;
@@ -184,11 +192,12 @@ public class DelayHandlerTests {
}
@Test(expected = TestTimedOutException.class)
- public void delayHeaderIsFutureDateAndTimesOut() {
+ public void delayHeaderIsFutureDateAndTimesOut() throws Exception {
DelayHandler delayHandler = new DelayHandler(0);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Date future = new Date(new Date().getTime() + 60 * 1000);
@@ -201,11 +210,12 @@ public class DelayHandlerTests {
}
@Test
- public void delayHeaderIsValidStringAndDefaultDelayWouldTimeout() {
+ public void delayHeaderIsValidStringAndDefaultDelayWouldTimeout() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setDelayHeaderName("delay");
ResultHandler resultHandler = new ResultHandler();
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(resultHandler);
Message> message = MessageBuilder.withPayload("test")
@@ -217,17 +227,18 @@ public class DelayHandlerTests {
}
@Test
- public void verifyShutdownWithoutWaitingByDefault() throws InterruptedException {
+ public void verifyShutdownWithoutWaitingByDefault() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
+ delayHandler.afterPropertiesSet();
delayHandler.handleMessage(new StringMessage("foo"));
delayHandler.destroy();
- final ScheduledThreadPoolExecutor scheduler = (ScheduledThreadPoolExecutor)
- new DirectFieldAccessor(delayHandler).getPropertyValue("scheduler");
+ final ThreadPoolTaskScheduler taskScheduler = (ThreadPoolTaskScheduler)
+ new DirectFieldAccessor(delayHandler).getPropertyValue("taskScheduler");
final CountDownLatch latch = new CountDownLatch(1);
new Thread(new Runnable() {
public void run() {
try {
- scheduler.awaitTermination(10000, TimeUnit.MILLISECONDS);
+ taskScheduler.getScheduledExecutor().awaitTermination(10000, TimeUnit.MILLISECONDS);
latch.countDown();
}
catch (InterruptedException e) {
@@ -240,18 +251,19 @@ public class DelayHandlerTests {
}
@Test
- public void verifyShutdownWithWait() throws InterruptedException {
+ public void verifyShutdownWithWait() throws Exception {
DelayHandler delayHandler = new DelayHandler(5000);
delayHandler.setWaitForTasksToCompleteOnShutdown(true);
+ delayHandler.afterPropertiesSet();
delayHandler.handleMessage(new StringMessage("foo"));
delayHandler.destroy();
- final ScheduledThreadPoolExecutor scheduler = (ScheduledThreadPoolExecutor)
- new DirectFieldAccessor(delayHandler).getPropertyValue("scheduler");
+ final ThreadPoolTaskScheduler taskScheduler = (ThreadPoolTaskScheduler)
+ new DirectFieldAccessor(delayHandler).getPropertyValue("taskScheduler");
final CountDownLatch latch = new CountDownLatch(1);
new Thread(new Runnable() {
public void run() {
try {
- scheduler.awaitTermination(10000, TimeUnit.MILLISECONDS);
+ taskScheduler.getScheduledExecutor().awaitTermination(10000, TimeUnit.MILLISECONDS);
latch.countDown();
}
catch (InterruptedException e) {
@@ -264,9 +276,10 @@ public class DelayHandlerTests {
}
@Test(expected = MessageDeliveryException.class)
- public void handlerThrowsExceptionWithNoDelay() {
+ public void handlerThrowsExceptionWithNoDelay() throws Exception {
DelayHandler delayHandler = new DelayHandler(0);
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
input.subscribe(delayHandler);
output.subscribe(new MessageHandler() {
public void handleMessage(Message> message) {
@@ -278,10 +291,11 @@ public class DelayHandlerTests {
}
@Test
- public void errorChannelHeaderAndHandlerThrowsExceptionWithDelay() {
+ public void errorChannelHeaderAndHandlerThrowsExceptionWithDelay() throws Exception {
DelayHandler delayHandler = new DelayHandler(0);
delayHandler.setDelayHeaderName("delay");
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
DirectChannel errorChannel = new DirectChannel();
ResultHandler resultHandler = new ResultHandler();
errorChannel.subscribe(resultHandler);
@@ -308,7 +322,7 @@ public class DelayHandlerTests {
}
@Test
- public void errorChannelNameHeaderAndHandlerThrowsExceptionWithDelay() {
+ public void errorChannelNameHeaderAndHandlerThrowsExceptionWithDelay() throws Exception {
String errorChannelName = "customErrorChannel";
StaticApplicationContext context = new StaticApplicationContext();
context.registerSingleton(errorChannelName, DirectChannel.class);
@@ -318,6 +332,7 @@ public class DelayHandlerTests {
delayHandler.setBeanFactory(context);
delayHandler.setDelayHeaderName("delay");
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
ResultHandler resultHandler = new ResultHandler();
customErrorChannel.subscribe(resultHandler);
input.subscribe(delayHandler);
@@ -343,7 +358,7 @@ public class DelayHandlerTests {
}
@Test
- public void defaultErrorChannelAndHandlerThrowsExceptionWithDelay() {
+ public void defaultErrorChannelAndHandlerThrowsExceptionWithDelay() throws Exception {
StaticApplicationContext context = new StaticApplicationContext();
context.registerSingleton(
IntegrationContextUtils.ERROR_CHANNEL_BEAN_NAME, DirectChannel.class);
@@ -354,6 +369,7 @@ public class DelayHandlerTests {
delayHandler.setBeanFactory(context);
delayHandler.setDelayHeaderName("delay");
delayHandler.setOutputChannel(output);
+ delayHandler.afterPropertiesSet();
ResultHandler resultHandler = new ResultHandler();
defaultErrorChannel.subscribe(resultHandler);
input.subscribe(delayHandler);