From 5c2a764a2ff846b0a824459f2760d663c068cafb Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 17 Jul 2017 11:01:24 -0400 Subject: [PATCH] GH-58, 80: Exclusive Consumer and Lazy Queues Resolves #58 Resolves #80 Add properties to support exclusive consumers and lazy queues. --- .../properties/RabbitCommonProperties.java | 26 ++++++++++++++++ .../properties/RabbitConsumerProperties.java | 14 +++++++++ .../RabbitExchangeQueueProvisioner.java | 10 +++++-- .../src/main/asciidoc/overview.adoc | 30 +++++++++++++++++++ .../rabbit/RabbitMessageChannelBinder.java | 1 + .../binder/rabbit/RabbitBinderTests.java | 16 ++++++---- 6 files changed, 89 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 e84bc3676..45a132fdd 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 @@ -149,6 +149,16 @@ public abstract class RabbitCommonProperties { */ private String prefix = ""; + /** + * True if the queue is provisioned as a lazy queue. + */ + private boolean lazy; + + /** + * True if the DLQ is provisioned as a lazy queue. + */ + private boolean dlqLazy; + public String getExchangeType() { return this.exchangeType; } @@ -342,4 +352,20 @@ public abstract class RabbitCommonProperties { this.prefix = prefix; } + public boolean isLazy() { + return this.lazy; + } + + public void setLazy(boolean lazy) { + this.lazy = lazy; + } + + public boolean isDlqLazy() { + return this.dlqLazy; + } + + public void setDlqLazy(boolean dlqLazy) { + this.dlqLazy = dlqLazy; + } + } 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 9b8ba62ef..618cb877c 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 @@ -50,6 +50,11 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { private long recoveryInterval = 5000; + /** + * True if the consumer is exclusive. + */ + private boolean exclusive; + public boolean isTransacted() { return transacted; } @@ -159,4 +164,13 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { public void setRecoveryInterval(long recoveryInterval) { this.recoveryInterval = recoveryInterval; } + + public boolean isExclusive() { + return this.exclusive; + } + + public void setExclusive(boolean exclusive) { + this.exclusive = exclusive; + } + } 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 6788f6cdd..63f5b8d21 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 @@ -349,7 +349,7 @@ public class RabbitExchangeQueueProvisioner implements ProvisioningProvider args, Integer expires, Integer maxLength, Integer maxLengthBytes, - Integer maxPriority, Integer ttl) { + Integer maxPriority, Integer ttl, boolean lazy) { if (expires != null) { args.put("x-expires", expires); } @@ -381,6 +382,9 @@ public class RabbitExchangeQueueProvisioner implements ProvisioningProvider properties = createConsumerProperties(); properties.getExtension().setRequeueRejected(true); properties.getExtension().setTransacted(true); + properties.getExtension().setExclusive(true); Binding consumerBinding = binder.bindConsumer("props.0", null, createBindableChannel("input", new BindingProperties()), properties); Lifecycle endpoint = extractEndpoint(consumerBinding); @@ -172,6 +173,7 @@ public class RabbitBinderTests extends assertThat(container.getAcknowledgeMode()).isEqualTo(AcknowledgeMode.AUTO); assertThat(container.getQueueNames()[0]).startsWith(properties.getExtension().getPrefix()); assertThat(TestUtils.getPropertyValue(container, "transactional", Boolean.class)).isTrue(); + assertThat(TestUtils.getPropertyValue(container, "exclusive", Boolean.class)).isTrue(); assertThat(TestUtils.getPropertyValue(container, "concurrentConsumers")).isEqualTo(1); assertThat(TestUtils.getPropertyValue(container, "maxConcurrentConsumers")).isNull(); assertThat(TestUtils.getPropertyValue(container, "defaultRequeueRejected", Boolean.class)).isTrue(); @@ -291,6 +293,7 @@ public class RabbitBinderTests extends extProps.setExchangeAutoDelete(true); extProps.setBindingRoutingKey("foo"); extProps.setExpires(30_000); + extProps.setLazy(true); extProps.setMaxLength(10_000); extProps.setMaxLengthBytes(100_000); extProps.setMaxPriority(10); @@ -302,6 +305,7 @@ public class RabbitBinderTests extends extProps.setDlqDeadLetterExchange("propsUser3"); extProps.setDlqDeadLetterRoutingKey("propsUser3"); extProps.setDlqExpires(60_000); + extProps.setDlqLazy(true); extProps.setDlqMaxLength(20_000); extProps.setDlqMaxLengthBytes(40_000); extProps.setDlqMaxPriority(8); @@ -353,6 +357,7 @@ public class RabbitBinderTests extends assertThat(args.get("x-message-ttl")).isEqualTo(2_000); 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"); queue = rmt.getClient().getQueue("/", "customDLQ"); @@ -370,6 +375,7 @@ public class RabbitBinderTests extends assertThat(args.get("x-message-ttl")).isEqualTo(1_000); 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"); } @Test