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 bae437ad7d..5e5d9de7aa 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 @@ -34,23 +34,28 @@ import org.springframework.integration.store.SimpleMessageStore; import org.springframework.util.Assert; /** - * Message handler that holds a buffer of correlated messages in a {@link MessageStore}. This class takes care of - * correlated groups of messages that can be completed in batches. It is useful for aggregating, resequencing, or custom - * implementations requiring correlation. + * Message handler that holds a buffer of correlated messages in a + * {@link MessageStore}. This class takes care of correlated groups of messages + * that can be completed in batches. It is useful for aggregating, resequencing, + * or custom implementations requiring correlation. *

- * To customize this handler inject {@link CorrelationStrategy}, {@link ReleaseStrategy}, and - * {@link MessageGroupProcessor} implementations as you require. + * To customize this handler inject {@link CorrelationStrategy}, + * {@link ReleaseStrategy}, and {@link MessageGroupProcessor} implementations as + * you require. *

- * By default the CorrelationStrategy will be a HeaderAttributeCorrelationStrategy and the ReleaseStrategy will be a + * By default the CorrelationStrategy will be a + * HeaderAttributeCorrelationStrategy and the ReleaseStrategy will be a * SequenceSizeReleaseStrategy. * * @author Iwein Fuld * @author Dave Syer * @since 2.0 */ -public class CorrelatingMessageHandler extends AbstractMessageHandler implements MessageProducer { +public class CorrelatingMessageHandler extends AbstractMessageHandler implements + MessageProducer { - private static final Log logger = LogFactory.getLog(CorrelatingMessageHandler.class); + private static final Log logger = LogFactory + .getLog(CorrelatingMessageHandler.class); public static final long DEFAULT_SEND_TIMEOUT = 1000L; @@ -76,24 +81,23 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements private final ConcurrentMap locks = new ConcurrentHashMap(); - public CorrelatingMessageHandler(MessageGroupProcessor processor, MessageGroupStore store, - CorrelationStrategy correlationStrategy, ReleaseStrategy releaseStrategy) { + public CorrelatingMessageHandler(MessageGroupProcessor processor, + MessageGroupStore store, CorrelationStrategy correlationStrategy, + ReleaseStrategy releaseStrategy) { Assert.notNull(store); Assert.notNull(processor); - this.messageStore = store; - store.registerMessageGroupExpiryCallback(new MessageGroupCallback() { - public void execute(MessageGroup group) { - forceComplete(group); - } - }); + setMessageStore(store); this.outputProcessor = processor; this.correlationStrategy = correlationStrategy == null ? new HeaderAttributeCorrelationStrategy( - MessageHeaders.CORRELATION_ID) : correlationStrategy; - this.releaseStrategy = releaseStrategy == null ? new SequenceSizeReleaseStrategy() : releaseStrategy; + MessageHeaders.CORRELATION_ID) + : correlationStrategy; + this.releaseStrategy = releaseStrategy == null ? new SequenceSizeReleaseStrategy() + : releaseStrategy; this.channelTemplate.setSendTimeout(DEFAULT_SEND_TIMEOUT); } - public CorrelatingMessageHandler(MessageGroupProcessor processor, MessageGroupStore store) { + public CorrelatingMessageHandler(MessageGroupProcessor processor, + MessageGroupStore store) { this(processor, store, null, null); } @@ -101,8 +105,13 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements this(processor, new SimpleMessageStore(0), null, null); } - public void setMessageStore(MessageGroupStore messageStore) { - this.messageStore = messageStore; + public void setMessageStore(MessageGroupStore store) { + this.messageStore = store; + store.registerMessageGroupExpiryCallback(new MessageGroupCallback() { + public void execute(MessageGroup group) { + forceComplete(group); + } + }); } public void setCorrelationStrategy(CorrelationStrategy correlationStrategy) { @@ -146,7 +155,8 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements Object correlationKey = correlationStrategy.getCorrelationKey(message); if (logger.isDebugEnabled()) { - logger.debug("Handling message with correlationKey [" + correlationKey + "]: " + message); + logger.debug("Handling message with correlationKey [" + + correlationKey + "]: " + message); } // TODO: INT-1117 - make the lock global? @@ -162,19 +172,24 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements if (releaseStrategy.canRelease(group)) { if (logger.isDebugEnabled()) { - logger.debug("Completing group with correlationKey [" + correlationKey + "]"); + logger.debug("Completing group with correlationKey [" + + correlationKey + "]"); } try { - outputProcessor.processAndSend(group, channelTemplate, this.resolveReplyChannel(message, - this.outputChannel)); + outputProcessor.processAndSend(group, channelTemplate, + this.resolveReplyChannel(message, + this.outputChannel)); } finally { - // Always clean up even if there was an exception processing messages + // Always clean up even if there was an exception + // processing messages if (group.isComplete() || group.getSequenceSize() == 0) { - // The group is complete or else there is no sequence so there is no more state to track + // The group is complete or else there is no + // sequence so there is no more state to track remove(group); } else { - // Mark these messages as processed, but do not remove the group from store + // Mark these messages as processed, but do not + // remove the group from store mark(group); } @@ -183,7 +198,8 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements } else if (group.isComplete()) { try { - // If not releasing any messages the group might still be complete + // If not releasing any messages the group might still + // be complete for (Message discard : group.getUnmarked()) { discardChannel.send(discard); } @@ -211,20 +227,28 @@ public class CorrelatingMessageHandler extends AbstractMessageHandler implements // last chance for normal completion try { if (releaseStrategy.canRelease(group)) { - outputProcessor.processAndSend(group, channelTemplate, resolveReplyChannel(group.getOne(), - this.outputChannel)); + outputProcessor.processAndSend(group, channelTemplate, + resolveReplyChannel(group.getOne(), + this.outputChannel)); } else { if (sendPartialResultOnExpiry) { if (logger.isInfoEnabled()) { - logger.info("Processing partially complete messages for key [" + correlationKey - + "] to: " + outputChannel); + logger + .info("Processing partially complete messages for key [" + + correlationKey + + "] to: " + + outputChannel); } - outputProcessor.processAndSend(group, channelTemplate, resolveReplyChannel(group.getOne(), - this.outputChannel)); + outputProcessor.processAndSend(group, + channelTemplate, resolveReplyChannel(group + .getOne(), this.outputChannel)); } else { if (logger.isInfoEnabled()) { - logger.info("Discarding partially complete messages for key [" + correlationKey - + "] to: " + discardChannel); + logger + .info("Discarding partially complete messages for key [" + + correlationKey + + "] to: " + + discardChannel); } for (Message message : group.getUnmarked()) { discardChannel.send(message); diff --git a/org.springframework.integration/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java b/org.springframework.integration/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java index 6c70db05b0..8706492606 100644 --- a/org.springframework.integration/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java +++ b/org.springframework.integration/src/test/java/org/springframework/integration/config/AggregatorWithMessageStoreParserTests.java @@ -65,6 +65,22 @@ public class AggregatorWithMessageStoreParserTests { } + @Test + public void testExpiry() { + + input.send(createMessage("123", "id1", 3, 1, null)); + assertEquals(1, messageGroupStore.getMessageGroup("id1").size()); + input.send(createMessage("456", "id1", 3, 2, null)); + assertEquals(2, messageGroupStore.getMessageGroup("id1").size()); + messageGroupStore.expireMessageGroups(-10000); + assertEquals("One and only one message should have been aggregated", 1, aggregatorBean + .getAggregatedMessages().size()); + Message aggregatedMessage = aggregatorBean.getAggregatedMessages().get("id1"); + assertEquals("The aggregated message payload is not correct", "123456789", aggregatedMessage + .getPayload()); + } + + private static Message createMessage(T payload, Object correlationId, int sequenceSize, int sequenceNumber, MessageChannel outputChannel) { return MessageBuilder.withPayload(payload)