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.
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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.
|
||||
+
|
||||
|
||||
@@ -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<String> 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;
|
||||
|
||||
@@ -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<Integer>(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
|
||||
|
||||
Reference in New Issue
Block a user