diff --git a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java index a3034e2a..8f1a6de7 100644 --- a/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java +++ b/spring-rabbit/src/main/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainer.java @@ -164,19 +164,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } int delta = this.concurrentConsumers - concurrentConsumers; this.concurrentConsumers = concurrentConsumers; - if (isActive() && this.consumers != null) { - if (delta > 0) { - Iterator consumerIterator = this.consumers.iterator(); - while (consumerIterator.hasNext() && delta > 0) { - BlockingQueueConsumer consumer = consumerIterator.next(); - consumer.basicCancel(true); - consumerIterator.remove(); - delta--; - } - } - else { - addAndStartConsumers(-delta); - } + if (isActive()) { + adjustConsumers(delta); } } } @@ -196,7 +185,15 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta "'maxConcurrentConsumers' value must be at least 'concurrentConsumers'"); Assert.isTrue(!isExclusive() || maxConcurrentConsumers == 1, "When the consumer is exclusive, the concurrency must be 1"); + Integer oldMax = this.maxConcurrentConsumers; this.maxConcurrentConsumers = maxConcurrentConsumers; + if (oldMax != null && isActive()) { + int delta = oldMax - maxConcurrentConsumers; + if (delta > 0) { // only decrease, not increase + adjustConsumers(delta); + } + } + } /** @@ -587,10 +584,45 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta return count; } + /** + * Adjust consumers depending on delta. + * @param delta a negative value increases, positive decreases. + * @since 1.7.8 + */ + protected void adjustConsumers(int delta) { + synchronized (this.consumersMonitor) { + if (isActive() && this.consumers != null) { + if (delta > 0) { + Iterator consumerIterator = this.consumers.iterator(); + while (consumerIterator.hasNext() && delta > 0 + && (this.maxConcurrentConsumers == null + || this.consumers.size() > this.maxConcurrentConsumers)) { + BlockingQueueConsumer consumer = consumerIterator.next(); + consumer.basicCancel(true); + consumerIterator.remove(); + delta--; + } + } + else { + addAndStartConsumers(-delta); + } + } + } + } + + + /** + * Start up to delta consumers, limited by {@link #setMaxConcurrentConsumers(int)}. + * @param delta the consumers to add. + */ protected void addAndStartConsumers(int delta) { synchronized (this.consumersMonitor) { if (this.consumers != null) { for (int i = 0; i < delta; i++) { + if (this.maxConcurrentConsumers != null + && this.consumers.size() >= this.maxConcurrentConsumers) { + break; + } BlockingQueueConsumer consumer = createBlockingQueueConsumer(); this.consumers.add(consumer); AsyncMessageProcessingConsumer processor = new AsyncMessageProcessingConsumer(consumer); diff --git a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerLongTests.java b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerLongTests.java index f47155fa..e97b11e4 100644 --- a/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerLongTests.java +++ b/spring-rabbit/src/test/java/org/springframework/amqp/rabbit/listener/SimpleMessageListenerContainerLongTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2017 the original author or authors. + * Copyright 2002-2018 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,7 +16,9 @@ package org.springframework.amqp.rabbit.listener; +import static org.hamcrest.Matchers.equalTo; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThat; import static org.junit.Assert.fail; import java.util.Set; @@ -44,12 +46,23 @@ import com.rabbitmq.client.ConnectionFactory; * @since 1.2.1 * */ -@RabbitAvailable(queues = SimpleMessageListenerContainerLongTests.QUEUE) +@RabbitAvailable(queues = { + SimpleMessageListenerContainerLongTests.QUEUE, + SimpleMessageListenerContainerLongTests.QUEUE2, + SimpleMessageListenerContainerLongTests.QUEUE3, + SimpleMessageListenerContainerLongTests.QUEUE4 +}) @LongRunning public class SimpleMessageListenerContainerLongTests { public static final String QUEUE = "SimpleMessageListenerContainerLongTests.queue"; + public static final String QUEUE2 = "SimpleMessageListenerContainerLongTests.queue2"; + + public static final String QUEUE3 = "SimpleMessageListenerContainerLongTests.queue3"; + + public static final String QUEUE4 = "SimpleMessageListenerContainerLongTests.queue4"; + private final Log logger = LogFactory.getLog(SimpleMessageListenerContainerLongTests.class); private final SingleConnectionFactory connectionFactory; @@ -132,6 +145,90 @@ public class SimpleMessageListenerContainerLongTests { connectionFactory.destroy(); } + @Test + public void testIncreaseMinAtMax() throws Exception { + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.connectionFactory); + container.setStartConsumerMinInterval(100); + container.setConsecutiveActiveTrigger(1); + container.setMessageListener(m -> { + try { + Thread.sleep(50); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + container.setQueueNames(QUEUE2); + container.setConcurrentConsumers(2); + container.setMaxConcurrentConsumers(5); + container.afterPropertiesSet(); + container.start(); + RabbitTemplate template = new RabbitTemplate(this.connectionFactory); + for (int i = 0; i < 20; i++) { + template.convertAndSend(QUEUE2, "foo"); + } + waitForNConsumers(container, 5); + container.setConcurrentConsumers(4); + Set consumers = (Set) TestUtils.getPropertyValue(container, "consumers"); + assertThat(consumers.size(), equalTo(5)); + } + + @Test + public void testDecreaseMinAtMax() throws Exception { + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.connectionFactory); + container.setStartConsumerMinInterval(100); + container.setConsecutiveActiveTrigger(1); + container.setMessageListener(m -> { + try { + Thread.sleep(50); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + container.setQueueNames(QUEUE3); + container.setConcurrentConsumers(2); + container.setMaxConcurrentConsumers(3); + container.afterPropertiesSet(); + container.start(); + RabbitTemplate template = new RabbitTemplate(this.connectionFactory); + for (int i = 0; i < 20; i++) { + template.convertAndSend(QUEUE3, "foo"); + } + waitForNConsumers(container, 3); + container.setConcurrentConsumers(1); + Set consumers = (Set) TestUtils.getPropertyValue(container, "consumers"); + assertThat(consumers.size(), equalTo(3)); + } + + @Test + public void testDecreaseMaxAtMax() throws Exception { + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.connectionFactory); + container.setStartConsumerMinInterval(100); + container.setConsecutiveActiveTrigger(1); + container.setMessageListener(m -> { + try { + Thread.sleep(50); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }); + container.setQueueNames(QUEUE4); + container.setConcurrentConsumers(2); + container.setMaxConcurrentConsumers(3); + container.afterPropertiesSet(); + container.start(); + RabbitTemplate template = new RabbitTemplate(this.connectionFactory); + for (int i = 0; i < 20; i++) { + template.convertAndSend(QUEUE4, "foo"); + } + waitForNConsumers(container, 3); + container.setConcurrentConsumers(1); + container.setMaxConcurrentConsumers(1); + Set consumers = (Set) TestUtils.getPropertyValue(container, "consumers"); + assertThat(consumers.size(), equalTo(1)); + } public void handleMessage(String foo) { logger.info(foo);