From ef7e3dedcf565fd30f7cad29b65de11e36017b6f Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Fri, 10 Mar 2017 14:34:50 -0500 Subject: [PATCH] GH-9: Remove request/replyHeaderPatterns Resolves #9 Replace with simply `headerPatterns`. The reply header patterns were never used and are removed. The request header patterns are deprecated in favor of header patterns. --- .../properties/RabbitConsumerProperties.java | 26 +++++++++----- .../properties/RabbitProducerProperties.java | 26 +++++++++----- .../src/main/asciidoc/overview.adoc | 24 +++++-------- .../rabbit/RabbitMessageChannelBinder.java | 17 +++++---- .../binder/rabbit/RabbitBinderTests.java | 36 +++++++++++++------ 5 files changed, 77 insertions(+), 52 deletions(-) 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 ae5999499..bdae3a24c 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 @@ -35,8 +35,6 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { private int prefetch = 1; - private String[] requestHeaderPatterns = new String[] {"STANDARD_REQUEST_HEADERS", "*"}; - private int txSize = 1; private boolean durableSubscription = true; @@ -45,7 +43,7 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { private boolean requeueRejected = false; - private String[] replyHeaderPatterns = new String[] {"STANDARD_REPLY_HEADERS", "*"}; + private String[] headerPatterns = new String[] {"*"}; private long recoveryInterval = 5000; @@ -84,12 +82,22 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { this.prefetch = prefetch; } + /** + * @deprecated - use {@link #getHeaderPatterns()}. + * @return the header patterns. + */ + @Deprecated public String[] getRequestHeaderPatterns() { - return requestHeaderPatterns; + return this.headerPatterns; } + /** + * @deprecated - use {@link #setHeaderPatterns(String[])}. + * @param requestHeaderPatterns + */ + @Deprecated public void setRequestHeaderPatterns(String[] requestHeaderPatterns) { - this.requestHeaderPatterns = requestHeaderPatterns; + this.headerPatterns = requestHeaderPatterns; } @Min(value = 1, message = "Tx Size should be greater than zero.") @@ -125,12 +133,12 @@ public class RabbitConsumerProperties extends RabbitCommonProperties { this.requeueRejected = requeueRejected; } - public String[] getReplyHeaderPatterns() { - return replyHeaderPatterns; + public String[] getHeaderPatterns() { + return headerPatterns; } - public void setReplyHeaderPatterns(String[] replyHeaderPatterns) { - this.replyHeaderPatterns = replyHeaderPatterns; + public void setHeaderPatterns(String[] replyHeaderPatterns) { + this.headerPatterns = replyHeaderPatterns; } public long getRecoveryInterval() { diff --git a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java index 5a420693d..6b2a7e733 100644 --- a/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java +++ b/spring-cloud-stream-binder-rabbit-core/src/main/java/org/springframework/cloud/stream/binder/rabbit/properties/RabbitProducerProperties.java @@ -26,8 +26,6 @@ import org.springframework.amqp.core.MessageDeliveryMode; */ public class RabbitProducerProperties extends RabbitCommonProperties { - private String[] requestHeaderPatterns = new String[] {"STANDARD_REQUEST_HEADERS", "*"}; - private boolean compress; private boolean batchingEnabled; @@ -42,7 +40,7 @@ public class RabbitProducerProperties extends RabbitCommonProperties { private MessageDeliveryMode deliveryMode = MessageDeliveryMode.PERSISTENT; - private String[] replyHeaderPatterns = new String[] {"STANDARD_REPLY_HEADERS", "*"}; + private String[] headerPatterns = new String[] {"*"}; /** * When using a delayed message exchange, a SpEL expression to determine the delay to apply to messages @@ -54,12 +52,22 @@ public class RabbitProducerProperties extends RabbitCommonProperties { */ private String routingKeyExpression; + /** + * @deprecated - use {@link #setHeaderPatterns(String[])}. + * @param requestHeaderPatterns the patterns. + */ + @Deprecated public void setRequestHeaderPatterns(String[] requestHeaderPatterns) { - this.requestHeaderPatterns = requestHeaderPatterns; + this.headerPatterns = requestHeaderPatterns; } + /** + * @deprecated - use {@link #getHeaderPatterns()}. + * @return the header patterns. + */ + @Deprecated public String[] getRequestHeaderPatterns() { - return requestHeaderPatterns; + return this.headerPatterns; } public void setCompress(boolean compress) { @@ -78,12 +86,12 @@ public class RabbitProducerProperties extends RabbitCommonProperties { return deliveryMode; } - public String[] getReplyHeaderPatterns() { - return replyHeaderPatterns; + public String[] getHeaderPatterns() { + return headerPatterns; } - public void setReplyHeaderPatterns(String[] replyHeaderPatterns) { - this.replyHeaderPatterns = replyHeaderPatterns; + public void setHeaderPatterns(String[] replyHeaderPatterns) { + this.headerPatterns = replyHeaderPatterns; } public boolean isBatchingEnabled() { 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 ca0ee5c51..caf8cc731 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 @@ -179,6 +179,10 @@ expires:: how long before an unused queue is deleted (ms) + Default: `no expiration` +headerPatterns:: + Patterns for headers to be mapped from inbound messages. ++ +Default: `['*']` (all headers). maxConcurrency:: the maximum number of consumers + @@ -211,14 +215,6 @@ requeueRejected:: Whether delivery failures should be requeued when retry is disabled or republishToDlq is false. + Default: `false`. -requestHeaderPatterns:: - The request headers to be transported. -+ -Default: `[STANDARD_REQUEST_HEADERS,'*']`. -replyHeaderPatterns:: - The reply headers to be transported. -+ -Default: `[STANDARD_REPLY_HEADERS,'*']`. republishToDlq:: By default, messages which fail after retries are exhausted are rejected. If a dead-letter queue (DLQ) is configured, RabbitMQ will route the failed message (unchanged) to the DLQ. @@ -358,6 +354,10 @@ expires:: Only applies if `requiredGroups` are provided and then only to those groups. + Default: `no expiration` +headerPatterns:: + Patterns for headers to be mapped to outbound messages. ++ +Default: `['*']` (all headers). maxLength:: maximum number of messages in the queue Only applies if `requiredGroups` are provided and then only to those groups. @@ -377,14 +377,6 @@ prefix:: A prefix to be added to the name of the `destination` exchange. + Default: "". -requestHeaderPatterns:: - The request headers to be transported. -+ -Default: `[STANDARD_REQUEST_HEADERS,'*']`. -replyHeaderPatterns:: - The reply headers to be transported. -+ -Default: `[STANDARD_REPLY_HEADERS,'*']`. routingKeyExpression:: A SpEL expression to determine the routing key to use when publishing messages. + 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 8c57c1a06..b11906a87 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 @@ -16,10 +16,9 @@ package org.springframework.cloud.stream.binder.rabbit; +import java.util.ArrayList; import java.util.Arrays; - -import com.rabbitmq.client.AMQP; -import com.rabbitmq.client.Envelope; +import java.util.List; import org.springframework.amqp.core.MessagePostProcessor; import org.springframework.amqp.core.MessageProperties; @@ -65,6 +64,9 @@ import org.springframework.scheduling.TaskScheduler; import org.springframework.util.Assert; import org.springframework.util.StringUtils; +import com.rabbitmq.client.AMQP; +import com.rabbitmq.client.Envelope; + /** * A {@link org.springframework.cloud.stream.binder.Binder} implementation backed by RabbitMQ. * @author Mark Fisher @@ -206,8 +208,10 @@ public class RabbitMessageChannelBinder endpoint.setDelayExpressionString(extendedProperties.getDelayExpression()); } DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.outboundMapper(); - mapper.setRequestHeaderNames(extendedProperties.getRequestHeaderPatterns()); - mapper.setReplyHeaderNames(extendedProperties.getReplyHeaderPatterns()); + List headerPatterns = new ArrayList<>(extendedProperties.getHeaderPatterns().length + 1); + headerPatterns.add("!" + BinderHeaders.PARTITION_HEADER); + headerPatterns.addAll(Arrays.asList(extendedProperties.getHeaderPatterns())); + mapper.setRequestHeaderNames(headerPatterns.toArray(new String[headerPatterns.size()])); endpoint.setHeaderMapper(mapper); endpoint.setDefaultDeliveryMode(extendedProperties.getDeliveryMode()); endpoint.setBeanFactory(this.getBeanFactory()); @@ -268,8 +272,7 @@ public class RabbitMessageChannelBinder adapter.setBeanFactory(this.getBeanFactory()); adapter.setBeanName("inbound." + baseQueueName); DefaultAmqpHeaderMapper mapper = DefaultAmqpHeaderMapper.inboundMapper(); - mapper.setRequestHeaderNames(properties.getExtension().getRequestHeaderPatterns()); - mapper.setReplyHeaderNames(properties.getExtension().getReplyHeaderPatterns()); + mapper.setRequestHeaderNames(properties.getExtension().getHeaderPatterns()); adapter.setHeaderMapper(mapper); adapter.afterPropertiesSet(); return adapter; 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 925c6c859..1755fecbd 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 @@ -52,6 +52,7 @@ import org.springframework.amqp.support.postprocessor.DelegatingDecompressingPos import org.springframework.amqp.utils.test.TestUtils; import org.springframework.beans.DirectFieldAccessor; import org.springframework.boot.autoconfigure.amqp.RabbitProperties; +import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; @@ -402,7 +403,8 @@ public class RabbitBinderTests extends producerBinding = binder.bindProducer("props.0", channel, producerProperties); endpoint = extractEndpoint(producerBinding); - assertThat(getEndpointRouting(endpoint)).isEqualTo("'props.0-' + headers['partition']"); + assertThat(getEndpointRouting(endpoint)) + .isEqualTo("'props.0-' + headers['" + BinderHeaders.PARTITION_HEADER + "']"); assertThat(TestUtils.getPropertyValue(endpoint, "delayExpression", SpelExpression.class) .getExpressionString()).isEqualTo("42"); mode = TestUtils.getPropertyValue(endpoint, "defaultDeliveryMode", MessageDeliveryMode.class); @@ -609,12 +611,16 @@ public class RabbitBinderTests extends org.springframework.amqp.core.Message received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 1); + assertThat(received.getMessageProperties().getReceivedRoutingKey()) + .isEqualTo("bindertest.partDLQ.0.dlqPartGrp-1"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); output.send(new GenericMessage<>(0)); received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 0); + assertThat(received.getMessageProperties().getReceivedRoutingKey()) + .isEqualTo("bindertest.partDLQ.0.dlqPartGrp-0"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); input0Binding.unbind(); input1Binding.unbind(); @@ -697,12 +703,16 @@ public class RabbitBinderTests extends org.springframework.amqp.core.Message received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 1); + assertThat(received.getMessageProperties().getHeaders().get("x-original-routingKey")) + .isEqualTo("partPubDLQ.0-1"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); output.send(new GenericMessage<>(0)); received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 0); + assertThat(received.getMessageProperties().getHeaders().get("x-original-routingKey")) + .isEqualTo("partPubDLQ.0-0"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); input0Binding.unbind(); input1Binding.unbind(); @@ -787,12 +797,16 @@ public class RabbitBinderTests extends org.springframework.amqp.core.Message received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 1); + assertThat(received.getMessageProperties().getReceivedRoutingKey()) + .isEqualTo("bindertest.partDLQ.1.dlqPartGrp-1"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); output.send(new GenericMessage(0)); received = template.receive(streamDLQName); assertThat(received).isNotNull(); - assertThat(received.getMessageProperties().getHeaders()).containsEntry("partition", 0); + assertThat(received.getMessageProperties().getReceivedRoutingKey()) + .isEqualTo("bindertest.partDLQ.1.dlqPartGrp-0"); + assertThat(received.getMessageProperties().getHeaders()).doesNotContainKey(BinderHeaders.PARTITION_HEADER); input0Binding.unbind(); input1Binding.unbind(); @@ -1055,8 +1069,8 @@ public class RabbitBinderTests extends private void verifyFooRequestProducer(Lifecycle endpoint) { List requestMatchers = TestUtils.getPropertyValue(endpoint, "headerMapper.requestHeaderMatcher.matchers", List.class); - assertThat(requestMatchers).hasSize(1); - assertThat(TestUtils.getPropertyValue(requestMatchers.get(0), "pattern")).isEqualTo("foo"); + assertThat(requestMatchers).hasSize(2); + assertThat(TestUtils.getPropertyValue(requestMatchers.get(1), "pattern")).isEqualTo("foo"); } @Override @@ -1077,8 +1091,8 @@ public class RabbitBinderTests extends @Override protected void checkRkExpressionForPartitionedModuleSpEL(Object endpoint) { - assertThat(getEndpointRouting(endpoint)) - .contains(getExpectedRoutingBaseDestination("'part.0'", "test") + " + '-' + headers['partition']"); + assertThat(getEndpointRouting(endpoint)).contains(getExpectedRoutingBaseDestination("'part.0'", "test") + + " + '-' + headers['" + BinderHeaders.PARTITION_HEADER + "']"); } @Override