INT-2291 backported

This commit is contained in:
Mark Fisher
2011-12-19 10:37:51 -05:00
parent 820a1e138a
commit da21487251
2 changed files with 31 additions and 28 deletions

View File

@@ -28,9 +28,7 @@
<dispatcher load-balancer="none" failover="true" task-executor="taskExecutor"/>
</channel>
<beans:bean id="taskExecutor" class="org.springframework.core.task.SimpleAsyncTaskExecutor">
<beans:property name="threadGroupName" value="dispatchers"/>
</beans:bean>
<beans:bean id="taskExecutor" class="java.util.concurrent.Executors" factory-method="newCachedThreadPool"/>
<!-- send async errors to null channel to avoid chatter -->
<bridge input-channel="errorChannel" output-channel="nullChannel"/>

View File

@@ -19,7 +19,6 @@ package org.springframework.integration.channel;
import java.util.List;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
@@ -61,7 +60,7 @@ public class MixedDispatcherConfigurationScenarioTests {
private static final int TOTAL_EXECUTIONS = 40;
private ExecutorService scheduler;
private ExecutorService executor;
private CountDownLatch allDone;
private CountDownLatch start;
@@ -91,9 +90,10 @@ public class MixedDispatcherConfigurationScenarioTests {
Mockito.reset(handlerA);
Mockito.reset(handlerB);
Mockito.reset(handlerC);
scheduler = Executors.newCachedThreadPool();
ac = new ClassPathXmlApplicationContext("MixedDispatcherConfigurationScenarioTests-context.xml",
MixedDispatcherConfigurationScenarioTests.class);
executor = ac.getBean("taskExecutor", ExecutorService.class);
allDone = new CountDownLatch(TOTAL_EXECUTIONS);
start = new CountDownLatch(1);
failed = new AtomicBoolean(false);
@@ -146,13 +146,13 @@ public class MixedDispatcherConfigurationScenarioTests {
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
assertTrue("not all messages were accepted", failed.get());
verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message);
@@ -196,13 +196,13 @@ public class MixedDispatcherConfigurationScenarioTests {
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
assertTrue("not all messages were accepted", failed.get());
verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message);
@@ -274,13 +274,13 @@ public class MixedDispatcherConfigurationScenarioTests {
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
assertTrue("not all messages were accepted", failed.get());
verify(handlerA, times(14)).handleMessage(message);
@@ -334,13 +334,13 @@ public class MixedDispatcherConfigurationScenarioTests {
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
assertTrue("not all messages were accepted", failed.get());
verify(handlerA, times(14)).handleMessage(message);
@@ -413,13 +413,13 @@ public class MixedDispatcherConfigurationScenarioTests {
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
assertFalse("not all messages were accepted", failed.get());
verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message);
@@ -428,14 +428,15 @@ public class MixedDispatcherConfigurationScenarioTests {
verify(exceptionRegistry, never()).add((Exception) anyObject());
}
@Test(timeout = 5000)
@Test(timeout = 5000)
public void failoverNoLoadBalancingWithExecutorConcurrent() throws Exception {
final ExecutorChannel channel = (ExecutorChannel) ac.getBean("noLoadBalancerFailoverExecutor");
final UnicastingDispatcher dispatcher = channel.getDispatcher();
dispatcher.addHandler(handlerA);
dispatcher.addHandler(handlerB);
dispatcher.addHandler(handlerC);
dispatcher.addHandler(handlerC);
final CountDownLatch taskSubmissionMonitor = new CountDownLatch(TOTAL_EXECUTIONS);
doAnswer(new Answer<Object>() {
public Object answer(InvocationOnMock invocation) {
@@ -461,21 +462,25 @@ public class MixedDispatcherConfigurationScenarioTests {
Runnable messageSenderTask = new Runnable() {
public void run() {
try {
start.await();
start.await();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
channel.send(message);
taskSubmissionMonitor.countDown();
}
};
for (int i = 0; i < TOTAL_EXECUTIONS; i++) {
scheduler.execute(messageSenderTask);
executor.execute(messageSenderTask);
}
start.countDown();
allDone.await();
scheduler.shutdown();
scheduler.awaitTermination(5, TimeUnit.SECONDS);
taskSubmissionMonitor.await();
executor.shutdown();
executor.awaitTermination(5, TimeUnit.SECONDS);
verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message);
verify(handlerB, times(TOTAL_EXECUTIONS)).handleMessage(message);