Customize concurrency at listener level
Prior to this commit, customizing the concurrency to use fo a given JMS listener involved to define it in a specific listener-container. As this is quite restrictive, users may stop using the XML namespace support altogether to fallback on regular abstract bean definition for the container. This commit adds a concurrency attribute to the jms and jca listener element as well as on the @JmsListener annotation. If the value is set, it takes precedence; otherwise the value provided by the factory is used. Issue: SPR-11988
This commit is contained in:
@@ -107,4 +107,13 @@ public @interface JmsListener {
|
||||
*/
|
||||
String selector() default "";
|
||||
|
||||
/**
|
||||
* The concurrency for the listener, if any.
|
||||
* <p>The concurrency limits can be a "lower-upper" String, e.g. "5-10", or a simple
|
||||
* upper limit String, e.g. "10" (the lower limit will be 1 in this case).
|
||||
* <p>The underlying container may or may not support all features. For instance, it
|
||||
* may not be able to scale: in that case only the upper value is used.
|
||||
*/
|
||||
String concurrency() default "";
|
||||
|
||||
}
|
||||
|
||||
@@ -181,6 +181,9 @@ public class JmsListenerAnnotationBeanPostProcessor implements BeanPostProcessor
|
||||
if (StringUtils.hasText(jmsListener.subscription())) {
|
||||
endpoint.setSubscription(jmsListener.subscription());
|
||||
}
|
||||
if (StringUtils.hasText(jmsListener.concurrency())) {
|
||||
endpoint.setConcurrency(jmsListener.concurrency());
|
||||
}
|
||||
|
||||
JmsListenerContainerFactory<?> factory = null;
|
||||
String containerFactoryBeanName = jmsListener.containerFactory();
|
||||
|
||||
@@ -42,6 +42,8 @@ public abstract class AbstractJmsListenerEndpoint implements JmsListenerEndpoint
|
||||
|
||||
private String selector;
|
||||
|
||||
private String concurrency;
|
||||
|
||||
|
||||
public void setId(String id) {
|
||||
this.id = id;
|
||||
@@ -95,6 +97,23 @@ public abstract class AbstractJmsListenerEndpoint implements JmsListenerEndpoint
|
||||
return this.selector;
|
||||
}
|
||||
|
||||
/**
|
||||
* Set a concurrency for the listener, if any.
|
||||
* <p>The concurrency limits can be a "lower-upper" String, e.g. "5-10", or a simple
|
||||
* upper limit String, e.g. "10" (the lower limit will be 1 in this case).
|
||||
* <p>The underlying container may or may not support all features. For instance, it
|
||||
* may not be able to scale: in that case only the upper value is used.
|
||||
*/
|
||||
public void setConcurrency(String concurrency) {
|
||||
this.concurrency = concurrency;
|
||||
}
|
||||
|
||||
/**
|
||||
* Return the concurrency for the listener, if any.
|
||||
*/
|
||||
public String getConcurrency() {
|
||||
return concurrency;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setupMessageContainer(MessageListenerContainer container) {
|
||||
@@ -121,6 +140,9 @@ public abstract class AbstractJmsListenerEndpoint implements JmsListenerEndpoint
|
||||
if (getSelector() != null) {
|
||||
container.setMessageSelector(getSelector());
|
||||
}
|
||||
if (getConcurrency() != null) {
|
||||
container.setConcurrency(getConcurrency());
|
||||
}
|
||||
setupMessageListener(container);
|
||||
}
|
||||
|
||||
@@ -140,6 +162,9 @@ public abstract class AbstractJmsListenerEndpoint implements JmsListenerEndpoint
|
||||
if (getSelector() != null) {
|
||||
activationSpecConfig.setMessageSelector(getSelector());
|
||||
}
|
||||
if (getConcurrency() != null) {
|
||||
activationSpecConfig.setConcurrency(getConcurrency());
|
||||
}
|
||||
setupMessageListener(container);
|
||||
}
|
||||
|
||||
|
||||
@@ -247,6 +247,15 @@ abstract class AbstractListenerContainerParser implements BeanDefinitionParser {
|
||||
}
|
||||
configDef.getPropertyValues().add("messageSelector", selector);
|
||||
}
|
||||
|
||||
if (ele.hasAttribute(CONCURRENCY_ATTRIBUTE)) {
|
||||
String concurrency = ele.getAttribute(CONCURRENCY_ATTRIBUTE);
|
||||
if (!StringUtils.hasText(concurrency)) {
|
||||
parserContext.getReaderContext().error(
|
||||
"Listener 'concurrency' attribute contains empty value.", ele);
|
||||
}
|
||||
configDef.getPropertyValues().add("concurrency", concurrency);
|
||||
}
|
||||
}
|
||||
|
||||
protected PropertyValues parseCommonContainerProperties(Element ele, ParserContext parserContext) {
|
||||
|
||||
@@ -70,10 +70,10 @@ class JcaListenerContainerParser extends AbstractListenerContainerParser {
|
||||
containerDef.setSource(context.getSource());
|
||||
containerDef.setBeanClassName("org.springframework.jms.listener.endpoint.JmsMessageEndpointManager");
|
||||
|
||||
containerDef.getPropertyValues().addPropertyValues(context.getContainerValues());
|
||||
applyContainerValues(context, containerDef);
|
||||
|
||||
|
||||
BeanDefinition activationSpec = getActivationSpecConfigBeanDefinition(context.getContainerValues());
|
||||
BeanDefinition activationSpec = getActivationSpecConfigBeanDefinition(containerDef.getPropertyValues());
|
||||
parseListenerConfiguration(context.getListenerElement(), context.getParserContext(), activationSpec);
|
||||
|
||||
String phase = context.getContainerElement().getAttribute(PHASE_ATTRIBUTE);
|
||||
@@ -84,6 +84,21 @@ class JcaListenerContainerParser extends AbstractListenerContainerParser {
|
||||
return containerDef;
|
||||
}
|
||||
|
||||
/**
|
||||
* The property values provided by the factory element contains a mutable property (the
|
||||
* activation spec config). To avoid changing the bean definition from the parent, a clone
|
||||
* bean definition is created for the container being configured.
|
||||
*/
|
||||
private void applyContainerValues(ListenerContainerParserContext context, RootBeanDefinition containerDef) {
|
||||
// Apply settings from the container
|
||||
containerDef.getPropertyValues().addPropertyValues(context.getContainerValues());
|
||||
|
||||
// Clone the activationSpecConfig property value as it is mutable
|
||||
PropertyValue pv = containerDef.getPropertyValues().getPropertyValue("activationSpecConfig");
|
||||
RootBeanDefinition activationSpecConfig = new RootBeanDefinition((RootBeanDefinition) pv.getValue());
|
||||
containerDef.getPropertyValues().add("activationSpecConfig", activationSpecConfig);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected boolean indicatesPubSub(PropertyValues propertyValues) {
|
||||
BeanDefinition configDef = getActivationSpecConfigBeanDefinition(propertyValues);
|
||||
|
||||
@@ -92,7 +92,6 @@ class JmsListenerContainerParser extends AbstractListenerContainerParser {
|
||||
|
||||
// Set all container values
|
||||
containerDef.getPropertyValues().addPropertyValues(context.getContainerValues());
|
||||
parseListenerConfiguration(context.getListenerElement(), context.getParserContext(), containerDef);
|
||||
|
||||
Element containerEle = context.getContainerElement();
|
||||
String containerType = containerEle.getAttribute(CONTAINER_TYPE_ATTRIBUTE);
|
||||
@@ -116,6 +115,9 @@ class JmsListenerContainerParser extends AbstractListenerContainerParser {
|
||||
containerDef.getPropertyValues().add("phase", phase);
|
||||
}
|
||||
|
||||
// Parse listener specific settings
|
||||
parseListenerConfiguration(context.getListenerElement(), context.getParserContext(), containerDef);
|
||||
|
||||
return containerDef;
|
||||
}
|
||||
|
||||
|
||||
@@ -147,6 +147,11 @@ public abstract class AbstractMessageListenerContainer
|
||||
private boolean acceptMessagesWhileStopping = false;
|
||||
|
||||
|
||||
/**
|
||||
* Specify concurrency limits.
|
||||
*/
|
||||
public abstract void setConcurrency(String concurrency);
|
||||
|
||||
/**
|
||||
* Set the destination to receive messages from.
|
||||
* <p>Alternatively, specify a "destinationName", to be dynamically
|
||||
|
||||
@@ -295,6 +295,7 @@ public class DefaultMessageListenerContainer extends AbstractPollingMessageListe
|
||||
* ({@link #setConcurrentConsumers}) and will slowly scale up to the maximum number
|
||||
* of consumers {@link #setMaxConcurrentConsumers} in case of increasing load.
|
||||
*/
|
||||
@Override
|
||||
public void setConcurrency(String concurrency) {
|
||||
try {
|
||||
int separatorIndex = concurrency.indexOf('-');
|
||||
|
||||
@@ -114,6 +114,7 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta
|
||||
* {@link DefaultMessageListenerContainer}. For this local listener container,
|
||||
* generally use {@link #setConcurrentConsumers} instead.
|
||||
*/
|
||||
@Override
|
||||
public void setConcurrency(String concurrency) {
|
||||
try {
|
||||
int separatorIndex = concurrency.indexOf('-');
|
||||
|
||||
@@ -590,6 +590,17 @@
|
||||
]]></xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="concurrency" type="xsd:string">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation><![CDATA[
|
||||
The number of concurrent sessions/consumers to start for this listener.
|
||||
Can either be a simple number indicating the maximum number (e.g. "5")
|
||||
or a range indicating the lower as well as the upper limit (e.g. "3-5").
|
||||
Note that a specified minimum is just a hint and might be ignored at runtime.
|
||||
Default is the value provided by the container.
|
||||
]]></xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
</xsd:complexType>
|
||||
|
||||
</xsd:schema>
|
||||
|
||||
@@ -25,8 +25,10 @@ import org.junit.Rule;
|
||||
import org.junit.Test;
|
||||
import org.junit.rules.ExpectedException;
|
||||
|
||||
import org.springframework.beans.DirectFieldAccessor;
|
||||
import org.springframework.jms.listener.DefaultMessageListenerContainer;
|
||||
import org.springframework.jms.listener.MessageListenerContainer;
|
||||
import org.springframework.jms.listener.SimpleMessageListenerContainer;
|
||||
import org.springframework.jms.listener.adapter.MessageListenerAdapter;
|
||||
import org.springframework.jms.listener.endpoint.JmsActivationSpecConfig;
|
||||
import org.springframework.jms.listener.endpoint.JmsMessageEndpointManager;
|
||||
@@ -48,12 +50,15 @@ public class JmsListenerEndpointTests {
|
||||
endpoint.setDestination("myQueue");
|
||||
endpoint.setSelector("foo = 'bar'");
|
||||
endpoint.setSubscription("mySubscription");
|
||||
endpoint.setConcurrency("5-10");
|
||||
endpoint.setMessageListener(messageListener);
|
||||
|
||||
endpoint.setupMessageContainer(container);
|
||||
assertEquals("myQueue", container.getDestinationName());
|
||||
assertEquals("foo = 'bar'", container.getMessageSelector());
|
||||
assertEquals("mySubscription", container.getDurableSubscriptionName());
|
||||
assertEquals(5, container.getConcurrentConsumers());
|
||||
assertEquals(10, container.getMaxConcurrentConsumers());
|
||||
assertEquals(messageListener, container.getMessageListener());
|
||||
}
|
||||
|
||||
@@ -65,6 +70,7 @@ public class JmsListenerEndpointTests {
|
||||
endpoint.setDestination("myQueue");
|
||||
endpoint.setSelector("foo = 'bar'");
|
||||
endpoint.setSubscription("mySubscription");
|
||||
endpoint.setConcurrency("10");
|
||||
endpoint.setMessageListener(messageListener);
|
||||
|
||||
endpoint.setupMessageContainer(container);
|
||||
@@ -72,9 +78,21 @@ public class JmsListenerEndpointTests {
|
||||
assertEquals("myQueue", config.getDestinationName());
|
||||
assertEquals("foo = 'bar'", config.getMessageSelector());
|
||||
assertEquals("mySubscription", config.getDurableSubscriptionName());
|
||||
assertEquals(10, config.getMaxConcurrency());
|
||||
assertEquals(messageListener, container.getMessageListener());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void setupConcurrencySimpleContainer() {
|
||||
SimpleMessageListenerContainer container = new SimpleMessageListenerContainer();
|
||||
MessageListener messageListener = new MessageListenerAdapter();
|
||||
SimpleJmsListenerEndpoint endpoint = new SimpleJmsListenerEndpoint();
|
||||
endpoint.setConcurrency("5-10"); // simple implementation only support max value
|
||||
endpoint.setMessageListener(messageListener);
|
||||
|
||||
endpoint.setupMessageContainer(container);
|
||||
assertEquals(10, new DirectFieldAccessor(container).getPropertyValue("concurrentConsumers"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void setupMessageContainerNoListener() {
|
||||
|
||||
@@ -102,13 +102,9 @@ public class JmsNamespaceHandlerTests {
|
||||
for (DefaultMessageListenerContainer container : containers.values()) {
|
||||
if (container.getConnectionFactory().equals(defaultConnectionFactory)) {
|
||||
defaultConnectionFactoryCount++;
|
||||
assertEquals(2, container.getConcurrentConsumers());
|
||||
assertEquals(3, container.getMaxConcurrentConsumers());
|
||||
}
|
||||
else if (container.getConnectionFactory().equals(explicitConnectionFactory)) {
|
||||
explicitConnectionFactoryCount++;
|
||||
assertEquals(3, container.getConcurrentConsumers());
|
||||
assertEquals(5, container.getMaxConcurrentConsumers());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -226,6 +222,34 @@ public class JmsNamespaceHandlerTests {
|
||||
assertEquals(DefaultMessageListenerContainer.DEFAULT_RECOVERY_INTERVAL, recoveryInterval3);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testConcurrency() {
|
||||
// JMS
|
||||
DefaultMessageListenerContainer listener0 = this.context
|
||||
.getBean(DefaultMessageListenerContainer.class.getName() + "#0", DefaultMessageListenerContainer.class);
|
||||
DefaultMessageListenerContainer listener1 = this.context
|
||||
.getBean("listener1", DefaultMessageListenerContainer.class);
|
||||
DefaultMessageListenerContainer listener2 = this.context
|
||||
.getBean("listener2", DefaultMessageListenerContainer.class);
|
||||
|
||||
assertEquals("Wrong concurrency on listener using placeholder", 2, listener0.getConcurrentConsumers());
|
||||
assertEquals("Wrong concurrency on listener using placeholder", 3, listener0.getMaxConcurrentConsumers());
|
||||
assertEquals("Wrong concurrency on listener1", 3, listener1.getConcurrentConsumers());
|
||||
assertEquals("Wrong max concurrency on listener1", 5, listener1.getMaxConcurrentConsumers());
|
||||
assertEquals("Wrong custom concurrency on listener2", 5, listener2.getConcurrentConsumers());
|
||||
assertEquals("Wrong custom max concurrency on listener2", 10, listener2.getMaxConcurrentConsumers());
|
||||
|
||||
// JCA
|
||||
JmsMessageEndpointManager listener3 = this.context
|
||||
.getBean("listener3", JmsMessageEndpointManager.class);
|
||||
JmsMessageEndpointManager listener4 = this.context
|
||||
.getBean("listener4", JmsMessageEndpointManager.class);
|
||||
assertEquals("Wrong concurrency on listener3", 5,
|
||||
listener3.getActivationSpecConfig().getMaxConcurrency());
|
||||
assertEquals("Wrong custom concurrency on listener4", 7,
|
||||
listener4.getActivationSpecConfig().getMaxConcurrency());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testErrorHandlers() {
|
||||
ErrorHandler expected = this.context.getBean("testErrorHandler", ErrorHandler.class);
|
||||
|
||||
@@ -11,7 +11,8 @@
|
||||
transaction-manager="testTransactionManager" error-handler="testErrorHandler"
|
||||
cache="connection" concurrency="3-5" prefetch="50" receive-timeout="100" back-off="testBackOff" phase="99">
|
||||
<jms:listener id="listener1" destination="testDestination" ref="testBean1" method="setName"/>
|
||||
<jms:listener id="listener2" destination="testDestination" ref="testBean2" method="setName" response-destination="responseDestination"/>
|
||||
<jms:listener id="listener2" destination="testDestination" ref="testBean2" method="setName"
|
||||
concurrency="5-10" response-destination="responseDestination"/>
|
||||
</jms:listener-container>
|
||||
|
||||
<!-- TODO: remove the task-executor reference once issue with blocking on stop is resolved -->
|
||||
@@ -31,7 +32,8 @@
|
||||
resource-adapter="testResourceAdapter" activation-spec-factory="testActivationSpecFactory"
|
||||
message-converter="testMessageConverter" concurrency="5" prefetch="50" phase="77">
|
||||
<jms:listener id="listener3" destination="testDestination" ref="testBean1" method="setName"/>
|
||||
<jms:listener id="listener4" destination="testDestination" ref="testBean2" method="setName" response-destination="responseDestination"/>
|
||||
<jms:listener id="listener4" destination="testDestination" ref="testBean2" method="setName"
|
||||
concurrency="7" response-destination="responseDestination"/>
|
||||
</jms:jca-listener-container>
|
||||
|
||||
<jms:jca-listener-container activation-spec-factory="testActivationSpecFactory">
|
||||
|
||||
@@ -41661,6 +41661,12 @@ describes all available attributes:
|
||||
|
||||
| selector
|
||||
| An optional message selector for this listener.
|
||||
|
||||
| concurrency
|
||||
| The number of concurrent sessions/consumers to start for this listener. Can either be
|
||||
| a simple number indicating the maximum number (e.g. "5") or a range indicating the lower
|
||||
| as well as the upper limit (e.g. "3-5"). Note that a specified minimum is just a hint
|
||||
| and might be ignored at runtime. Default is the value provided by the container
|
||||
|===
|
||||
|
||||
The `<listener-container/>` element also accepts several optional attributes. This
|
||||
|
||||
Reference in New Issue
Block a user