From 86b8a3ba5d676f34d4475e932490888a9cc541db Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 22 Oct 2019 15:01:53 -0400 Subject: [PATCH] GH-242: Support multiple queue bindings Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-rabbit/issues/242 Allow a consumer binding to bind its queue with multiple routing keys. Add `bindingQueueDelimiter` to treat `bindingRoutingKey` as a list. Resolves #266 --- docs/src/main/asciidoc/overview.adoc | 15 ++++- .../properties/RabbitCommonProperties.java | 13 ++++ .../RabbitExchangeQueueProvisioner.java | 60 +++++++++++++++---- .../binder/rabbit/RabbitBinderTests.java | 11 +++- 4 files changed, 82 insertions(+), 17 deletions(-) diff --git a/docs/src/main/asciidoc/overview.adoc b/docs/src/main/asciidoc/overview.adoc index 3a5105c30..2518a9bb9 100644 --- a/docs/src/main/asciidoc/overview.adoc +++ b/docs/src/main/asciidoc/overview.adoc @@ -139,9 +139,14 @@ Whether to automatically declare the DLQ and bind it to the binder DLX. Default: `false`. bindingRoutingKey:: The routing key with which to bind the queue to the exchange (if `bindQueue` is `true`). -For partitioned destinations, `-` is appended. +Can be multiple keys - see `bindingRoutingKeyDelimiter`. +For partitioned destinations, `-` is appended to each key. + Default: `#`. +bindingRoutingKeyDelimiter:: +When this is not null, 'bindingRoutingKey' is considered to be a list of keys delimited by this value; often a comma is used. ++ +Default: `null`. bindQueue:: Whether to declare the queue and bind it to the destination exchange. Set it to `false` if you have set up your own infrastructure and have previously created and bound the queue. @@ -420,10 +425,16 @@ The batch timeout when batching is enabled. Default: `5000`. bindingRoutingKey:: The routing key with which to bind the queue to the exchange (if `bindQueue` is `true`). -Only applies to non-partitioned destinations. +Can be multiple keys - see `bindingRoutingKeyDelimiter`. +For partitioned destinations, `-n` is appended to each key. Only applies if `requiredGroups` are provided and then only to those groups. + Default: `#`. +bindingRoutingKeyDelimiter:: +When this is not null, 'bindingRoutingKey' is considered to be a list of keys delimited by this value; often a comma is used. +Only applies if `requiredGroups` are provided and then only to those groups. ++ +Default: `null`. bindQueue:: Whether to declare the queue and bind it to the destination exchange. Set it to `false` if you have set up your own infrastructure and have previously created and bound the queue. 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 dd5317199..91874014e 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 @@ -77,6 +77,11 @@ public abstract class RabbitCommonProperties { */ private String bindingRoutingKey; + /** + * when not null, treat 'bindingRoutingKey' as a delimited list of keys to bind. + */ + private String bindingRoutingKeyDelimiter; + /** * default time to live to apply to the queue when declared (ms). */ @@ -270,6 +275,14 @@ public abstract class RabbitCommonProperties { this.bindingRoutingKey = routingKey; } + public String getBindingRoutingKeyDelimiter() { + return this.bindingRoutingKeyDelimiter; + } + + public void setBindingRoutingKeyDelimiter(String bindingRoutingKeyDelimiter) { + this.bindingRoutingKeyDelimiter = bindingRoutingKeyDelimiter; + } + public Integer getTtl() { return this.ttl; } 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 5d5fad262..9e1c0ad8c 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 @@ -53,6 +53,7 @@ import org.springframework.cloud.stream.provisioning.ProvisioningProvider; import org.springframework.context.ApplicationListener; import org.springframework.context.support.GenericApplicationContext; import org.springframework.util.Assert; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -112,8 +113,16 @@ public class RabbitExchangeQueueProvisioner Queue queue = new Queue(baseQueueName, true, false, false, queueArgs( baseQueueName, producerProperties.getExtension(), false)); declareQueue(baseQueueName, queue); - binding = notPartitionedBinding(exchange, queue, - producerProperties.getExtension()); + String[] routingKeys = bindingRoutingKeys(producerProperties.getExtension()); + if (ObjectUtils.isEmpty(routingKeys)) { + binding = notPartitionedBinding(exchange, queue, null, producerProperties.getExtension()); + } + else { + for (String routingKey : routingKeys) { + binding = notPartitionedBinding(exchange, queue, routingKey, + producerProperties.getExtension()); + } + } } } else { @@ -132,8 +141,17 @@ public class RabbitExchangeQueueProvisioner String prefix = producerProperties.getExtension().getPrefix(); String destination = StringUtils.isEmpty(prefix) ? exchangeName : exchangeName.substring(prefix.length()); - binding = partitionedBinding(destination, exchange, queue, + String[] routingKeys = bindingRoutingKeys(producerProperties.getExtension()); + if (ObjectUtils.isEmpty(routingKeys)) { + binding = partitionedBinding(destination, exchange, queue, null, producerProperties.getExtension(), i); + } + else { + for (String routingKey : routingKeys) { + binding = partitionedBinding(destination, exchange, queue, routingKey, + producerProperties.getExtension(), i); + } + } } } } @@ -217,8 +235,15 @@ public class RabbitExchangeQueueProvisioner Binding binding = null; if (properties.getExtension().isBindQueue()) { declareQueue(queueName, queue); - binding = declareConsumerBindings(name, properties, exchange, partitioned, - queue); + String[] routingKeys = bindingRoutingKeys(properties.getExtension()); + if (ObjectUtils.isEmpty(routingKeys)) { + binding = declareConsumerBindings(name, null, properties, exchange, partitioned, queue); + } + else { + for (String routingKey : routingKeys) { + binding = declareConsumerBindings(name, routingKey, properties, exchange, partitioned, queue); + } + } } if (durable) { autoBindDLQ(applyPrefix(properties.getExtension().getPrefix(), baseQueueName), @@ -238,21 +263,23 @@ public class RabbitExchangeQueueProvisioner + (StringUtils.hasText(group) ? group : "default"); } - private Binding declareConsumerBindings(String name, + private Binding declareConsumerBindings(String name, String routingKey, ExtendedConsumerProperties properties, Exchange exchange, boolean partitioned, Queue queue) { + if (partitioned) { - return partitionedBinding(name, exchange, queue, properties.getExtension(), + return partitionedBinding(name, exchange, queue, routingKey, properties.getExtension(), properties.getInstanceIndex()); } else { - return notPartitionedBinding(exchange, queue, properties.getExtension()); + return notPartitionedBinding(exchange, queue, routingKey, properties.getExtension()); } } - private Binding partitionedBinding(String destination, Exchange exchange, Queue queue, + private Binding partitionedBinding(String destination, Exchange exchange, Queue queue, String rk, RabbitCommonProperties extendedProperties, int index) { - String bindingKey = extendedProperties.getBindingRoutingKey(); + + String bindingKey = rk; if (bindingKey == null) { bindingKey = destination; } @@ -286,9 +313,10 @@ public class RabbitExchangeQueueProvisioner } } - private Binding notPartitionedBinding(Exchange exchange, Queue queue, + private Binding notPartitionedBinding(Exchange exchange, Queue queue, String rk, RabbitCommonProperties extendedProperties) { - String routingKey = extendedProperties.getBindingRoutingKey(); + + String routingKey = rk; if (routingKey == null) { routingKey = "#"; } @@ -321,6 +349,14 @@ public class RabbitExchangeQueueProvisioner } } + private String[] bindingRoutingKeys(RabbitCommonProperties extendedProperties) { + /* + * When the delimiter is null, we get a String[1] containing the original. + */ + return StringUtils.delimitedListToStringArray(extendedProperties.getBindingRoutingKey(), + extendedProperties.getBindingRoutingKeyDelimiter()); + } + /** * If so requested, declare the DLX/DLQ and bind it. The DLQ is bound to the DLX with * a routing key of the original queue name because we use default exchange routing by 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 ad1ebc6b4..bac722127 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 @@ -487,7 +487,8 @@ public class RabbitBinderTests extends RabbitTestBinder binder = getBinder(); ExtendedConsumerProperties properties = createConsumerProperties(); properties.getExtension().setExchangeType(ExchangeTypes.DIRECT); - properties.getExtension().setBindingRoutingKey("foo"); + properties.getExtension().setBindingRoutingKey("foo,bar"); + properties.getExtension().setBindingRoutingKeyDelimiter(","); properties.getExtension().setQueueNameGroupOnly(true); // properties.getExtension().setDelayedExchange(true); // requires delayed message // exchange plugin; tested locally @@ -509,10 +510,14 @@ public class RabbitBinderTests extends Thread.sleep(100); bindings = client.getBindingsBySource("/", "propsUser2"); } - assertThat(bindings.size()).isEqualTo(1); + assertThat(bindings.size()).isEqualTo(2); assertThat(bindings.get(0).getSource()).isEqualTo("propsUser2"); assertThat(bindings.get(0).getDestination()).isEqualTo(group); - assertThat(bindings.get(0).getRoutingKey()).isEqualTo("foo"); + assertThat(bindings.get(0).getRoutingKey()).isIn("foo", "bar"); + assertThat(bindings.get(1).getSource()).isEqualTo("propsUser2"); + assertThat(bindings.get(1).getDestination()).isEqualTo(group); + assertThat(bindings.get(1).getRoutingKey()).isIn("foo", "bar"); + assertThat(bindings.get(1).getRoutingKey()).isNotEqualTo(bindings.get(0).getRoutingKey()); ExchangeInfo exchange = client.getExchange("/", "propsUser2"); while (n++ < 100 && exchange == null) {