INT-1525: Add treatment for releasePartialSequences to ResequencerParser and CorrelatingMessageHandler

This commit is contained in:
Iwein Fuld
2010-10-15 17:20:58 +02:00
parent d68c9ba12b
commit 0ddc1fe667
4 changed files with 23 additions and 2 deletions

View File

@@ -142,6 +142,13 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
this.sendPartialResultOnExpiry = sendPartialResultOnExpiry;
}
public void setReleasePartialSequences(boolean releasePartialSequences){
Assert.isInstanceOf(SequenceSizeReleaseStrategy.class, this.releaseStrategy,
"Release strategy of type [" + this.releaseStrategy.getClass().getSimpleName()
+ "] cannot release partial sequences. Use the default SequenceSizeReleaseStrategy instead.");
((SequenceSizeReleaseStrategy)this.releaseStrategy).setReleasePartialSequences(releasePartialSequences);
}
@Override
public String getComponentType() {
return "aggregator";
@@ -162,6 +169,9 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements
synchronized (lock) {
MessageGroup group = messageStore.getMessageGroup(correlationKey);
if (group.canAdd(message)) {
if (logger.isTraceEnabled()) {
logger.trace("Adding message to group [ " + group + "]");
}
group = store(correlationKey, message);
if (releaseStrategy.canRelease(group)) {
Collection<Message> completedMessages = null;

View File

@@ -26,6 +26,7 @@ import org.w3c.dom.Element;
*
* @author Marius Bogoevici
* @author Dave Syer
* @author Iwein Fuld
*/
public class ResequencerParser extends AbstractConsumerEndpointParser {
@@ -84,6 +85,7 @@ public class ResequencerParser extends AbstractConsumerEndpointParser {
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, DISCARD_CHANNEL_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, SEND_TIMEOUT_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, SEND_PARTIAL_RESULT_ON_EXPIRY_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, RELEASE_PARTIAL_SEQUENCES_ATTRIBUTE);
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "auto-startup");
return builder;
}

View File

@@ -77,7 +77,7 @@ public class ResequencerParserTests {
"The ResequencerEndpoint is not configured with the appropriate 'send partial results on timeout' flag",
true, getPropertyValue(resequencer, "sendPartialResultOnExpiry"));
assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag",
false, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences"));
true, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences"));
}
@Test
@@ -90,6 +90,15 @@ public class ResequencerParserTests {
.getBean("testCorrelationStrategy"), getPropertyValue(resequencer, "correlationStrategy"));
}
@Test
public void shouldSetReleasePartialSequencesFlag(){
EventDrivenConsumer endpoint = (EventDrivenConsumer) context.getBean("completelyDefinedResequencer");
CorrelatingMessageHandler resequencer = TestUtils.getPropertyValue(endpoint, "handler",
CorrelatingMessageHandler.class);
assertEquals("The ResequencerEndpoint is not configured with the appropriate 'release partial sequences' flag",
true, getPropertyValue(getPropertyValue(resequencer, "releaseStrategy"), "releasePartialSequences"));
}
@Test
public void testCorrelationStrategyRefAndMethod() throws Exception {
EventDrivenConsumer endpoint = (EventDrivenConsumer) context

View File

@@ -35,7 +35,7 @@
discard-channel="discardChannel"
send-timeout="86420000"
send-partial-result-on-expiry="true"
release-partial-sequences="false"/>
release-partial-sequences="true"/>
<resequencer id="resequencerWithCorrelationStrategyRefOnly"
input-channel="inputChannel3"