diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java index ba204e081c..5823ae6fda 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java @@ -33,6 +33,7 @@ import org.springframework.integration.Message; import org.springframework.integration.MessageDeliveryException; import org.springframework.integration.MessageDispatchingException; import org.springframework.integration.MessagingException; +import org.springframework.integration.context.IntegrationProperties; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.SubscribableChannel; import org.springframework.integration.dispatcher.AbstractDispatcher; @@ -43,6 +44,7 @@ import org.springframework.util.Assert; /** * @author Mark Fisher * @author Gary Russell + * @author Artem Bilan * @since 2.1 */ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel implements SubscribableChannel, SmartLifecycle, DisposableBean { @@ -51,11 +53,11 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple private final SimpleMessageListenerContainer container; - private volatile MessageDispatcher dispatcher; + private volatile AbstractDispatcher dispatcher; private final boolean isPubSub; - private volatile int maxSubscribers = Integer.MAX_VALUE; + private volatile Integer maxSubscribers; public AbstractSubscribableAmqpChannel(String channelName, SimpleMessageListenerContainer container, AmqpTemplate amqpTemplate) { this(channelName, container, amqpTemplate, false); @@ -79,6 +81,9 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple */ public void setMaxSubscribers(int maxSubscribers) { this.maxSubscribers = maxSubscribers; + if (this.dispatcher != null) { + this.dispatcher.setMaxSubscribers(this.maxSubscribers); + } } public boolean subscribe(MessageHandler handler) { @@ -93,9 +98,13 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple public void onInit() throws Exception { super.onInit(); this.dispatcher = this.createDispatcher(); - if (this.dispatcher instanceof AbstractDispatcher) { - ((AbstractDispatcher) this.dispatcher).setMaxSubscribers(this.maxSubscribers); + if (this.maxSubscribers == null) { + this.maxSubscribers = this.getIntegrationProperty(this.isPubSub ? + IntegrationProperties.CHANNELS_MAX_BROADCAST_SUBSCRIBERS : + IntegrationProperties.CHANNELS_MAX_UNICAST_SUBSCRIBERS, + Integer.class); } + this.setMaxSubscribers(this.maxSubscribers); AmqpAdmin admin = new RabbitAdmin(this.container.getConnectionFactory()); Queue queue = this.initializeQueue(admin, this.channelName); this.container.setQueues(queue); @@ -110,7 +119,7 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple } } - protected abstract MessageDispatcher createDispatcher(); + protected abstract AbstractDispatcher createDispatcher(); protected abstract Queue initializeQueue(AmqpAdmin admin, String channelName); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java index ded7d382e2..181f9d40d9 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2011 the original author or authors. + * Copyright 2002-2013 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -20,7 +20,7 @@ import org.springframework.amqp.core.AmqpAdmin; import org.springframework.amqp.core.AmqpTemplate; import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; -import org.springframework.integration.dispatcher.MessageDispatcher; +import org.springframework.integration.dispatcher.AbstractDispatcher; import org.springframework.integration.dispatcher.RoundRobinLoadBalancingStrategy; import org.springframework.integration.dispatcher.UnicastingDispatcher; @@ -57,7 +57,7 @@ public class PointToPointSubscribableAmqpChannel extends AbstractSubscribableAmq } @Override - protected MessageDispatcher createDispatcher() { + protected AbstractDispatcher createDispatcher() { UnicastingDispatcher unicastingDispatcher = new UnicastingDispatcher(); unicastingDispatcher.setLoadBalancingStrategy(new RoundRobinLoadBalancingStrategy()); return unicastingDispatcher; diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java index a3b979e70a..4f7af317ff 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2012 the original author or authors. + * Copyright 2002-2013 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -23,8 +23,8 @@ import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.FanoutExchange; import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; +import org.springframework.integration.dispatcher.AbstractDispatcher; import org.springframework.integration.dispatcher.BroadcastingDispatcher; -import org.springframework.integration.dispatcher.MessageDispatcher; /** * @author Mark Fisher @@ -65,7 +65,7 @@ public class PublishSubscribeAmqpChannel extends AbstractSubscribableAmqpChannel } @Override - protected MessageDispatcher createDispatcher() { + protected AbstractDispatcher createDispatcher() { return new BroadcastingDispatcher(true); } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java index d09a04070b..934afde5ad 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpChannelFactoryBean.java @@ -121,7 +121,7 @@ public class AmqpChannelFactoryBean extends AbstractFactoryBean T getIntegrationProperty(String key, Class tClass) { + return this.defaultConversionService.convert(this.integrationProperties.getProperty(key), tClass); } @Override diff --git a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java index c54cb0dae3..38faa7a16d 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/context/IntegrationProperties.java @@ -16,16 +16,96 @@ package org.springframework.integration.context; +import java.io.IOException; +import java.util.Properties; + +import org.springframework.beans.factory.config.PropertiesFactoryBean; +import org.springframework.core.io.Resource; +import org.springframework.core.io.support.PathMatchingResourcePatternResolver; +import org.springframework.core.io.support.ResourcePatternResolver; + /** - * Convention Enumeration to represent keys from 'META-INF/spring.integration.properties'. + * Utility class to encapsulate infrastructure Integration properties constants and + * their default values from resources 'META-INF/spring.integration.default.properties'. * * @author Artem Bilan * @since 3.0 */ -public interface IntegrationProperties { +public final class IntegrationProperties { - String LATE_REPLY_LOGGING_LEVEL = "messagingTemplate.lateReply.logging.level"; + /** + * Specifies whether to allow create automatically {@link org.springframework.integration.channel.DirectChannel} + * beans for non-declared channels or not. + */ + public static final String CHANNELS_AUTOCREATE = "channels.autoCreate"; - String CHANNELS_AUTOCREATE = "channels.autoCreate"; + /** + * Specifies the value for {@link org.springframework.integration.dispatcher.UnicastingDispatcher#maxSubscribers} + * in case of point-to-point channels (e.g. {@link org.springframework.integration.channel.ExecutorChannel}), + * if the attribute {@code max-subscribers} isn't configured on the channel component. + */ + public static final String CHANNELS_MAX_UNICAST_SUBSCRIBERS = "channels.maxUnicastSubscribers"; + + /** + * Specifies the value for {@link org.springframework.integration.dispatcher.BroadcastingDispatcher#maxSubscribers} + * in case of point-to-point channels (e.g. {@link org.springframework.integration.channel.PublishSubscribeChannel}), + * if the attribute {@code max-subscribers} isn't configured on the channel component. + */ + public static final String CHANNELS_MAX_BROADCAST_SUBSCRIBERS = "channels.maxBroadcastSubscribers"; + + /** + * Specifies the value of {@link org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler#poolSize} + * for {@code taskScheduler} bean initialized but Integration infrastructure. + * @see org.springframework.integration.config.xml.DefaultConfiguringBeanFactoryPostProcessor#registerTaskScheduler + */ + public static final String TASKSCHEDULER_POOLSIZE = "taskScheduler.poolSize"; + +// TODO public static final String LATE_REPLY_LOGGING_LEVEL = "messagingTemplate.lateReply.logging.level"; + + private static Properties defaults; + + static { + String resourcePattern = "classpath*:META-INF/spring.integration.default.properties"; + try { + ResourcePatternResolver resourceResolver = new PathMatchingResourcePatternResolver(IntegrationProperties.class.getClassLoader()); + Resource[] defaultResources = resourceResolver.getResources(resourcePattern); + + PropertiesFactoryBean propertiesFactoryBean = new PropertiesFactoryBean(); + propertiesFactoryBean.setLocations(defaultResources); + propertiesFactoryBean.afterPropertiesSet(); + defaults = propertiesFactoryBean.getObject(); + } + catch (IOException e) { + throw new IllegalStateException("Can't load '" + resourcePattern + "' resources.", e); + } + } + + /** + * @return {@link Properties} with default values for Integration properties + * from resources 'META-INF/spring.integration.default.properties'. + */ + public static Properties defaults() { + return defaults; + } + + /** + * Build the bean property definition expression to resolve the value + * from Integration properties within the bean building phase. + * + * @param key the Integration property key. + * @return the bean property definition expression. + * @throws IllegalArgumentException if provided {@code key} isn't an Integration property. + */ + public static String getExpressionFor(String key) { + if (defaults.containsKey(key)) { + return "#{T(" + IntegrationContextUtils.class.getName() + ").getIntegrationProperties(beanFactory).getProperty('" + key + "')}"; + } + else { + throw new IllegalArgumentException("The provided key [" + key + "] isn't the one of Integration properties: " + defaults.keySet()); + } + } + + private IntegrationProperties() { + } } diff --git a/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties b/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties index 199b366130..3633615b97 100644 --- a/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties +++ b/spring-integration-core/src/main/resources/META-INF/spring.integration.default.properties @@ -1,2 +1,4 @@ channels.autoCreate=true -messagingTemplate.lateReply.logging.level=warn +channels.maxUnicastSubscribers=0x7fffffff +channels.maxBroadcastSubscribers=0x7fffffff +taskScheduler.poolSize=10 diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersDefaultConfigurationTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersDefaultConfigurationTests-context.xml index 7d4a9c1182..902587d6ad 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersDefaultConfigurationTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersDefaultConfigurationTests-context.xml @@ -7,6 +7,8 @@ http://www.springframework.org/schema/task http://www.springframework.org/schema/task/spring-task.xsd http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd"> + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersOverrideDefaultTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersOverrideDefaultTests-context.xml index 388be5668a..b4a39d9398 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersOverrideDefaultTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersOverrideDefaultTests-context.xml @@ -1,15 +1,16 @@ + xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" + xmlns:int="http://www.springframework.org/schema/integration" + xmlns:util="http://www.springframework.org/schema/util" + xsi:schemaLocation="http://www.springframework.org/schema/integration http://www.springframework.org/schema/integration/spring-integration.xsd + http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd + http://www.springframework.org/schema/util http://www.springframework.org/schema/util/spring-util.xsd"> - - - - - + + 456 + 789 + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersTests.java index 79c7701824..21cba2662b 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/xml/DispatcherMaxSubscribersTests.java @@ -23,11 +23,15 @@ import org.springframework.integration.test.util.TestUtils; /** * @author Gary Russell + * @author Artem Bilan * @since 2.2 * */ public abstract class DispatcherMaxSubscribersTests { + @Autowired + private MessageChannel autoCreateChannel; + @Autowired private MessageChannel defaultChannel; @@ -54,20 +58,17 @@ public abstract class DispatcherMaxSubscribersTests { } protected void doTestUnicast(int val1, int val2, int val3, int val4, int val5) { - Integer defaultMax = TestUtils.getPropertyValue( - TestUtils.getPropertyValue(defaultChannel, "dispatcher"), "maxSubscribers", Integer.class); + Integer autoCreateMax = TestUtils.getPropertyValue(autoCreateChannel, "dispatcher.maxSubscribers", Integer.class); + assertEquals(val1, autoCreateMax.intValue()); + Integer defaultMax = TestUtils.getPropertyValue(defaultChannel, "dispatcher.maxSubscribers", Integer.class); assertEquals(val1, defaultMax.intValue()); - Integer defaultMax2 = TestUtils.getPropertyValue( - TestUtils.getPropertyValue(defaultChannel2, "dispatcher"), "maxSubscribers", Integer.class); + Integer defaultMax2 = TestUtils.getPropertyValue(defaultChannel2, "dispatcher.maxSubscribers", Integer.class); assertEquals(val2, defaultMax2.intValue()); - Integer explicitMax = TestUtils.getPropertyValue( - TestUtils.getPropertyValue(explicitChannel, "dispatcher"), "maxSubscribers", Integer.class); + Integer explicitMax = TestUtils.getPropertyValue(explicitChannel, "dispatcher.maxSubscribers", Integer.class); assertEquals(val3, explicitMax.intValue()); - Integer execMax = TestUtils.getPropertyValue( - TestUtils.getPropertyValue(executorChannel, "dispatcher"), "maxSubscribers", Integer.class); + Integer execMax = TestUtils.getPropertyValue(executorChannel, "dispatcher.maxSubscribers", Integer.class); assertEquals(val4, execMax.intValue()); - Integer explicitExecMax = TestUtils.getPropertyValue( - TestUtils.getPropertyValue(explicitExecutorChannel, "dispatcher"), "maxSubscribers", Integer.class); + Integer explicitExecMax = TestUtils.getPropertyValue(explicitExecutorChannel, "dispatcher.maxSubscribers", Integer.class); assertEquals(val5, explicitExecMax.intValue()); } @@ -81,4 +82,4 @@ public abstract class DispatcherMaxSubscribersTests { Integer explicitMin = TestUtils.getPropertyValue(pubSubExplicitChannel, "dispatcher.minSubscribers", Integer.class); assertEquals(1, explicitMin.intValue()); } -} \ No newline at end of file +} diff --git a/spring-integration-core/src/test/java/org/springframework/integration/context/IntegrationContextTests.java b/spring-integration-core/src/test/java/org/springframework/integration/context/IntegrationContextTests.java index 64d9ee572e..3ad521a9bc 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/context/IntegrationContextTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/context/IntegrationContextTests.java @@ -17,7 +17,6 @@ package org.springframework.integration.context; import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertSame; import java.util.Properties; @@ -26,6 +25,8 @@ import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.integration.test.util.TestUtils; +import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -38,17 +39,22 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; public class IntegrationContextTests { @Autowired - @Qualifier(IntegrationContextUtils.INTEGRATION_PROPERTIES_BEAN_NAME) + @Qualifier(IntegrationContextUtils.INTEGRATION_GLOBAL_PROPERTIES_BEAN_NAME) private Properties integrationProperties; @Autowired @Qualifier("fooService") private IntegrationObjectSupport serviceActivator; + @Autowired + private ThreadPoolTaskScheduler taskScheduler; + @Test public void testIntegrationContextComponents() { - assertEquals("error", this.integrationProperties.get(IntegrationProperties.LATE_REPLY_LOGGING_LEVEL)); - assertSame(this.integrationProperties, this.serviceActivator.getIntegrationProperties()); + //TODO INT-3005 assertEquals("error", this.integrationProperties.get(IntegrationProperties.LATE_REPLY_LOGGING_LEVEL)); + assertEquals("20", this.integrationProperties.get(IntegrationProperties.TASKSCHEDULER_POOLSIZE)); + assertEquals(this.integrationProperties, this.serviceActivator.getIntegrationProperties()); + assertEquals(20, TestUtils.getPropertyValue(this.taskScheduler, "poolSize")); } } diff --git a/spring-integration-core/src/test/resources/META-INF/spring.integration.properties b/spring-integration-core/src/test/resources/META-INF/spring.integration.properties index 38ffd39c66..bd78272249 100644 --- a/spring-integration-core/src/test/resources/META-INF/spring.integration.properties +++ b/spring-integration-core/src/test/resources/META-INF/spring.integration.properties @@ -1,2 +1,5 @@ #channels.autoCreate=false +#channels.maxUnicastSubscribers=1 +#channels.maxBroadcastSubscribers=1 messagingTemplate.lateReply.logging.level=error +taskScheduler.poolSize=20 diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java index 41886087f7..0c420e0433 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/SubscribableJmsChannel.java @@ -26,6 +26,7 @@ import org.springframework.integration.Message; import org.springframework.integration.MessageDeliveryException; import org.springframework.integration.MessageDispatchingException; import org.springframework.integration.MessagingException; +import org.springframework.integration.context.IntegrationProperties; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.SubscribableChannel; import org.springframework.integration.dispatcher.AbstractDispatcher; @@ -51,7 +52,7 @@ public class SubscribableJmsChannel extends AbstractJmsChannel implements Subscr private volatile boolean initialized; - private volatile int maxSubscribers = Integer.MAX_VALUE; + private volatile Integer maxSubscribers; public SubscribableJmsChannel(AbstractMessageListenerContainer container, JmsTemplate jmsTemplate) { super(jmsTemplate); @@ -105,6 +106,12 @@ public class SubscribableJmsChannel extends AbstractJmsChannel implements Subscr unicastingDispatcher.setLoadBalancingStrategy(new RoundRobinLoadBalancingStrategy()); this.dispatcher = unicastingDispatcher; } + if (this.maxSubscribers == null) { + this.maxSubscribers = this.getIntegrationProperty(isPubSub ? + IntegrationProperties.CHANNELS_MAX_BROADCAST_SUBSCRIBERS : + IntegrationProperties.CHANNELS_MAX_UNICAST_SUBSCRIBERS, + Integer.class); + } this.dispatcher.setMaxSubscribers(this.maxSubscribers); } diff --git a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelParser.java b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelParser.java index 9223835a5f..24acc8f2b8 100644 --- a/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelParser.java +++ b/spring-integration-jms/src/main/java/org/springframework/integration/jms/config/JmsChannelParser.java @@ -18,12 +18,13 @@ package org.springframework.integration.jms.config; import javax.jms.Session; +import org.w3c.dom.Element; + import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractChannelParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; import org.springframework.util.StringUtils; -import org.w3c.dom.Element; /** * Parser for the 'channel' and 'publish-subscribe-channel' elements of the @@ -58,12 +59,13 @@ public class JmsChannelParser extends AbstractChannelParser { builder.addPropertyReference("connectionFactory", connectionFactory); if ("channel".equals(element.getLocalName())) { this.parseDestination(element, parserContext, builder, "queue"); - this.setMaxSubscribersProperty(parserContext, builder, element, IntegrationNamespaceUtils.DEFAULT_MAX_UNICAST_SUBSCRIBERS_PROPERTY_NAME); } else if ("publish-subscribe-channel".equals(element.getLocalName())) { this.parseDestination(element, parserContext, builder, "topic"); - this.setMaxSubscribersProperty(parserContext, builder, element, IntegrationNamespaceUtils.DEFAULT_MAX_BROADCAST_SUBSCRIBERS_PROPERTY_NAME); } + + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "max-subscribers"); + String containerType = element.getAttribute(CONTAINER_TYPE_ATTRIBUTE); String containerClass = element.getAttribute(CONTAINER_CLASS_ATTRIBUTE); if (!StringUtils.hasText(containerClass)) { diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/channel/SubscribableRedisChannel.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/channel/SubscribableRedisChannel.java index c28ea42f53..d81221b79a 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/channel/SubscribableRedisChannel.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/channel/SubscribableRedisChannel.java @@ -34,6 +34,7 @@ import org.springframework.integration.MessageDeliveryException; import org.springframework.integration.MessageDispatchingException; import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.channel.MessagePublishingErrorHandler; +import org.springframework.integration.context.IntegrationProperties; import org.springframework.integration.core.MessageHandler; import org.springframework.integration.core.SubscribableChannel; import org.springframework.integration.dispatcher.AbstractDispatcher; @@ -61,6 +62,8 @@ public class SubscribableRedisChannel extends AbstractMessageChannel implements private final AbstractDispatcher dispatcher = new BroadcastingDispatcher(true); + private volatile Integer maxSubscribers; + private volatile boolean initialized; // defaults @@ -97,6 +100,7 @@ public class SubscribableRedisChannel extends AbstractMessageChannel implements * @param maxSubscribers */ public void setMaxSubscribers(int maxSubscribers) { + this.maxSubscribers = maxSubscribers; this.dispatcher.setMaxSubscribers(maxSubscribers); } @@ -120,6 +124,10 @@ public class SubscribableRedisChannel extends AbstractMessageChannel implements return; } super.onInit(); + if (this.maxSubscribers == null) { + Integer maxSubscribers = this.getIntegrationProperty(IntegrationProperties.CHANNELS_MAX_BROADCAST_SUBSCRIBERS, Integer.class); + this.setMaxSubscribers(maxSubscribers); + } if (this.messageConverter == null){ this.messageConverter = new SimpleMessageConverter(); } diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisChannelParser.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisChannelParser.java index 42432169bb..369646a017 100644 --- a/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisChannelParser.java +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/config/RedisChannelParser.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2012 the original author or authors. + * Copyright 2002-2013 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -16,13 +16,14 @@ package org.springframework.integration.redis.config; +import org.w3c.dom.Element; + import org.springframework.beans.factory.support.BeanDefinitionBuilder; import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.config.xml.AbstractChannelParser; import org.springframework.integration.config.xml.IntegrationNamespaceUtils; import org.springframework.integration.redis.channel.SubscribableRedisChannel; import org.springframework.util.StringUtils; -import org.w3c.dom.Element; /** * Parser for the 'channel' and 'publish-subscribe-channel' elements of the @@ -52,7 +53,8 @@ public class RedisChannelParser extends AbstractChannelParser { // The following 2 attributes should be added once configurable on the RedisMessageListenerContainer // IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "phase"); // IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup"); - this.setMaxSubscribersProperty(parserContext, builder, element, IntegrationNamespaceUtils.DEFAULT_MAX_BROADCAST_SUBSCRIBERS_PROPERTY_NAME); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "max-subscribers"); + return builder; }