From 31f02523d747e628172dc2765c9d1be53bbdc420 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 9 Apr 2019 11:34:06 -0400 Subject: [PATCH] GH-212: Configurable anonymous queue name prefix Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/212 Enable configuration of the prefix part of anonymous auto-delete queues. --- docs/src/main/asciidoc/overview.adoc | 6 ++++++ .../properties/RabbitConsumerProperties.java | 13 ++++++++++++ .../RabbitExchangeQueueProvisioner.java | 14 +++++++++---- .../binder/rabbit/RabbitBinderTests.java | 20 +++++++++++++++++++ 4 files changed, 49 insertions(+), 4 deletions(-) diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 6ccf3af80..80eed476e 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -122,6 +122,12 @@ acknowledgeMode:: The acknowledge mode. + Default: `AUTO`. +anonymousGroupPrefix:: +When the binding has no `group` property, an anonymous, auto-delete queue is bound to the destination exchange. +The default naming stragegy for such queues results in a queue named `anonymous.`. +Set this property to change the prefix to something other than the default. ++ +Default: `anonymous.`. autoBindDlq:: Whether to automatically declare the DLQ and bind it to the binder DLX. + diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java index 28310c45d..8900f4dd3 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitConsumerProperties.java @@ -120,6 +120,11 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { */ private ContainerType containerType = ContainerType.SIMPLE; + /** + * Prefix for anonymous queue names (when no group is provided). + */ + private String anonymousGroupPrefix = "anonymous."; + public boolean isTransacted() { return transacted; } @@ -286,4 +291,12 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { this.containerType = containerType; } + public String getAnonymousGroupPrefix() { + return this.anonymousGroupPrefix; + } + + public void setAnonymousGroupPrefix(String anonymousGroupPrefix) { + this.anonymousGroupPrefix = anonymousGroupPrefix; + } + } 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 c35897982..7e066e232 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 @@ -66,8 +66,6 @@ public class RabbitExchangeQueueProvisioner ProvisioningProvider, ExtendedProducerProperties> { // @checkstyle:on - private static final Base64UrlNamingStrategy ANONYMOUS_GROUP_NAME_GENERATOR = new Base64UrlNamingStrategy( - "anonymous."); /** * The delimiter between a group and index when constructing a binder @@ -163,15 +161,23 @@ public class RabbitExchangeQueueProvisioner private ConsumerDestination doProvisionConsumerDestination(String name, String group, ExtendedConsumerProperties properties) { + boolean anonymous = !StringUtils.hasText(group); + Base64UrlNamingStrategy anonQueueNameGenerator = null; + if (anonymous) { + anonQueueNameGenerator = new Base64UrlNamingStrategy( + properties.getExtension().getAnonymousGroupPrefix() == null + ? "" + : properties.getExtension().getAnonymousGroupPrefix()); + } String baseQueueName; if (properties.getExtension().isQueueNameGroupOnly()) { - baseQueueName = anonymous ? ANONYMOUS_GROUP_NAME_GENERATOR.generateName() + baseQueueName = anonymous ? anonQueueNameGenerator.generateName() : group; } else { baseQueueName = groupedName(name, - anonymous ? ANONYMOUS_GROUP_NAME_GENERATOR.generateName() : group); + anonymous ? anonQueueNameGenerator.generateName() : group); } if (this.logger.isInfoEnabled()) { this.logger.info("declaring queue for inbound: " + baseQueueName 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 be97ebced..ffe148389 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 @@ -461,6 +461,26 @@ public class RabbitBinderTests extends assertThat(container.isRunning()).isFalse(); } + @Test + public void testAnonWithBuiltInExchangeCustomPrefix() throws Exception { + RabbitTestBinder binder = getBinder(); + ExtendedConsumerProperties properties = createConsumerProperties(); + properties.getExtension().setDeclareExchange(false); + properties.getExtension().setQueueNameGroupOnly(true); + properties.getExtension().setAnonymousGroupPrefix("customPrefix."); + + Binding consumerBinding = binder.bindConsumer("amq.topic", null, + createBindableChannel("input", new BindingProperties()), properties); + Lifecycle endpoint = extractEndpoint(consumerBinding); + SimpleMessageListenerContainer container = TestUtils.getPropertyValue(endpoint, + "messageListenerContainer", SimpleMessageListenerContainer.class); + String queueName = container.getQueueNames()[0]; + assertThat(queueName).startsWith("customPrefix."); + assertThat(container.isRunning()).isTrue(); + consumerBinding.unbind(); + assertThat(container.isRunning()).isFalse(); + } + @SuppressWarnings("deprecation") @Test public void testConsumerPropertiesWithUserInfrastructureCustomExchangeAndRK()