From 5b41dc151339e6c91ff35bb7068e26e21a315d57 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 6 Oct 2016 13:39:28 -0400 Subject: [PATCH] INT-4131: Add delayExpression to AMQP Outbounds JIRA: https://jira.spring.io/browse/INT-4131 Specify an expression on the outbound endpoints to set the `x-delay` header when using the RabbitMQ Delayed Message Exchange plugin. Polishing - PR Comments - Add `setDelay` - Port `FunctionExpression` from DSL - Add `SupplierExpression` Javadoc Fixes Use ValueExpression for delay * Simple polishing for `SupplierExpression` * Mention plain `delay` property in the `amqp.adoc` Conflicts: spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/OutboundEndpointTests.java Resolved. Backport Polishing Java 6 cpmpatibility; deprecation warnings. --- .../AmqpOutboundChannelAdapterParser.java | 2 + .../config/AmqpOutboundGatewayParser.java | 2 + .../AbstractAmqpOutboundEndpoint.java | 56 +++++++++++ .../amqp/outbound/AmqpOutboundEndpoint.java | 2 + .../outbound/AsyncAmqpOutboundGateway.java | 7 +- .../config/spring-integration-amqp-4.3.xsd | 14 +++ ...boundChannelAdapterParserTests-context.xml | 1 + ...AmqpOutboundChannelAdapterParserTests.java | 3 + ...AmqpOutboundGatewayParserTests-context.xml | 2 + .../AmqpOutboundGatewayParserTests.java | 3 + .../amqp/outbound/OutboundEndpointTests.java | 95 +++++++++++++++---- .../support/DefaultAmqpHeaderMapperTests.java | 2 + .../ExpressionEvaluatingMessageProcessor.java | 4 +- src/reference/asciidoc/amqp.adoc | 7 ++ src/reference/asciidoc/whats-new.adoc | 5 + 15 files changed, 181 insertions(+), 24 deletions(-) diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java index 6f84a66db1..0fcbf1fa36 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java @@ -81,6 +81,8 @@ public class AmqpOutboundChannelAdapterParser extends AbstractOutboundChannelAda IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-ack-channel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-nack-channel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "return-channel"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "delay-expression", + "delayExpressionString"); return builder.getBeanDefinition(); } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java index bfd4453847..ac5dafd556 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java @@ -100,6 +100,8 @@ public class AmqpOutboundGatewayParser extends AbstractConsumerEndpointParser { mapperBuilder.setFactoryMethod("outboundMapper"); IntegrationNamespaceUtils.configureHeaderMapper(element, builder, parserContext, mapperBuilder, null); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "delay-expression", + "delayExpressionString"); return builder; } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java index 325d519763..3425e05deb 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java @@ -31,6 +31,7 @@ import org.springframework.expression.Expression; import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; import org.springframework.integration.channel.NullChannel; +import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; @@ -77,6 +78,10 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin private volatile ConnectionFactory connectionFactory; + private volatile Expression delayExpression; + + private volatile ExpressionEvaluatingMessageProcessor delayGenerator; + private volatile boolean running; public void setHeaderMapper(AmqpHeaderMapper headerMapper) { @@ -188,6 +193,44 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin this.lazyConnect = lazyConnect; } + /** + * Set the value to set in the {@code x-delay} header when using the + * RabbitMQ delayed message exchange plugin. By default, the {@link AmqpHeaders#DELAY} + * header (if present) is mapped; setting the delay here overrides that value. + * @param delay the delay. + * @since 4.3.5 + */ + public void setDelay(int delay) { + this.delayExpression = new ValueExpression(delay); + } + + /** + * Set the SpEL expression to calculate the {@code x-delay} header when using the + * RabbitMQ delayed message exchange plugin. By default, the {@link AmqpHeaders#DELAY} + * header (if present) is mapped; setting the expression here overrides that value. + * @param delayExpression the expression. + * @since 4.3.5 + */ + public void setDelayExpression(Expression delayExpression) { + this.delayExpression = delayExpression; + } + + /** + * Set the SpEL expression to calculate the {@code x-delay} header when using the + * RabbitMQ delayed message exchange plugin. By default, the {@link AmqpHeaders#DELAY} + * header (if present) is mapped; setting the expression here overrides that value. + * @param delayExpression the expression. + * @since 4.3.5 + */ + public void setDelayExpressionString(String delayExpression) { + if (delayExpression == null) { + this.delayExpression = null; + } + else { + this.delayExpression = EXPRESSION_PARSER.parseExpression(delayExpression); + } + } + protected final void setConnectionFactory(ConnectionFactory connectionFactory) { this.connectionFactory = connectionFactory; } @@ -284,6 +327,13 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin Assert.state(this.confirmNackChannel == null || nullChannel != null, "A 'confirmCorrelationExpression' is required when specifying a 'confirmNackChannel'"); } + if (this.delayExpression != null) { + this.delayGenerator = new ExpressionEvaluatingMessageProcessor(this.delayExpression, + Integer.class); + if (beanFactory != null) { + this.delayGenerator.setBeanFactory(beanFactory); + } + } endpointInit(); } @@ -365,6 +415,12 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin return routingKey; } + protected void addDelayProperty(Message message, org.springframework.amqp.core.Message amqpMessage) { + if (this.delayGenerator != null) { + amqpMessage.getMessageProperties().setDelay(this.delayGenerator.processMessage(message)); + } + } + protected Message buildReplyMessage(MessageConverter converter, org.springframework.amqp.core.Message amqpReplyMessage) { Object replyObject = converter.fromMessage(amqpReplyMessage); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java index 51b36c3bc0..8c7ebd0c64 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpoint.java @@ -130,6 +130,7 @@ public class AmqpOutboundEndpoint extends AbstractAmqpOutboundEndpoint MessageConverter converter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); org.springframework.amqp.core.Message amqpMessage = MappingUtils.mapMessage(requestMessage, converter, getHeaderMapper(), getDefaultDeliveryMode()); + addDelayProperty(requestMessage, amqpMessage); ((RabbitTemplate) this.amqpTemplate).send(exchangeName, routingKey, amqpMessage, correlationData); } else { @@ -153,6 +154,7 @@ public class AmqpOutboundEndpoint extends AbstractAmqpOutboundEndpoint MessageConverter converter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter(); org.springframework.amqp.core.Message amqpMessage = MappingUtils.mapMessage(requestMessage, converter, getHeaderMapper(), getDefaultDeliveryMode()); + addDelayProperty(requestMessage, amqpMessage); org.springframework.amqp.core.Message amqpReplyMessage = ((RabbitTemplate) this.amqpTemplate).sendAndReceive(exchangeName, routingKey, amqpMessage, correlationData); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java index b3573911ca..78f03c9e4a 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AsyncAmqpOutboundGateway.java @@ -60,10 +60,11 @@ public class AsyncAmqpOutboundGateway extends AbstractAmqpOutboundEndpoint { @Override protected Object handleRequestMessage(Message requestMessage) { + org.springframework.amqp.core.Message amqpMessage = MappingUtils.mapMessage(requestMessage, + this.messageConverter, getHeaderMapper(), getDefaultDeliveryMode()); + addDelayProperty(requestMessage, amqpMessage); RabbitMessageFuture future = this.template.sendAndReceive(generateExchangeName(requestMessage), - generateRoutingKey(requestMessage), - MappingUtils.mapMessage(requestMessage, this.messageConverter, getHeaderMapper(), - getDefaultDeliveryMode())); + generateRoutingKey(requestMessage), amqpMessage); future.addCallback(new FutureCallback(requestMessage)); CorrelationData correlationData = generateCorrelationData(requestMessage); if (correlationData != null && future.getConfirm() != null) { diff --git a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.3.xsd b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.3.xsd index 58549adb1a..a1efb88e23 100644 --- a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.3.xsd +++ b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-4.3.xsd @@ -542,6 +542,20 @@ property set to TRUE. + + + + + + + + + + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml index 5fd66b7d5b..ac14e3e4b2 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests-context.xml @@ -25,6 +25,7 @@ exchange-name="outboundchanneladapter.test.1" default-delivery-mode="NON_PERSISTENT" lazy-connect="false" + delay-expression="42" mapped-request-headers="foo*"/> diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java index 90575a39ae..49f375a67b 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java @@ -126,6 +126,9 @@ public class AmqpOutboundChannelAdapterParserTests { AmqpOutboundEndpoint.class); assertNotNull(TestUtils.getPropertyValue(endpoint, "defaultDeliveryMode")); assertFalse(TestUtils.getPropertyValue(endpoint, "lazyConnect", Boolean.class)); + assertEquals("42", + TestUtils.getPropertyValue(endpoint, "delayExpression", org.springframework.expression.Expression.class) + .getExpressionString()); Field amqpTemplateField = ReflectionUtils.findField(AmqpOutboundEndpoint.class, "amqpTemplate"); amqpTemplateField.setAccessible(true); diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml index d674a90ec8..f7c1ef8b0b 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml @@ -15,6 +15,7 @@ exchange-name="si.test.exchange" routing-key="si.test.binding" amqp-template="amqpTemplate" + delay-expression="42" auto-startup="false" order="5" return-channel="returnChannel"> @@ -92,6 +93,7 @@ exchange-name="si.test.exchange" routing-key="si.test.binding" async-template="asyncTemplate" + delay-expression="42" auto-startup="false" order="5" return-channel="returnChannel"> diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java index 3b1c907295..8cb23b773a 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java @@ -94,6 +94,9 @@ public class AmqpOutboundGatewayParserTests { assertEquals(Long.valueOf(777), sendTimeout); assertTrue(TestUtils.getPropertyValue(gateway, "lazyConnect", Boolean.class)); + assertEquals("42", + TestUtils.getPropertyValue(gateway, "delayExpression", org.springframework.expression.Expression.class) + .getExpressionString()); } @SuppressWarnings("rawtypes") diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/OutboundEndpointTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/OutboundEndpointTests.java index caf5a4483f..c441c46ac8 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/OutboundEndpointTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/OutboundEndpointTests.java @@ -16,29 +16,39 @@ package org.springframework.integration.amqp.outbound; +import static org.hamcrest.Matchers.equalTo; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertNull; +import static org.junit.Assert.assertThat; +import static org.mockito.BDDMockito.willAnswer; +import static org.mockito.BDDMockito.willDoNothing; import static org.mockito.Matchers.any; import static org.mockito.Matchers.anyString; -import static org.mockito.Mockito.doAnswer; +import static org.mockito.Matchers.eq; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import java.util.concurrent.atomic.AtomicReference; import org.junit.Test; -import org.mockito.invocation.InvocationOnMock; -import org.mockito.stubbing.Answer; +import org.mockito.ArgumentCaptor; import org.springframework.amqp.core.Message; +import org.springframework.amqp.rabbit.AsyncRabbitTemplate; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; +import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer; import org.springframework.amqp.rabbit.support.CorrelationData; +import org.springframework.beans.factory.BeanFactory; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; +import org.springframework.integration.channel.NullChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.scheduling.TaskScheduler; /** * @author Gary Russell @@ -47,6 +57,61 @@ import org.springframework.messaging.MessageHeaders; */ public class OutboundEndpointTests { + @Test + public void testDelayExpression() { + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + RabbitTemplate amqpTemplate = spy(new RabbitTemplate(connectionFactory)); + AmqpOutboundEndpoint endpoint = new AmqpOutboundEndpoint(amqpTemplate); + willDoNothing() + .given(amqpTemplate).send(anyString(), anyString(), any(Message.class), any(CorrelationData.class)); + willAnswer(invocation -> invocation.getArgumentAt(2, Message.class)) + .given(amqpTemplate) + .sendAndReceive(anyString(), anyString(), any(Message.class), any(CorrelationData.class)); + endpoint.setExchangeName("foo"); + endpoint.setRoutingKey("bar"); + endpoint.setDelayExpressionString("42"); + endpoint.setBeanFactory(mock(BeanFactory.class)); + endpoint.afterPropertiesSet(); + endpoint.handleMessage(new GenericMessage<>("foo")); + ArgumentCaptor captor = ArgumentCaptor.forClass(Message.class); + verify(amqpTemplate).send(eq("foo"), eq("bar"), captor.capture(), any(CorrelationData.class)); + assertThat(captor.getValue().getMessageProperties().getDelay(), equalTo(42)); + endpoint.setExpectReply(true); + endpoint.setOutputChannel(new NullChannel()); + endpoint.handleMessage(new GenericMessage<>("foo")); + verify(amqpTemplate).sendAndReceive(eq("foo"), eq("bar"), captor.capture(), any(CorrelationData.class)); + assertThat(captor.getValue().getMessageProperties().getDelay(), equalTo(42)); + + endpoint.setDelay(23); + endpoint.setRoutingKey("baz"); + endpoint.afterPropertiesSet(); + endpoint.handleMessage(new GenericMessage<>("foo")); + verify(amqpTemplate).sendAndReceive(eq("foo"), eq("baz"), captor.capture(), any(CorrelationData.class)); + assertThat(captor.getValue().getMessageProperties().getDelay(), equalTo(23)); + } + + @Test + public void testAsyncDelayExpression() { + ConnectionFactory connectionFactory = mock(ConnectionFactory.class); + AsyncRabbitTemplate amqpTemplate = spy(new AsyncRabbitTemplate(new RabbitTemplate(connectionFactory), + new SimpleMessageListenerContainer(connectionFactory), "replyTo")); + amqpTemplate.setTaskScheduler(mock(TaskScheduler.class)); + AsyncAmqpOutboundGateway gateway = new AsyncAmqpOutboundGateway(amqpTemplate); + willAnswer( + invocation -> amqpTemplate.new RabbitMessageFuture("foo", invocation.getArgumentAt(2, Message.class))) + .given(amqpTemplate).sendAndReceive(anyString(), anyString(), any(Message.class)); + gateway.setExchangeName("foo"); + gateway.setRoutingKey("bar"); + gateway.setDelayExpressionString("42"); + gateway.setBeanFactory(mock(BeanFactory.class)); + gateway.setOutputChannel(new NullChannel()); + gateway.afterPropertiesSet(); + ArgumentCaptor captor = ArgumentCaptor.forClass(Message.class); + gateway.handleMessage(new GenericMessage<>("foo")); + verify(amqpTemplate).sendAndReceive(eq("foo"), eq("bar"), captor.capture()); + assertThat(captor.getValue().getMessageProperties().getDelay(), equalTo(42)); + } + @Test public void testHeaderMapperWinsAdapter() { ConnectionFactory connectionFactory = mock(ConnectionFactory.class); @@ -54,14 +119,10 @@ public class OutboundEndpointTests { AmqpOutboundEndpoint endpoint = new AmqpOutboundEndpoint(amqpTemplate); final AtomicReference amqpMessage = new AtomicReference(); - doAnswer(new Answer() { - - @Override - public Object answer(InvocationOnMock invocation) throws Throwable { - amqpMessage.set((Message) invocation.getArguments()[2]); - return null; - } - }).when(amqpTemplate).send(anyString(), anyString(), any(Message.class), + willAnswer(invocation -> { + amqpMessage.set((Message) invocation.getArguments()[2]); + return null; + }).given(amqpTemplate).send(anyString(), anyString(), any(Message.class), any(CorrelationData.class)); org.springframework.messaging.Message message = MessageBuilder.withPayload("foo") .setHeader(MessageHeaders.CONTENT_TYPE, "bar") @@ -82,14 +143,10 @@ public class OutboundEndpointTests { endpoint.setHeaderMapper(mapper); final AtomicReference amqpMessage = new AtomicReference(); - doAnswer(new Answer() { - - @Override - public Object answer(InvocationOnMock invocation) throws Throwable { - amqpMessage.set((Message) invocation.getArguments()[2]); - return null; - } - }).when(amqpTemplate) + willAnswer(invocation -> { + amqpMessage.set((Message) invocation.getArguments()[2]); + return null; + }).given(amqpTemplate) .doSendAndReceiveWithTemporary(anyString(), anyString(), any(Message.class), any(CorrelationData.class)); org.springframework.messaging.Message message = MessageBuilder.withPayload("foo") .setHeader(MessageHeaders.CONTENT_TYPE, "bar") diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java index 8f69dc2e62..0309e978e2 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/support/DefaultAmqpHeaderMapperTests.java @@ -49,6 +49,7 @@ import org.springframework.util.MimeTypeUtils; */ public class DefaultAmqpHeaderMapperTests { + @SuppressWarnings("deprecation") @Test public void fromHeaders() { DefaultAmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.outboundMapper(); @@ -149,6 +150,7 @@ public class DefaultAmqpHeaderMapperTests { } + @SuppressWarnings("deprecation") @Test public void toHeaders() { DefaultAmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.inboundMapper(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/ExpressionEvaluatingMessageProcessor.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/ExpressionEvaluatingMessageProcessor.java index 40c7787a08..4b89402af1 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/ExpressionEvaluatingMessageProcessor.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/ExpressionEvaluatingMessageProcessor.java @@ -1,5 +1,5 @@ /* - * Copyright 2002-2014 the original author or authors. + * Copyright 2002-2016 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -69,7 +69,7 @@ public class ExpressionEvaluatingMessageProcessor extends AbstractMessageProc */ @Override public T processMessage(Message message) { - return this.evaluateExpression(this.expression, message, this.expectedType); + return evaluateExpression(this.expression, message, this.expectedType); } @Override diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 5b3e7376c1..610f9ce3d4 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -1149,7 +1149,14 @@ To configure a default user id for outbound messages, configure it on a `RabbitT Similarly, to set the user id property on replies, inject an appropriately configured template into the inbound gateway. See the http://docs.spring.io/spring-amqp/reference/html/_reference.html#template-user-id[Spring AMQP documentation] for more information. +[[amqp-delay]] +=== Delayed Message Exchange +Spring AMQP supports the http://docs.spring.io/spring-amqp/reference/html/_reference.html#delayed-message-exchange[RabbitMQ Delayed Message Exchange Plugin]. +For inbound messages, the `x-delay` header is mapped to the `AmqpHeaders.RECEIVED_DELAY` header. +Setting the `AMQPHeaders.DELAY` header will cause the corresponding `x-delay` header to be set in outbound messages. +You can also specify the `delay` and `delayExpression` properties on outbound endpoints (`delay-expression` when using XML configuration). +This takes precedence over the `AmqpHeaders.DELAY` header. [[amqp-channels]] === AMQP Backed Message Channels diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index 017b471025..72c376537d 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -295,3 +295,8 @@ See <> for more information. The `BarrierMessageHandler` now supports a discard channel to which late-arriving trigger messages are sent. See <> for more information. + +==== AMQP Changes + +The AMQP outbound endpoints now support setting a delay expression for when using the RabbitMQ Delayed Message Exchange plugin. +See <> for more information.