diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java index 5aec030814..1b309906de 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java @@ -245,11 +245,11 @@ public class KafkaMessageListenerContainer implements SmartLifecycle { /** * The maximum number of messages that are buffered by each concurrent {@link MessageListener} runner. * Increasing the value may increase throughput, but also increases the memory consumption. - * Must be a power of 2. + * Must be a positive number and a power of 2. * @param queueSize the queue size */ public void setQueueSize(int queueSize) { - Assert.isTrue(Integer.bitCount(queueSize) == 1, "'queueSize' must be a power of 2"); + Assert.isTrue(queueSize > 0 && Integer.bitCount(queueSize) == 1, "'queueSize' must be a positive number and a power of 2"); this.queueSize = queueSize; } diff --git a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java index 6c811abfb8..935f38537e 100644 --- a/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java +++ b/spring-integration-kafka/src/test/java/org/springframework/integration/kafka/listener/KafkaMessageDrivenChannelAdapterTests.java @@ -89,6 +89,15 @@ public class KafkaMessageDrivenChannelAdapterTests extends AbstractMessageListen catch (IllegalArgumentException e) { assertThat(e.getMessage(), containsString("power of 2")); } + + try { + kafkaMessageListenerContainer.setQueueSize(Integer.MIN_VALUE); + fail("expected exception"); + } + catch (IllegalArgumentException e) { + assertThat(e.getMessage(), containsString("positive number")); + } + kafkaMessageListenerContainer.setQueueSize(1024); int expectedMessageCount = 100;