Revert "INT-1119: remove reaper attributes from aggregator and resequencer"

This reverts commit 15fde633db17d982daf7e2cf6b5d657b1e0a87f9.
This commit is contained in:
David Syer
2010-05-05 15:16:59 +00:00
parent 3825e6e758
commit 97436b3f3b
11 changed files with 66 additions and 6 deletions

View File

@@ -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();

View File

@@ -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;
}

View File

@@ -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;

View File

@@ -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);

View File

@@ -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;
}

View File

@@ -1719,6 +1719,9 @@
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="send-partial-result-on-timeout" type="xsd:string" />
<xsd:attribute name="tracked-correlation-id-capacity" type="xsd:string" />
<xsd:attribute name="reaper-interval" type="xsd:string" />
<xsd:attribute name="timeout" type="xsd:string" />
</xsd:extension>
</xsd:complexContent>
</xsd:complexType>
@@ -1763,6 +1766,9 @@
</xsd:attribute>
<xsd:attribute name="release-partial-sequences" type="xsd:string" />
<xsd:attribute name="send-partial-result-on-timeout" type="xsd:string" />
<xsd:attribute name="tracked-correlation-id-capacity" type="xsd:string" />
<xsd:attribute name="reaper-interval" type="xsd:string" />
<xsd:attribute name="timeout" type="xsd:string" />
</xsd:extension>
</xsd:complexContent>
</xsd:complexType>

View File

@@ -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);

View File

@@ -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);

View File

@@ -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"/>
<channel id="aggregatorWithReferenceAndMethodInput"/>
<aggregator id="aggregatorWithReferenceAndMethod"

View File

@@ -40,8 +40,8 @@ public class TestAnnotatedEndpointWithCustomizedAggregator {
inputChannel = "inputChannel",
outputChannel = "outputChannel",
discardChannel = "discardChannel",
sendPartialResultsOnTimeout = true,
sendTimeout = 98765432)
reaperInterval = 1234, sendPartialResultsOnTimeout = true,
sendTimeout = 98765432, timeout = 4567890, trackedCorrelationIdCapacity = 42)
public Message<?> aggregatingMethod(List<Message<?>> messages) {
List<Message<?>> sortableList = new ArrayList<Message<?>>(messages);
Collections.sort(sortableList, new MessageSequenceComparator());

View File

@@ -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"/>
<resequencer id="resequencerWithCorrelationStrategyRefOnly"