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 298ebdf49..75f6cc8ba 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 @@ -103,6 +103,12 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { */ private Long failedDeclarationRetryInterval; + /** + * Used to create the consumer tags; will be appended by '#n' where 'n' increments for + * each consumer created. + */ + private String consumerTagPrefix; + public boolean isTransacted() { return transacted; } @@ -245,4 +251,12 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { this.failedDeclarationRetryInterval = failedDeclarationRetryInterval; } + public String getConsumerTagPrefix() { + return this.consumerTagPrefix; + } + + public void setConsumerTagPrefix(String consumerTagPrefix) { + this.consumerTagPrefix = consumerTagPrefix; + } + } 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 2442eb016..6025c2006 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 @@ -130,6 +130,11 @@ Whether to bind the queue to the destination exchange. Set it to `false` if you have set up your own infrastructure and have previously created and bound the queue. + Default: `true`. +consumerTagPrefix:: +Used to create the consumer tag(s); will be appended by `#n` where `n` increments for each consumer created. +Example: `${spring.application.name}-${spring.cloud.stream.bindings.input.group}-${spring.cloud.stream.instance-index}`. ++ +Default: none - the broker will generate random consumer tags. deadLetterQueueName:: The name of the DLQ + diff --git a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java index d8980563a..c20d8e516 100644 --- a/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java +++ b/spring-cloud-stream-binder-rabbit/src/main/java/org/springframework/cloud/stream/binder/rabbit/RabbitMessageChannelBinder.java @@ -22,6 +22,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; import org.springframework.amqp.AmqpRejectAndDontRequeueException; import org.springframework.amqp.core.Message; @@ -397,7 +398,12 @@ public class RabbitMessageChannelBinder else if (getApplicationContext() != null) { listenerContainer.setApplicationEventPublisher(getApplicationContext()); } - this.getContainerCustomizer().configure(listenerContainer, consumerDestination.getName(), group); + getContainerCustomizer().configure(listenerContainer, consumerDestination.getName(), group); + if (StringUtils.hasText(properties.getExtension().getConsumerTagPrefix())) { + final AtomicInteger index = new AtomicInteger(); + listenerContainer.setConsumerTagStrategy(q -> + properties.getExtension().getConsumerTagPrefix() + "#" + index.getAndIncrement()); + } listenerContainer.afterPropertiesSet(); AmqpInboundChannelAdapter adapter = new AmqpInboundChannelAdapter(listenerContainer); 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 e5de8a3ab..5629bb726 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 @@ -431,6 +431,8 @@ public class RabbitBinderTests extends extProps.setDlqOverflowBehavior("reject-publish"); extProps.setDlqMaxPriority(8); extProps.setDlqTtl(1_000); + extProps.setConsumerTagPrefix("testConsumerTag"); + extProps.setExclusive(true); Binding consumerBinding = binder.bindConsumer("propsUser3", "infra", createBindableChannel("input", new BindingProperties()), properties); @@ -438,8 +440,6 @@ public class RabbitBinderTests extends SimpleMessageListenerContainer container = TestUtils.getPropertyValue(endpoint, "messageListenerContainer", SimpleMessageListenerContainer.class); assertThat(container.isRunning()).isTrue(); - consumerBinding.unbind(); - assertThat(container.isRunning()).isFalse(); RabbitManagementTemplate rmt = new RabbitManagementTemplate(); List bindings = rmt.getBindingsForExchange("/", "propsUser3"); int n = 0; @@ -462,10 +462,9 @@ public class RabbitBinderTests extends assertThat(exchange.isDurable()).isEqualTo(false); assertThat(exchange.isAutoDelete()).isEqualTo(true); -// Queue queue = rmt.getQueue("propsUser3"); AMQP-698 QueueInfo queue = rmt.getClient().getQueue("/", "propsUser3.infra"); n = 0; - while (n++ < 100 && queue == null) { + while (n++ < 100 && queue == null || queue.getConsumerCount() == 0) { Thread.sleep(100); queue = rmt.getClient().getQueue("/", "propsUser3.infra"); } @@ -480,6 +479,7 @@ public class RabbitBinderTests extends assertThat(args.get("x-dead-letter-exchange")).isEqualTo("customDLX"); assertThat(args.get("x-dead-letter-routing-key")).isEqualTo("customDLRK"); assertThat(args.get("x-queue-mode")).isEqualTo("lazy"); + assertThat(queue.getExclusiveConsumerTag()).isEqualTo("testConsumerTag#0"); queue = rmt.getClient().getQueue("/", "customDLQ"); @@ -499,6 +499,9 @@ public class RabbitBinderTests extends assertThat(args.get("x-dead-letter-exchange")).isEqualTo("propsUser3"); assertThat(args.get("x-dead-letter-routing-key")).isEqualTo("propsUser3"); assertThat(args.get("x-queue-mode")).isEqualTo("lazy"); + + consumerBinding.unbind(); + assertThat(container.isRunning()).isFalse(); } @Test