diff --git a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java index f6a446d126..79433a52a3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/bus/MessageBus.java @@ -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. *

Default is 'true'; set this to 'false' to allow for manual startup diff --git a/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java b/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java index 23af0b77be..7be5bbade8 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/bus/MessageBusTests.java @@ -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;