From a74090206f9ef8232042bce36972e77de711d16c Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Thu, 10 Dec 2009 03:48:21 +0000 Subject: [PATCH] INT-914 DelayHandler now uses Spring's TaskScheduler rather than ScheduledExecutorService. --- .../integration/handler/DelayHandler.java | 56 ++++++++++------- .../config/xml/DelayerParserTests-context.xml | 7 ++- .../config/xml/DelayerParserTests.java | 5 +- .../config/xml/DelayerUsageTests-context.xml | 7 ++- .../handler/DelayHandlerTests.java | 62 ++++++++++++------- 5 files changed, 83 insertions(+), 54 deletions(-) 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);