From e9ad73ad2c2e465670ac68016e24f47d6c22b172 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Fri, 4 Jul 2014 15:45:40 +0300 Subject: [PATCH] INT-3434: AMQP Don't Declare Queue if Supplied JIRA: https://jira.spring.io/browse/INT-3434 Previously `PointToPointSubscribableAmqpChannel` declared Queue always. Change the logic to use `admin.declareQueue` only if `queueName` isn't specified. In addition it is redundant to get deal with entire `Queue` object - just use the `queueName` for the `SimpleMessageListenerContainer` Polishing Declare queue the old way if not already present in broker, even if the name is supplied. --- .../AbstractSubscribableAmqpChannel.java | 20 ++++++----- .../PointToPointSubscribableAmqpChannel.java | 14 ++++---- .../channel/PublishSubscribeAmqpChannel.java | 8 ++--- .../amqp/channel/ChannelTests.java | 35 +++++++++++++++++++ .../DispatcherHasNoSubscribersTests.java | 26 ++++++-------- 5 files changed, 69 insertions(+), 34 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java index c8efbd1b90..8a43f0c4be 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/AbstractSubscribableAmqpChannel.java @@ -22,7 +22,6 @@ import org.apache.commons.logging.LogFactory; import org.springframework.amqp.core.AmqpAdmin; import org.springframework.amqp.core.AmqpTemplate; import org.springframework.amqp.core.MessageListener; -import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitAdmin; import org.springframework.amqp.rabbit.core.RabbitTemplate; @@ -49,7 +48,8 @@ import org.springframework.util.Assert; * @author Artem Bilan * @since 2.1 */ -abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel implements SubscribableChannel, SmartLifecycle, DisposableBean { +abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel + implements SubscribableChannel, SmartLifecycle, DisposableBean { private final String channelName; @@ -65,7 +65,8 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple private final ConnectionFactory connectionFactory; - public AbstractSubscribableAmqpChannel(String channelName, SimpleMessageListenerContainer container, AmqpTemplate amqpTemplate) { + public AbstractSubscribableAmqpChannel(String channelName, SimpleMessageListenerContainer container, + AmqpTemplate amqpTemplate) { this(channelName, container, amqpTemplate, false); } @@ -123,8 +124,8 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple Integer.class); } this.setMaxSubscribers(this.maxSubscribers); - Queue queue = this.initializeQueue(this.admin, this.channelName); - this.container.setQueues(queue); + String queue = this.obtainQueueName(this.admin, this.channelName); + this.container.setQueueNames(queue); MessageConverter converter = (this.getAmqpTemplate() instanceof RabbitTemplate) ? ((RabbitTemplate) this.getAmqpTemplate()).getMessageConverter() : new SimpleMessageConverter(); @@ -139,7 +140,7 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple protected abstract AbstractDispatcher createDispatcher(); - protected abstract Queue initializeQueue(AmqpAdmin admin, String channelName); + protected abstract String obtainQueueName(AmqpAdmin admin, String channelName); private static class DispatchingMessageListener implements MessageListener { @@ -198,7 +199,8 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple } } catch (Exception e) { - throw new MessagingException("Failure occured in AMQP listener while attempting to convert and dispatch Message.", e); + throw new MessagingException("Failure occured in AMQP listener " + + "while attempting to convert and dispatch Message.", e); } } } @@ -210,7 +212,7 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple @Override public boolean isAutoStartup() { - return (this.container != null) ? this.container.isAutoStartup() : false; + return (this.container != null) && this.container.isAutoStartup(); } @Override @@ -220,7 +222,7 @@ abstract class AbstractSubscribableAmqpChannel extends AbstractAmqpChannel imple @Override public boolean isRunning() { - return (this.container != null) ? this.container.isRunning() : false; + return (this.container != null) && this.container.isRunning(); } @Override diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java index bbf8815887..a16d5794f3 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PointToPointSubscribableAmqpChannel.java @@ -26,6 +26,7 @@ import org.springframework.integration.dispatcher.UnicastingDispatcher; /** * @author Mark Fisher + * @author Artem Bilan * @since 2.1 */ public class PointToPointSubscribableAmqpChannel extends AbstractSubscribableAmqpChannel { @@ -33,7 +34,8 @@ public class PointToPointSubscribableAmqpChannel extends AbstractSubscribableAmq private volatile String queueName; - public PointToPointSubscribableAmqpChannel(String channelName, SimpleMessageListenerContainer container, AmqpTemplate amqpTemplate) { + public PointToPointSubscribableAmqpChannel(String channelName, SimpleMessageListenerContainer container, + AmqpTemplate amqpTemplate) { super(channelName, container, amqpTemplate); } @@ -41,7 +43,6 @@ public class PointToPointSubscribableAmqpChannel extends AbstractSubscribableAmq /** * Provide a Queue name to be used. If this is not provided, * the Queue's name will be the same as the channel name. - * * @param queueName The queue name. */ public void setQueueName(String queueName) { @@ -49,13 +50,14 @@ public class PointToPointSubscribableAmqpChannel extends AbstractSubscribableAmq } @Override - protected Queue initializeQueue(AmqpAdmin admin, String channelName) { + protected String obtainQueueName(AmqpAdmin admin, String channelName) { if (this.queueName == null) { this.queueName = channelName; } - Queue queue = new Queue(this.queueName); - admin.declareQueue(queue); - return queue; + if (admin.getQueueProperties(this.queueName) == null) { + admin.declareQueue(new Queue(this.queueName)); + } + return this.queueName; } @Override diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java index ab4dbfed38..bc9512c389 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/channel/PublishSubscribeAmqpChannel.java @@ -46,7 +46,8 @@ public class PublishSubscribeAmqpChannel extends AbstractSubscribableAmqpChannel private volatile boolean initialized; - public PublishSubscribeAmqpChannel(String channelName, SimpleMessageListenerContainer container, AmqpTemplate amqpTemplate) { + public PublishSubscribeAmqpChannel(String channelName, SimpleMessageListenerContainer container, + AmqpTemplate amqpTemplate) { super(channelName, container, amqpTemplate, true); } @@ -56,7 +57,6 @@ public class PublishSubscribeAmqpChannel extends AbstractSubscribableAmqpChannel * FanoutExchange will be declared implicitly, and its name will be the same * as the channel name prefixed by "si.fanout.". In either case, an effectively * anonymous Queue will be declared automatically. - * * @param exchange The fanout exchange. */ public void setExchange(FanoutExchange exchange) { @@ -64,7 +64,7 @@ public class PublishSubscribeAmqpChannel extends AbstractSubscribableAmqpChannel } @Override - protected Queue initializeQueue(AmqpAdmin admin, String channelName) { + protected String obtainQueueName(AmqpAdmin admin, String channelName) { if (this.exchange == null) { String exchangeName = "si.fanout." + channelName; this.exchange = new FanoutExchange(exchangeName); @@ -80,7 +80,7 @@ public class PublishSubscribeAmqpChannel extends AbstractSubscribableAmqpChannel } } this.initialized = true; - return this.queue; + return this.queue.getName(); } private void doDeclares() { diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/ChannelTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/ChannelTests.java index 23415437d4..9ea3bad92c 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/ChannelTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/ChannelTests.java @@ -16,6 +16,8 @@ package org.springframework.integration.amqp.channel; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; +import static org.mockito.Mockito.mock; import java.util.Collection; import java.util.concurrent.CyclicBarrier; @@ -25,7 +27,12 @@ import org.junit.ClassRule; import org.junit.Test; import org.junit.runner.RunWith; +import org.springframework.amqp.core.AmqpTemplate; +import org.springframework.amqp.core.Queue; import org.springframework.amqp.rabbit.connection.CachingConnectionFactory; +import org.springframework.amqp.rabbit.core.RabbitAdmin; +import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; +import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.integration.amqp.rule.BrokerRunning; @@ -82,4 +89,32 @@ public class ChannelTests { assertEquals(0, TestUtils.getPropertyValue(factory, "connectionListener.delegates", Collection.class).size()); } + /* + * Verify queue is declared if not present and not declared if it is already present. + */ + @Test + public void channelDeclarationTests() { + RabbitAdmin admin = new RabbitAdmin(this.factory); + admin.deleteQueue("implicit"); + SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(this.factory); + AmqpTemplate amqpTemplate = mock(AmqpTemplate.class); + PointToPointSubscribableAmqpChannel channel = new PointToPointSubscribableAmqpChannel("implicit", container, + amqpTemplate); + channel.setBeanFactory(mock(BeanFactory.class)); + channel.afterPropertiesSet(); + assertNotNull(admin.getQueueProperties("implicit")); + admin.deleteQueue("implicit"); + + admin.deleteQueue("explicit"); + channel.setQueueName("explicit"); + channel.afterPropertiesSet(); + assertNotNull(admin.getQueueProperties("explicit")); + + admin.deleteQueue("explicit"); + admin.declareQueue(new Queue("explicit", false)); // verify no declaration if exists with non-standard props + channel.afterPropertiesSet(); + assertNotNull(admin.getQueueProperties("explicit")); + admin.deleteQueue("explicit"); + } + } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/DispatcherHasNoSubscribersTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/DispatcherHasNoSubscribersTests.java index b01290c3f7..dee08be04b 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/DispatcherHasNoSubscribersTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/channel/DispatcherHasNoSubscribersTests.java @@ -13,23 +13,21 @@ * See the License for the specific language governing permissions and * limitations under the License. */ + package org.springframework.integration.amqp.channel; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertTrue; -import static org.junit.Assert.fail; +import static org.junit.Assert.*; import static org.mockito.Matchers.any; import static org.mockito.Matchers.anyBoolean; import static org.mockito.Matchers.anyString; -import static org.mockito.Mockito.doAnswer; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.when; +import static org.mockito.Mockito.*; import java.util.ArrayList; import java.util.List; import java.util.Map; +import com.rabbitmq.client.AMQP.Queue.DeclareOk; +import com.rabbitmq.client.Channel; import org.apache.commons.logging.Log; import org.junit.Test; import org.mockito.invocation.InvocationOnMock; @@ -48,9 +46,6 @@ import org.springframework.beans.factory.BeanFactory; import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.MessageDeliveryException; -import com.rabbitmq.client.AMQP.Queue.DeclareOk; -import com.rabbitmq.client.Channel; - /** * @author Gary Russell @@ -65,7 +60,8 @@ public class DispatcherHasNoSubscribersTests { final Channel channel = mock(Channel.class); DeclareOk declareOk = mock(DeclareOk.class); when(declareOk.getQueue()).thenReturn("noSubscribersChannel"); - when(channel.queueDeclare(anyString(), anyBoolean(), anyBoolean(), anyBoolean(), any(Map.class))).thenReturn(declareOk); + when(channel.queueDeclare(anyString(), anyBoolean(), anyBoolean(), anyBoolean(), any(Map.class))) + .thenReturn(declareOk); Connection connection = mock(Connection.class); doAnswer(new Answer() { @Override @@ -78,8 +74,8 @@ public class DispatcherHasNoSubscribersTests { container.setConnectionFactory(connectionFactory); AmqpTemplate amqpTemplate = mock(AmqpTemplate.class); - PointToPointSubscribableAmqpChannel amqpChannel = new PointToPointSubscribableAmqpChannel("noSubscribersChannel", - container, amqpTemplate); + PointToPointSubscribableAmqpChannel amqpChannel = + new PointToPointSubscribableAmqpChannel("noSubscribersChannel", container, amqpTemplate); amqpChannel.setBeanName("noSubscribersChannel"); amqpChannel.setBeanFactory(mock(BeanFactory.class)); amqpChannel.afterPropertiesSet(); @@ -112,9 +108,9 @@ public class DispatcherHasNoSubscribersTests { PublishSubscribeAmqpChannel amqpChannel = new PublishSubscribeAmqpChannel("noSubscribersChannel", container, amqpTemplate) { @Override - protected Queue initializeQueue(AmqpAdmin admin, + protected String obtainQueueName(AmqpAdmin admin, String channelName) { - return queue; + return queue.getName(); }}; amqpChannel.setBeanName("noSubscribersChannel"); amqpChannel.setBeanFactory(mock(BeanFactory.class));