From 628707ee32c0d5b78e2b15706e3bc55cb40cfc91 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 6 Apr 2012 11:37:27 -0400 Subject: [PATCH] AMQP-223 Fix Consumer Thread Management If doStart() was called multiple times, multiple threads ran in each consumer. There is a check to prevent creating multiple consumers in this case, but the thread management had no such check. --- .../SimpleMessageListenerContainer.java | 17 ++++++++++++++--- 1 file changed, 14 insertions(+), 3 deletions(-) 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 a7c8b034..36ab6f51 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 @@ -292,9 +292,17 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta protected void doStart() throws Exception { super.doStart(); synchronized (this.consumersMonitor) { - initializeConsumers(); + int newConsumers = initializeConsumers(); if (this.consumers == null) { - logger.info("Consumers were initialized and then cleared (presumably the container was stopped concurrently)"); + if (logger.isInfoEnabled()) { + logger.info("Consumers were initialized and then cleared (presumably the container was stopped concurrently)"); + } + return; + } + if (newConsumers <= 0) { + if (logger.isInfoEnabled()) { + logger.info("Consumers are already running"); + } return; } Set processors = new HashSet(); @@ -343,7 +351,8 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta } - protected void initializeConsumers() { + protected int initializeConsumers() { + int count = 0; synchronized (this.consumersMonitor) { if (this.consumers == null) { cancellationLock.reset(); @@ -351,9 +360,11 @@ public class SimpleMessageListenerContainer extends AbstractMessageListenerConta for (int i = 0; i < this.concurrentConsumers; i++) { BlockingQueueConsumer consumer = createBlockingQueueConsumer(); this.consumers.add(consumer); + count++; } } } + return count; } protected boolean isChannelLocallyTransacted(Channel channel) {