GH-36: Add SqsMessageDrivenCA.queueStopTimeout
Fixes GH-36 (https://github.com/spring-projects/spring-integration-aws/issues/36) * Add ability to set `queueStopTimeout` for `SimpleMessageListenerContainer` in `SqsMessageDrivenChannelAdapter` * Update unit tests to test `queueStopTimeout` * Ensure that we don't override the default `queueStopTimeout` with `0` `long` * Upgrade to `SI-4.2.8`
This commit is contained in:
committed by
Artem Bilan
parent
d2aa5224ec
commit
7be94b7690
@@ -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'
|
||||
|
||||
|
||||
@@ -32,6 +32,7 @@ import org.springframework.util.StringUtils;
|
||||
* The parser for the {@code <int-aws:sqs-message-driven-channel-adapter>}.
|
||||
*
|
||||
* @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");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
@@ -560,6 +560,14 @@
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="queue-stop-timeout">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
Configure the maximum number of milliseconds the method waits for a queue
|
||||
to stop before interrupting the current thread. Defaults to 10000.
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
</xsd:attribute>
|
||||
<xsd:attribute name="wait-time-out">
|
||||
<xsd:annotation>
|
||||
<xsd:documentation>
|
||||
|
||||
@@ -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"/>
|
||||
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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");
|
||||
|
||||
Reference in New Issue
Block a user