INT-1830 added support for configuring concurrency on jms channel
This commit is contained in:
@@ -43,6 +43,7 @@ import org.springframework.transaction.PlatformTransactionManager;
|
||||
import org.springframework.util.Assert;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
import org.springframework.util.ErrorHandler;
|
||||
import org.springframework.util.StringUtils;
|
||||
|
||||
/**
|
||||
* @author Mark Fisher
|
||||
@@ -71,7 +72,8 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
|
||||
private volatile String clientId;
|
||||
|
||||
private volatile Integer concurrentConsumers;
|
||||
private volatile String concurrency;
|
||||
//private volatile Integer concurrentConsumers;
|
||||
|
||||
private volatile ConnectionFactory connectionFactory;
|
||||
|
||||
@@ -91,7 +93,7 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
|
||||
private volatile Integer idleTaskExecutionLimit;
|
||||
|
||||
private volatile Integer maxConcurrentConsumers;
|
||||
//private volatile Integer maxConcurrentConsumers;
|
||||
|
||||
private volatile Integer maxMessagesPerTask;
|
||||
|
||||
@@ -195,9 +197,12 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
this.clientId = clientId;
|
||||
}
|
||||
|
||||
public void setConcurrentConsumers(int concurrentConsumers) {
|
||||
this.concurrentConsumers = concurrentConsumers;
|
||||
public void setConcurrency(String concurrency) {
|
||||
this.concurrency = concurrency;
|
||||
}
|
||||
// public void setConcurrentConsumers(int concurrentConsumers) {
|
||||
// this.concurrentConsumers = concurrentConsumers;
|
||||
// }
|
||||
|
||||
public void setConnectionFactory(ConnectionFactory connectionFactory) {
|
||||
this.connectionFactory = connectionFactory;
|
||||
@@ -241,9 +246,9 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
this.idleTaskExecutionLimit = idleTaskExecutionLimit;
|
||||
}
|
||||
|
||||
public void setMaxConcurrentConsumers(int maxConcurrentConsumers) {
|
||||
this.maxConcurrentConsumers = maxConcurrentConsumers;
|
||||
}
|
||||
// public void setMaxConcurrentConsumers(int maxConcurrentConsumers) {
|
||||
// this.maxConcurrentConsumers = maxConcurrentConsumers;
|
||||
// }
|
||||
|
||||
public void setMaxMessagesPerTask(int maxMessagesPerTask) {
|
||||
this.maxMessagesPerTask = maxMessagesPerTask;
|
||||
@@ -380,20 +385,28 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
container.setSessionAcknowledgeMode(this.sessionAcknowledgeMode);
|
||||
container.setSessionTransacted(this.sessionTransacted);
|
||||
container.setSubscriptionDurable(this.subscriptionDurable);
|
||||
|
||||
int[] conc = parseConcurrency(concurrency);
|
||||
|
||||
if (container instanceof DefaultMessageListenerContainer) {
|
||||
DefaultMessageListenerContainer dmlc = (DefaultMessageListenerContainer) container;
|
||||
if (this.cacheLevelName != null) {
|
||||
dmlc.setCacheLevelName(this.cacheLevelName);
|
||||
}
|
||||
if (this.concurrentConsumers != null) {
|
||||
dmlc.setConcurrentConsumers(this.concurrentConsumers);
|
||||
|
||||
if (conc != null) {
|
||||
if (containerType.isAssignableFrom(DefaultMessageListenerContainer.class)) {
|
||||
dmlc.setConcurrentConsumers(conc[0]);
|
||||
if (conc.length == 2){
|
||||
dmlc.setMaxConcurrentConsumers(conc[1]);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (this.idleTaskExecutionLimit != null) {
|
||||
dmlc.setIdleTaskExecutionLimit(this.idleTaskExecutionLimit);
|
||||
}
|
||||
if (this.maxConcurrentConsumers != null) {
|
||||
dmlc.setMaxConcurrentConsumers(this.maxConcurrentConsumers);
|
||||
}
|
||||
|
||||
if (this.maxMessagesPerTask != null) {
|
||||
dmlc.setMaxMessagesPerTask(this.maxMessagesPerTask);
|
||||
}
|
||||
@@ -415,8 +428,8 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
}
|
||||
else if (container instanceof SimpleMessageListenerContainer) {
|
||||
SimpleMessageListenerContainer smlc = (SimpleMessageListenerContainer) container;
|
||||
if (this.concurrentConsumers != null) {
|
||||
smlc.setConcurrentConsumers(this.concurrentConsumers);
|
||||
if (conc != null) {
|
||||
smlc.setConcurrentConsumers(conc[0]);
|
||||
}
|
||||
smlc.setPubSubNoLocal(this.pubSubNoLocal);
|
||||
smlc.setTaskExecutor(this.taskExecutor);
|
||||
@@ -466,4 +479,20 @@ public class JmsChannelFactoryBean extends AbstractFactoryBean<AbstractJmsChanne
|
||||
((SubscribableJmsChannel) this.channel).destroy();
|
||||
}
|
||||
}
|
||||
|
||||
private int[] parseConcurrency(String concurrency) {
|
||||
if (!StringUtils.hasText(concurrency)) {
|
||||
return null;
|
||||
}
|
||||
int separatorIndex = concurrency.indexOf('-');
|
||||
if (separatorIndex != -1) {
|
||||
int[] result = new int[2];
|
||||
result[0] = Integer.parseInt(concurrency.substring(0, separatorIndex));
|
||||
result[1] = Integer.parseInt(concurrency.substring(separatorIndex + 1, concurrency.length()));
|
||||
return result;
|
||||
}
|
||||
else {
|
||||
return new int[] {1, Integer.parseInt(concurrency)};
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -32,6 +32,7 @@ import org.springframework.util.StringUtils;
|
||||
* Spring Integration JMS namespace.
|
||||
*
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakusky
|
||||
* @since 2.0
|
||||
*/
|
||||
public class JmsChannelParser extends AbstractChannelParser {
|
||||
@@ -108,16 +109,8 @@ public class JmsChannelParser extends AbstractChannelParser {
|
||||
builder.addPropertyValue("sessionAcknowledgeMode", acknowledgeMode);
|
||||
}
|
||||
}
|
||||
int[] concurrency = parseConcurrency(element, parserContext);
|
||||
if (concurrency != null) {
|
||||
if (containerType.startsWith("default")) {
|
||||
builder.addPropertyValue("concurrentConsumers", concurrency[0]);
|
||||
builder.addPropertyValue("maxConcurrentConsumers", concurrency[1]);
|
||||
}
|
||||
else {
|
||||
builder.addPropertyValue("concurrentConsumers", concurrency[1]);
|
||||
}
|
||||
}
|
||||
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "concurrency");
|
||||
|
||||
String prefetch = element.getAttribute("prefetch");
|
||||
if (StringUtils.hasText(prefetch)) {
|
||||
if (containerType.startsWith("default")) {
|
||||
@@ -178,29 +171,4 @@ public class JmsChannelParser extends AbstractChannelParser {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private int[] parseConcurrency(Element ele, ParserContext parserContext) {
|
||||
String concurrency = ele.getAttribute("concurrency");
|
||||
if (!StringUtils.hasText(concurrency)) {
|
||||
return null;
|
||||
}
|
||||
try {
|
||||
int separatorIndex = concurrency.indexOf('-');
|
||||
if (separatorIndex != -1) {
|
||||
int[] result = new int[2];
|
||||
result[0] = Integer.parseInt(concurrency.substring(0, separatorIndex));
|
||||
result[1] = Integer.parseInt(concurrency.substring(separatorIndex + 1, concurrency.length()));
|
||||
return result;
|
||||
}
|
||||
else {
|
||||
return new int[] {1, Integer.parseInt(concurrency)};
|
||||
}
|
||||
}
|
||||
catch (NumberFormatException ex) {
|
||||
parserContext.getReaderContext().error("Invalid concurrency value [" + concurrency + "]: only " +
|
||||
"single maximum integer (e.g. \"5\") and minimum-maximum combo (e.g. \"3-5\") supported.", ele, ex);
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -2,10 +2,13 @@
|
||||
<beans xmlns="http://www.springframework.org/schema/beans"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xmlns:jms="http://www.springframework.org/schema/integration/jms"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/beans
|
||||
http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/integration/jms
|
||||
http://www.springframework.org/schema/integration/jms/spring-integration-jms.xsd">
|
||||
xmlns:context="http://www.springframework.org/schema/context"
|
||||
xsi:schemaLocation="http://www.springframework.org/schema/jms http://www.springframework.org/schema/jms/spring-jms-3.0.xsd
|
||||
http://www.springframework.org/schema/integration/jms http://www.springframework.org/schema/integration/jms/spring-integration-jms.xsd
|
||||
http://www.springframework.org/schema/beans http://www.springframework.org/schema/beans/spring-beans.xsd
|
||||
http://www.springframework.org/schema/context http://www.springframework.org/schema/context/spring-context-3.0.xsd">
|
||||
|
||||
<context:property-placeholder location="classpath:org/springframework/integration/jms/config/channel.properties"/>
|
||||
|
||||
<jms:channel id="queueReferenceChannel" queue="testQueue"/>
|
||||
|
||||
@@ -56,5 +59,8 @@
|
||||
</bean>
|
||||
|
||||
<alias name="connectionFactory" alias="connFact"/>
|
||||
|
||||
|
||||
<jms:channel id="withPlaceholders" queue="${queue}" concurrency="${concurrency}"/>
|
||||
|
||||
</beans>
|
||||
|
||||
@@ -38,6 +38,7 @@ import org.springframework.integration.channel.ChannelInterceptor;
|
||||
import org.springframework.integration.channel.interceptor.ChannelInterceptorAdapter;
|
||||
import org.springframework.integration.jms.PollableJmsChannel;
|
||||
import org.springframework.integration.jms.SubscribableJmsChannel;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.jms.core.JmsTemplate;
|
||||
import org.springframework.jms.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.jms.listener.DefaultMessageListenerContainer;
|
||||
@@ -69,7 +70,10 @@ public class JmsChannelParserTests {
|
||||
|
||||
@Autowired
|
||||
private MessageChannel topicNameChannel;
|
||||
|
||||
|
||||
@Autowired
|
||||
private MessageChannel withPlaceholders;
|
||||
|
||||
@Autowired
|
||||
private MessageChannel topicNameWithResolverChannel;
|
||||
|
||||
@@ -218,6 +222,14 @@ public class JmsChannelParserTests {
|
||||
JmsTemplate jmsTemplate = (JmsTemplate) accessor.getPropertyValue("jmsTemplate");
|
||||
assertEquals("foo", jmsTemplate.getDefaultDestinationName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void withPlaceholders() {
|
||||
DefaultMessageListenerContainer container = TestUtils.getPropertyValue(withPlaceholders, "container", DefaultMessageListenerContainer.class);
|
||||
System.out.println(container.getDestination());
|
||||
System.out.println(container.getConcurrentConsumers());
|
||||
System.out.println(container.getMaxConcurrentConsumers());
|
||||
}
|
||||
|
||||
|
||||
static class TestDestinationResolver implements DestinationResolver {
|
||||
|
||||
@@ -0,0 +1,2 @@
|
||||
queue=testQueue
|
||||
concurrency=5-25
|
||||
Reference in New Issue
Block a user