From 3c20034cb24480a322802581888f086a4a2dc30f Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 10 Oct 2017 11:29:57 -0400 Subject: [PATCH] GH-93: Consume from existing queue Resolves: https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/93 The Rabbit binder consumes from a queue named `.`. Add a property to omit the `.` part so users can use SCSt to consume from existing queue(s). * Polishing - PR Comments --- .../rabbit/properties/RabbitCommonProperties.java | 13 +++++++++++++ .../RabbitExchangeQueueProvisioner.java | 10 ++++++---- .../src/main/asciidoc/overview.adoc | 11 +++++++++++ .../stream/binder/rabbit/RabbitBinderTests.java | 7 +++++-- .../stream/binder/rabbit/RabbitTestBinder.java | 14 ++++++++++++-- 5 files changed, 47 insertions(+), 8 deletions(-) diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java index 45a132fdd..b44d4f801 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitCommonProperties.java @@ -54,6 +54,11 @@ public abstract class RabbitCommonProperties { */ private boolean delayedExchange = false; + /** + * set to true to name the queue with only the group; default is destination.group + */ + private boolean queueNameGroupOnly = false; + /** * whether to bind a queue (or queues when partitioned) to the exchange */ @@ -199,6 +204,14 @@ public abstract class RabbitCommonProperties { this.delayedExchange = delayedExchange; } + public boolean isQueueNameGroupOnly() { + return this.queueNameGroupOnly; + } + + public void setQueueNameGroupOnly(boolean queueNameGroupOnly) { + this.queueNameGroupOnly = queueNameGroupOnly; + } + public boolean isBindQueue() { return this.bindQueue; } diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java index e964b5229..041a005e4 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/provisioning/RabbitExchangeQueueProvisioner.java @@ -95,7 +95,8 @@ public class RabbitExchangeQueueProvisioner implements ApplicationListener properties) { + public ConsumerDestination provisionConsumerDestination(String name, String group, + ExtendedConsumerProperties properties) { boolean anonymous = !StringUtils.hasText(group); - String baseQueueName = anonymous ? groupedName(name, ANONYMOUS_GROUP_NAME_GENERATOR.generateName()) - : groupedName(name, group); + String baseQueueName = anonymous ? groupedName(name, ANONYMOUS_GROUP_NAME_GENERATOR.generateName()) + : properties.getExtension().isQueueNameGroupOnly() ? group : groupedName(name, group); if (this.logger.isInfoEnabled()) { this.logger.info("declaring queue for inbound: " + baseQueueName + ", bound to: " + name); } diff --git a/spring-cloud-stream-binder-rabbit-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-rabbit-docs/src/main/asciidoc/overview.adoc index 2dfa2cb66..6920094e9 100644 --- a/spring-cloud-stream-binder-rabbit-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-rabbit-docs/src/main/asciidoc/overview.adoc @@ -232,6 +232,11 @@ prefix:: A prefix to be added to the name of the `destination` and queues. + Default: "". +queueNameGroupOnly:: + When true, consume from a queue with a name equal to the `group`; otherwise the queue name is `destination.group`. + This is useful, for example, when using Spring Cloud Stream to consume from an existing RabbitMQ queue. ++ +Default: false. recoveryInterval:: The interval between connection recovery attempts, in milliseconds. + @@ -419,6 +424,12 @@ prefix:: A prefix to be added to the name of the `destination` exchange. + Default: "". +queueNameGroupOnly:: + When true, consume from a queue with a name equal to the `group`; otherwise the queue name is `destination.group`. + This is useful, for example, when using Spring Cloud Stream to consume from an existing RabbitMQ queue. + Only applies if `requiredGroups` are provided and then only to those groups. ++ +Default: false. routingKeyExpression:: A SpEL expression to determine the routing key to use when publishing messages. For a fixed routing key, use a literal expression, e.g. `routingKeyExpression='my.routingKey'` in a properties file, or `routingKeyExpression: '''my.routingKey'''` in a YAML file. diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java index 4b5854776..c309d07e6 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitBinderTests.java @@ -334,9 +334,11 @@ public class RabbitBinderTests extends ExtendedConsumerProperties properties = createConsumerProperties(); properties.getExtension().setExchangeType(ExchangeTypes.DIRECT); properties.getExtension().setBindingRoutingKey("foo"); + properties.getExtension().setQueueNameGroupOnly(true); // properties.getExtension().setDelayedExchange(true); // requires delayed message exchange plugin; tested locally - Binding consumerBinding = binder.bindConsumer("propsUser2", "infra", + String group = "infra"; + Binding consumerBinding = binder.bindConsumer("propsUser2", group, createBindableChannel("input", new BindingProperties()), properties); Lifecycle endpoint = extractEndpoint(consumerBinding); SimpleMessageListenerContainer container = TestUtils.getPropertyValue(endpoint, "messageListenerContainer", @@ -344,6 +346,7 @@ public class RabbitBinderTests extends assertThat(container.isRunning()).isTrue(); consumerBinding.unbind(); assertThat(container.isRunning()).isFalse(); + assertThat(container.getQueueNames()[0]).isEqualTo(group); RabbitManagementTemplate rmt = new RabbitManagementTemplate(); List bindings = rmt.getBindingsForExchange("/", "propsUser2"); int n = 0; @@ -353,7 +356,7 @@ public class RabbitBinderTests extends } assertThat(bindings.size()).isEqualTo(1); assertThat(bindings.get(0).getExchange()).isEqualTo("propsUser2"); - assertThat(bindings.get(0).getDestination()).isEqualTo("propsUser2.infra"); + assertThat(bindings.get(0).getDestination()).isEqualTo(group); assertThat(bindings.get(0).getRoutingKey()).isEqualTo("foo"); // // TODO: AMQP-696 diff --git a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java index e3e289711..6a80b06c8 100644 --- a/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/test/java/org/springframework/cloud/stream/binder/rabbit/RabbitTestBinder.java @@ -75,7 +75,12 @@ public class RabbitTestBinder extends AbstractTestBinder bindConsumer(String name, String group, MessageChannel moduleInputChannel, ExtendedConsumerProperties properties) { if (group != null) { - this.queues.add(properties.getExtension().getPrefix() + name + ("." + group)); + if (properties.getExtension().isQueueNameGroupOnly()) { + this.queues.add(properties.getExtension().getPrefix() + group); + } + else { + this.queues.add(properties.getExtension().getPrefix() + name + ("." + group)); + } } this.exchanges.add(properties.getExtension().getPrefix() + name); this.prefixes.add(properties.getExtension().getPrefix()); @@ -90,7 +95,12 @@ public class RabbitTestBinder extends AbstractTestBinder