INT-1119: remove reaper attributes from aggregator and resequencer
This commit is contained in:
@@ -31,7 +31,6 @@ 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;
|
||||
|
||||
/**
|
||||
@@ -112,17 +111,6 @@ 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;
|
||||
@@ -203,7 +191,6 @@ 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();
|
||||
|
||||
@@ -58,27 +58,9 @@ 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;
|
||||
|
||||
}
|
||||
|
||||
@@ -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.ReleaseStrategy;
|
||||
import org.springframework.integration.annotation.CorrelationStrategy;
|
||||
import org.springframework.integration.annotation.ReleaseStrategy;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
import org.springframework.integration.message.MessageHandler;
|
||||
import org.springframework.integration.store.SimpleMessageStore;
|
||||
@@ -66,9 +66,6 @@ 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;
|
||||
|
||||
@@ -30,6 +30,7 @@ import org.w3c.dom.Element;
|
||||
* @author Marius Bogoevici
|
||||
* @author Mark Fisher
|
||||
* @author Oleg Zhurakousky
|
||||
* @author Dave Syer
|
||||
*/
|
||||
public class AggregatorParser extends AbstractConsumerEndpointParser {
|
||||
|
||||
@@ -49,10 +50,6 @@ 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";
|
||||
@@ -97,10 +94,7 @@ 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);
|
||||
|
||||
@@ -62,9 +62,6 @@ 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;
|
||||
}
|
||||
|
||||
@@ -1719,9 +1719,6 @@
|
||||
</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>
|
||||
@@ -1766,9 +1763,6 @@
|
||||
</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>
|
||||
|
||||
Reference in New Issue
Block a user