diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java
index ecee0a38ae..7f0a8fe94e 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/CorrelatingMessageHandler.java
@@ -31,6 +31,7 @@ import org.springframework.integration.store.MessageGroupCallback;
import org.springframework.integration.store.MessageGroupStore;
import org.springframework.integration.store.MessageStore;
import org.springframework.integration.store.SimpleMessageStore;
+import org.springframework.scheduling.TaskScheduler;
import org.springframework.util.Assert;
/**
@@ -111,6 +112,17 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
this.releaseStrategy = releaseStrategy;
}
+ public void setTaskScheduler(TaskScheduler taskScheduler) {
+ super.setTaskScheduler(taskScheduler);
+ }
+
+ // TODO: INT-958 - remove unused property setters
+ public void setTimeout(long timeout) {
+ }
+
+ public void setReaperInterval(long reaperInterval) {
+ }
+
public void setOutputChannel(MessageChannel outputChannel) {
Assert.notNull(outputChannel, "'outputChannel' must not be null");
this.outputChannel = outputChannel;
@@ -191,6 +203,7 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
}
+ // TODO: INT-958 - arrange for this to be called if user desires, e.g. periodically
private final boolean forceComplete(MessageGroup group) {
Object correlationKey = group.getCorrelationKey();
diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Aggregator.java b/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Aggregator.java
index af26a1bfc5..a8d30d2345 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Aggregator.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/annotation/Aggregator.java
@@ -58,9 +58,27 @@ public @interface Aggregator {
*/
long sendTimeout() default CorrelatingMessageHandler.DEFAULT_SEND_TIMEOUT;
+ /**
+ * maximum time to wait for completion (in milliseconds)
+ */
+ long timeout() default CorrelatingMessageHandler.DEFAULT_TIMEOUT;
+
/**
* indicates whether to send an incomplete aggregate on timeout
*/
boolean sendPartialResultsOnTimeout() default false;
+ /**
+ * interval for the task that checks for timed-out aggregates
+ */
+ long reaperInterval() default CorrelatingMessageHandler.DEFAULT_REAPER_INTERVAL;
+
+ /**
+ * maximum number of correlation IDs to maintain so that received messages
+ * may be recognized as belonging to an aggregate that has already completed
+ * or timed out
+ */
+ // TODO: INT-958 - remove / deal with tracked id capacity
+ int trackedCorrelationIdCapacity() default 42;
+
}
diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java
index a63c361074..dfa2f104c6 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/annotation/AggregatorAnnotationPostProcessor.java
@@ -22,13 +22,13 @@ import java.util.concurrent.atomic.AtomicReference;
import org.springframework.beans.factory.ListableBeanFactory;
import org.springframework.core.annotation.AnnotationUtils;
+import org.springframework.integration.aggregator.ReleaseStrategyAdapter;
import org.springframework.integration.aggregator.CorrelatingMessageHandler;
import org.springframework.integration.aggregator.CorrelationStrategyAdapter;
import org.springframework.integration.aggregator.MethodInvokingMessageGroupProcessor;
-import org.springframework.integration.aggregator.ReleaseStrategyAdapter;
import org.springframework.integration.annotation.Aggregator;
-import org.springframework.integration.annotation.CorrelationStrategy;
import org.springframework.integration.annotation.ReleaseStrategy;
+import org.springframework.integration.annotation.CorrelationStrategy;
import org.springframework.integration.core.MessageChannel;
import org.springframework.integration.message.MessageHandler;
import org.springframework.integration.store.SimpleMessageStore;
@@ -66,6 +66,9 @@ public class AggregatorAnnotationPostProcessor extends AbstractMethodAnnotationP
}
handler.setSendTimeout(annotation.sendTimeout());
handler.setSendPartialResultOnTimeout(annotation.sendPartialResultsOnTimeout());
+ handler.setReaperInterval(annotation.reaperInterval());
+ handler.setTimeout(annotation.timeout());
+ // handler.setTrackedCorrelationIdCapacity(annotation.trackedCorrelationIdCapacity());
handler.setBeanFactory(this.beanFactory);
handler.afterPropertiesSet();
return handler;
diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AggregatorParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AggregatorParser.java
index ad39b259d6..2b267e9b85 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AggregatorParser.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/AggregatorParser.java
@@ -30,7 +30,6 @@ import org.w3c.dom.Element;
* @author Marius Bogoevici
* @author Mark Fisher
* @author Oleg Zhurakousky
- * @author Dave Syer
*/
public class AggregatorParser extends AbstractConsumerEndpointParser {
@@ -50,6 +49,10 @@ public class AggregatorParser extends AbstractConsumerEndpointParser {
private static final String SEND_PARTIAL_RESULT_ON_TIMEOUT_ATTRIBUTE = "send-partial-result-on-timeout";
+ private static final String REAPER_INTERVAL_ATTRIBUTE = "reaper-interval";
+
+ private static final String TIMEOUT_ATTRIBUTE = "timeout";
+
private static final String RELEASE_STRATEGY_PROPERTY = "releaseStrategy";
private static final String CORRELATION_STRATEGY_PROPERTY = "correlationStrategy";
@@ -94,7 +97,10 @@ public class AggregatorParser extends AbstractConsumerEndpointParser {
SEND_TIMEOUT_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
SEND_PARTIAL_RESULT_ON_TIMEOUT_ATTRIBUTE);
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element,
+ REAPER_INTERVAL_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, TIMEOUT_ATTRIBUTE);
this.injectPropertyWithBean(RELEASE_STRATEGY_REF_ATTRIBUTE,
RELEASE_STRATEGY_METHOD_ATTRIBUTE, RELEASE_STRATEGY_PROPERTY,
"ReleaseStrategyAdapter", element, builder, parserContext);
diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java
index c9a1ac5fd8..ea8d1b1c4e 100644
--- a/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java
+++ b/org.springframework.integration/src/main/java/org/springframework/integration/config/xml/ResequencerParser.java
@@ -62,6 +62,9 @@ public class ResequencerParser extends AbstractConsumerEndpointParser {
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "discard-channel");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-timeout");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "send-partial-result-on-timeout");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reaper-interval");
+ // IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "tracked-correlation-id-capacity");
+ IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "timeout");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup");
return builder;
}
diff --git a/org.springframework.integration/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd b/org.springframework.integration/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd
index fe61961acb..75b55f2f9c 100644
--- a/org.springframework.integration/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd
+++ b/org.springframework.integration/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd
@@ -1719,6 +1719,9 @@
+
+
+
@@ -1763,6 +1766,9 @@
+
+
+
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java
index 961d5f16fd..ade0854b1e 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/AggregatorTests.java
@@ -69,6 +69,8 @@ public class AggregatorTests {
@Test
public void testShouldNotSendPartialResultOnTimeoutByDefault() throws InterruptedException {
QueueChannel discardChannel = new QueueChannel();
+ this.aggregator.setTimeout(50);
+ this.aggregator.setReaperInterval(10);
this.aggregator.setDiscardChannel(discardChannel);
QueueChannel replyChannel = new QueueChannel();
Message> message = createMessage(3, "ABC", 2, 1, replyChannel, null);
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java
index 06d2a33e99..71f81f07eb 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/aggregator/ResequencerTests.java
@@ -133,6 +133,7 @@ public class ResequencerTests {
this.resequencer.setSendPartialResultOnTimeout(false);
this.processor.setReleasePartialSequences(false);
this.resequencer.setDiscardChannel(discardChannel);
+ this.resequencer.setTimeout(90000);
this.resequencer.handleMessage(message1);
this.resequencer.handleMessage(message2);
assertEquals(1, store.expireMessageGroups(-10000));
@@ -161,6 +162,7 @@ public class ResequencerTests {
this.resequencer.setSendPartialResultOnTimeout(false);
this.processor.setReleasePartialSequences(false);
this.resequencer.setDiscardChannel(discardChannel);
+ this.resequencer.setTimeout(90000);
this.resequencer.handleMessage(message1);
this.resequencer.handleMessage(message2);
// this.resequencer.discardBarrier(this.resequencer.barriers.get("ABC"));
@@ -179,6 +181,7 @@ public class ResequencerTests {
this.resequencer.setSendPartialResultOnTimeout(false);
this.processor.setReleasePartialSequences(false);
this.resequencer.setDiscardChannel(discardChannel);
+ this.resequencer.setTimeout(90000);
this.resequencer.handleMessage(message1);
// this.resequencer.discardBarrier(this.resequencer.barriers.get("ABC"));
Message> reply1 = discardChannel.receive(0);
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml b/org.springframework.integration/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml
index c216cd080d..025c6e034f 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml
@@ -27,7 +27,10 @@
release-strategy="releaseStrategy"
correlation-strategy="correlationStrategy"
send-timeout="86420000"
- send-partial-result-on-timeout="true"/>
+ send-partial-result-on-timeout="true"
+ reaper-interval="135"
+ tracked-correlation-id-capacity="99"
+ timeout="42"/>
aggregatingMethod(List> messages) {
List> sortableList = new ArrayList>(messages);
Collections.sort(sortableList, new MessageSequenceComparator());
diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/resequencerParserTests.xml b/org.springframework.integration/src/test/java/org/springframework/integration/config/resequencerParserTests.xml
index bf839f9574..9bf8624b39 100644
--- a/org.springframework.integration/src/test/java/org/springframework/integration/config/resequencerParserTests.xml
+++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/resequencerParserTests.xml
@@ -31,6 +31,9 @@
discard-channel="discardChannel"
send-timeout="86420000"
send-partial-result-on-timeout="true"
+ reaper-interval="135"
+ tracked-correlation-id-capacity="99"
+ timeout="42"
release-partial-sequences="false"/>