diff --git a/build.gradle b/build.gradle index 2be5a90..bccb769 100644 --- a/build.gradle +++ b/build.gradle @@ -56,7 +56,7 @@ ext { servletApiVersion = '3.1.0' slf4jVersion = '1.7.21' springCloudAwsVersion = '1.1.0.RELEASE' - springIntegrationVersion = '4.2.6.RELEASE' + springIntegrationVersion = '4.2.8.RELEASE' idPrefix = 'aws' diff --git a/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java b/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java index 3617f31..6d63bb0 100644 --- a/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java +++ b/src/main/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParser.java @@ -32,6 +32,7 @@ import org.springframework.util.StringUtils; * The parser for the {@code }. * * @author Artem Bilan + * @author Patrick Fitzsimons */ public class SqsMessageDrivenChannelAdapterParser extends AbstractSingleBeanDefinitionParser { @@ -81,6 +82,7 @@ public class SqsMessageDrivenChannelAdapterParser extends AbstractSingleBeanDefi IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "max-number-of-messages"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "visibility-timeout"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "wait-time-out"); + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "queue-stop-timeout"); } } diff --git a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java index 546a795..c969ad5 100644 --- a/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java +++ b/src/main/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapter.java @@ -51,8 +51,7 @@ import com.amazonaws.services.sqs.AmazonSQSAsync; * @see SimpleMessageListenerContainer * @see QueueMessageHandler */ -public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport - implements DisposableBean { +public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport implements DisposableBean { private final SimpleMessageListenerContainerFactory simpleMessageListenerContainerFactory = new SimpleMessageListenerContainerFactory(); @@ -61,6 +60,8 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport private SimpleMessageListenerContainer listenerContainer; + private Long queueStopTimeout; + private SqsMessageDeletionPolicy messageDeletionPolicy = SqsMessageDeletionPolicy.NO_REDRIVE; public SqsMessageDrivenChannelAdapter(AmazonSQSAsync amazonSqs, String... queues) { @@ -99,6 +100,10 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport this.simpleMessageListenerContainerFactory.setDestinationResolver(destinationResolver); } + public void setQueueStopTimeout(long queueStopTimeout) { + this.queueStopTimeout = queueStopTimeout; + } + public void setMessageDeletionPolicy(SqsMessageDeletionPolicy messageDeletionPolicy) { Assert.notNull(messageDeletionPolicy, "'messageDeletionPolicy' must not be null."); this.messageDeletionPolicy = messageDeletionPolicy; @@ -108,6 +113,9 @@ public class SqsMessageDrivenChannelAdapter extends MessageProducerSupport protected void onInit() { super.onInit(); this.listenerContainer = this.simpleMessageListenerContainerFactory.createSimpleMessageListenerContainer(); + if (this.queueStopTimeout != null) { + this.listenerContainer.setQueueStopTimeout(this.queueStopTimeout); + } this.listenerContainer.setMessageHandler(new IntegrationQueueMessageHandler()); try { this.listenerContainer.afterPropertiesSet(); diff --git a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-1.0.xsd b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-1.0.xsd index d59d0d9..a3bfdf6 100644 --- a/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-1.0.xsd +++ b/src/main/resources/org/springframework/integration/aws/config/spring-integration-aws-1.0.xsd @@ -560,6 +560,14 @@ + + + + Configure the maximum number of milliseconds the method waits for a queue + to stop before interrupting the current thread. Defaults to 10000. + + + diff --git a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml index d11dfa8..ce760ae 100644 --- a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml +++ b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests-context.xml @@ -32,6 +32,7 @@ visibility-timeout="200" wait-time-out="40" send-timeout="2000" + queue-stop-timeout="11000" destination-resolver="destinationResolver" resource-id-resolver="resourceIdResolver"/> diff --git a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java index 25c2a6c..3cb5a42 100644 --- a/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java +++ b/src/test/java/org/springframework/integration/aws/config/xml/SqsMessageDrivenChannelAdapterParserTests.java @@ -96,6 +96,7 @@ public class SqsMessageDrivenChannelAdapterParserTests { assertThat(TestUtils.getPropertyValue(listenerContainer, "maxNumberOfMessages")).isEqualTo(5); assertThat(TestUtils.getPropertyValue(listenerContainer, "visibilityTimeout")).isEqualTo(200); assertThat(TestUtils.getPropertyValue(listenerContainer, "waitTimeOut")).isEqualTo(40); + assertThat(TestUtils.getPropertyValue(listenerContainer, "queueStopTimeout")).isEqualTo(11000L); assertThat(TestUtils.getPropertyValue(listenerContainer, "autoStartup")).isEqualTo(false); @SuppressWarnings("rawtypes") diff --git a/src/test/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapterTests.java b/src/test/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapterTests.java index 184df0e..6432fc1 100644 --- a/src/test/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapterTests.java +++ b/src/test/java/org/springframework/integration/aws/inbound/SqsMessageDrivenChannelAdapterTests.java @@ -31,6 +31,7 @@ import org.springframework.integration.aws.support.AwsHeaders; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; import org.springframework.integration.core.MessageProducer; +import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.PollableChannel; import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @@ -54,8 +55,14 @@ public class SqsMessageDrivenChannelAdapterTests { @Autowired private PollableChannel inputChannel; + @Autowired + private SqsMessageDrivenChannelAdapter sqsMessageDrivenChannelAdapter; + @Test public void testSqsMessageDrivenChannelAdapter() { + assertThat(TestUtils.getPropertyValue(this.sqsMessageDrivenChannelAdapter, + "listenerContainer.queueStopTimeout")) + .isEqualTo(10000L); org.springframework.messaging.Message receive = this.inputChannel.receive(1000); assertThat(receive).isNotNull(); assertThat((String) receive.getPayload()).isIn("messageContent", "messageContent2");