From 9fa4a6b80e72c1168a6f35f71143432866b84efa Mon Sep 17 00:00:00 2001 From: Mark Fisher Date: Tue, 20 May 2008 20:32:04 +0000 Subject: [PATCH] Added a "configure-async-event-multicaster" attribute to the element and a corresponding 'configureAsynEventMulticaster' property to the MessageBus class. When set to 'true', the 'taskScheduler' of the MessageBus will also be configured as the TaskExecutor for the ApplicationContext's ApplicationEventMulticaster. The default value for this property is 'false' (INT-58). --- .../integration/bus/MessageBus.java | 34 +++++++++++++++- .../integration/config/MessageBusParser.java | 2 +- .../config/spring-integration-core-1.0.xsd | 1 + .../config/MessageBusParserTests.java | 40 +++++++++++++++++++ .../messageBusWithAsyncEventMulticaster.xml | 12 ++++++ ...messageBusWithoutAsyncEventMulticaster.xml | 12 ++++++ 6 files changed, 99 insertions(+), 2 deletions(-) create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithAsyncEventMulticaster.xml create mode 100644 spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithoutAsyncEventMulticaster.xml 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 10c5e99424..14b4ef6661 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 @@ -27,13 +27,17 @@ import java.util.concurrent.ScheduledThreadPoolExecutor; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; + import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.ApplicationEvent; import org.springframework.context.ApplicationListener; import org.springframework.context.Lifecycle; +import org.springframework.context.event.ApplicationEventMulticaster; import org.springframework.context.event.ContextRefreshedEvent; +import org.springframework.context.event.SimpleApplicationEventMulticaster; +import org.springframework.context.support.AbstractApplicationContext; import org.springframework.integration.ConfigurationException; import org.springframework.integration.channel.ChannelRegistry; import org.springframework.integration.channel.ChannelRegistryAware; @@ -91,7 +95,9 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio private volatile ConcurrencyPolicy defaultConcurrencyPolicy; - private volatile boolean autoCreateChannels; + private volatile boolean configureAsyncEventMulticaster = false; + + private volatile boolean autoCreateChannels = false; private volatile boolean autoStartup = true; @@ -159,6 +165,17 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio this.autoCreateChannels = autoCreateChannels; } + /** + * Set whether the bus should configure its asynchronous task executor + * to also be used by the ApplicationContext's 'applicationEventMulticaster'. + * This will only apply if the multicaster defined within the context + * is an instance of SimpleApplicationEventMulticaster (the default). + * This property is 'false' by default. + */ + public void setConfigureAsyncEventMulticaster(boolean configureAsyncEventMulticaster) { + this.configureAsyncEventMulticaster = configureAsyncEventMulticaster; + } + @SuppressWarnings("unchecked") private void registerChannels(ApplicationContext context) { Map channelBeans = @@ -479,10 +496,25 @@ public class MessageBus implements ChannelRegistry, EndpointRegistry, Applicatio if (event instanceof ContextRefreshedEvent) { ApplicationContext context = ((ContextRefreshedEvent) event).getApplicationContext(); this.registerEndpoints(context); + if (this.configureAsyncEventMulticaster) { + this.initialize(); + this.doConfigureAsyncEventMulticaster(context); + } if (this.autoStartup) { this.start(); } } } + private void doConfigureAsyncEventMulticaster(ApplicationContext context) { + String multicasterBeanName = AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME; + if (context.containsBean(multicasterBeanName)) { + ApplicationEventMulticaster multicaster = + (ApplicationEventMulticaster) context.getBean(multicasterBeanName); + if (multicaster instanceof SimpleApplicationEventMulticaster) { + ((SimpleApplicationEventMulticaster) multicaster).setTaskExecutor(this.taskScheduler); + } + } + } + } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java index e4cd9d5da8..30822f11ce 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/MessageBusParser.java @@ -111,7 +111,7 @@ public class MessageBusParser extends AbstractSimpleBeanDefinitionParser { @Override protected void doParse(Element element, ParserContext parserContext, BeanDefinitionBuilder builder) { super.doParse(element, parserContext, builder); - addPostProcessors(parserContext); + this.addPostProcessors(parserContext); } /** diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd index 93e7b7dee5..d2d0b3a149 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/spring-integration-core-1.0.xsd @@ -32,6 +32,7 @@ + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java index 59568afdc6..f63eda029d 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/MessageBusParserTests.java @@ -22,10 +22,14 @@ import static org.junit.Assert.assertTrue; import org.junit.Test; +import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanCreationException; import org.springframework.beans.factory.BeanDefinitionStoreException; import org.springframework.context.ApplicationContext; +import org.springframework.context.event.SimpleApplicationEventMulticaster; +import org.springframework.context.support.AbstractApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; +import org.springframework.core.task.SyncTaskExecutor; import org.springframework.integration.ConfigurationException; import org.springframework.integration.bus.MessageBus; import org.springframework.integration.bus.TestMessageBusAwareImpl; @@ -33,6 +37,7 @@ import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.endpoint.TargetEndpoint; import org.springframework.integration.handler.TestHandlers; +import org.springframework.integration.scheduling.SimpleMessagingTaskScheduler; import org.springframework.integration.scheduling.Subscription; /** @@ -154,4 +159,39 @@ public class MessageBusParserTests { assertEquals(QueueChannel.class, context.getBean("specifiedTypeChannel").getClass()); } + @Test + public void testMulticasterIsSyncByDefault() { + ApplicationContext context = new ClassPathXmlApplicationContext( + "messageBusWithDefaults.xml", this.getClass()); + SimpleApplicationEventMulticaster multicaster = (SimpleApplicationEventMulticaster) + context.getBean(AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME); + DirectFieldAccessor accessor = new DirectFieldAccessor(multicaster); + Object taskExecutor = accessor.getPropertyValue("taskExecutor"); + assertEquals(SyncTaskExecutor.class, taskExecutor.getClass()); + } + + @Test + public void testAsyncMulticasterExplicitlySetToFalse() throws Exception { + AbstractApplicationContext context = new ClassPathXmlApplicationContext( + "messageBusWithoutAsyncEventMulticaster.xml", this.getClass()); + context.refresh(); + SimpleApplicationEventMulticaster multicaster = (SimpleApplicationEventMulticaster) + context.getBean(AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME); + DirectFieldAccessor accessor = new DirectFieldAccessor(multicaster); + Object taskExecutor = accessor.getPropertyValue("taskExecutor"); + assertEquals(SyncTaskExecutor.class, taskExecutor.getClass()); + } + + @Test + public void testAsyncMulticaster() throws Exception { + AbstractApplicationContext context = new ClassPathXmlApplicationContext( + "messageBusWithAsyncEventMulticaster.xml", this.getClass()); + context.refresh(); + SimpleApplicationEventMulticaster multicaster = (SimpleApplicationEventMulticaster) + context.getBean(AbstractApplicationContext.APPLICATION_EVENT_MULTICASTER_BEAN_NAME); + DirectFieldAccessor accessor = new DirectFieldAccessor(multicaster); + Object taskExecutor = accessor.getPropertyValue("taskExecutor"); + assertEquals(SimpleMessagingTaskScheduler.class, taskExecutor.getClass()); + } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithAsyncEventMulticaster.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithAsyncEventMulticaster.xml new file mode 100644 index 0000000000..681f25c01f --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithAsyncEventMulticaster.xml @@ -0,0 +1,12 @@ + + + + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithoutAsyncEventMulticaster.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithoutAsyncEventMulticaster.xml new file mode 100644 index 0000000000..763f811cc1 --- /dev/null +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/messageBusWithoutAsyncEventMulticaster.xml @@ -0,0 +1,12 @@ + + + + + +