Added setter for 'defaultConcurrencyPolicy' (INT-106).
This commit is contained in:
@@ -114,6 +114,15 @@ public class MessageBus implements ChannelRegistry, ApplicationContextAware, Lif
|
||||
this.executor = executor;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify the default concurrency policy to be used for any endpoint that
|
||||
* is registered without an explicitly provided policy of its own.
|
||||
*/
|
||||
public void setDefaultConcurrencyPolicy(ConcurrencyPolicy defaultConcurrencyPolicy) {
|
||||
Assert.notNull(defaultConcurrencyPolicy, "'defaultConcurrencyPolicy' must not be null");
|
||||
this.defaultConcurrencyPolicy = defaultConcurrencyPolicy;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set whether to automatically start the bus after initialization.
|
||||
* <p>Default is 'true'; set this to 'false' to allow for manual startup
|
||||
|
||||
@@ -36,10 +36,12 @@ import org.springframework.integration.adapter.SourceAdapter;
|
||||
import org.springframework.integration.channel.DispatcherPolicy;
|
||||
import org.springframework.integration.channel.MessageChannel;
|
||||
import org.springframework.integration.channel.SimpleChannel;
|
||||
import org.springframework.integration.endpoint.ConcurrencyPolicy;
|
||||
import org.springframework.integration.handler.MessageHandler;
|
||||
import org.springframework.integration.message.ErrorMessage;
|
||||
import org.springframework.integration.message.GenericMessage;
|
||||
import org.springframework.integration.message.Message;
|
||||
import org.springframework.integration.message.MessageDeliveryException;
|
||||
import org.springframework.integration.message.StringMessage;
|
||||
import org.springframework.integration.scheduling.Subscription;
|
||||
|
||||
@@ -178,6 +180,41 @@ public class MessageBusTests {
|
||||
bus.stop();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDefaultConcurrencyPolicy() throws InterruptedException {
|
||||
MessageBus bus = new MessageBus();
|
||||
bus.setDefaultConcurrencyPolicy(new ConcurrencyPolicy(1, 3));
|
||||
final CountDownLatch latch = new CountDownLatch(3);
|
||||
MessageHandler testHandler = new MessageHandler() {
|
||||
public Message<?> handle(Message<?> message) {
|
||||
latch.countDown();
|
||||
try {
|
||||
Thread.sleep(5000);
|
||||
}
|
||||
catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
};
|
||||
DispatcherPolicy dispatcherPolicy = new DispatcherPolicy();
|
||||
dispatcherPolicy.setRejectionLimit(1);
|
||||
dispatcherPolicy.setRetryInterval(0);
|
||||
SimpleChannel testChannel = new SimpleChannel(0, dispatcherPolicy);
|
||||
bus.registerChannel("testChannel", testChannel);
|
||||
bus.registerHandler("testHandler", testHandler, new Subscription(testChannel));
|
||||
bus.start();
|
||||
for (int i = 0; i < 4; i++) {
|
||||
assertTrue(testChannel.send(new StringMessage("test-"+ i), 100));
|
||||
}
|
||||
latch.await(1000, TimeUnit.MILLISECONDS);
|
||||
assertEquals(0, latch.getCount());
|
||||
MessageChannel errorChannel = bus.getErrorChannel();
|
||||
Message<?> errorMessage = errorChannel.receive(500);
|
||||
assertNotNull(errorMessage);
|
||||
assertEquals(MessageDeliveryException.class, errorMessage.getPayload().getClass());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testMultipleMessageBusBeans() {
|
||||
boolean exceptionThrown = false;
|
||||
|
||||
Reference in New Issue
Block a user