INT-914 DelayHandler now uses Spring's TaskScheduler rather than ScheduledExecutorService.

This commit is contained in:
Mark Fisher
2009-12-10 03:48:21 +00:00
parent ac0c848aeb
commit a74090206f
5 changed files with 83 additions and 54 deletions

View File

@@ -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.
* <p>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()
* <p>
* 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();
}
}

View File

@@ -2,8 +2,11 @@
<beans:beans xmlns="http://www.springframework.org/schema/integration"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:beans="http://www.springframework.org/schema/beans"
xmlns:task="http://www.springframework.org/schema/task"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/task
http://www.springframework.org/schema/task/spring-task.xsd
http://www.springframework.org/schema/integration
http://www.springframework.org/schema/integration/spring-integration.xsd">
@@ -28,8 +31,6 @@
default-delay="0"
scheduler="testScheduler"/>
<beans:bean id="testScheduler" class="org.springframework.scheduling.concurrent.ScheduledExecutorFactoryBean">
<beans:property name="poolSize" value="7"/>
</beans:bean>
<task:scheduler id="testScheduler" pool-size="7"/>
</beans:beans>

View File

@@ -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"));
}
}

View File

@@ -2,8 +2,11 @@
<beans:beans xmlns="http://www.springframework.org/schema/integration"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:beans="http://www.springframework.org/schema/beans"
xmlns:task="http://www.springframework.org/schema/task"
xsi:schemaLocation="http://www.springframework.org/schema/beans
http://www.springframework.org/schema/beans/spring-beans.xsd
http://www.springframework.org/schema/task
http://www.springframework.org/schema/task/spring-task.xsd
http://www.springframework.org/schema/integration
http://www.springframework.org/schema/integration/spring-integration.xsd">
@@ -35,9 +38,7 @@
send-timeout="20000"
scheduler="multiThreadScheduler"/>
<beans:bean id="multiThreadScheduler" class="org.springframework.scheduling.concurrent.ScheduledExecutorFactoryBean">
<beans:property name="poolSize" value="5"/>
</beans:bean>
<task:scheduler id="multiThreadScheduler" pool-size="5"/>
<service-activator input-channel="outputB" output-channel="outputB1" method="processMessage" ref="sampleHandler"/>

View File

@@ -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);