From d0a00c68fde6649eaf49694c8dd276f68e5186fe Mon Sep 17 00:00:00 2001 From: Iwein Fuld Date: Sun, 17 Oct 2010 20:00:08 +0200 Subject: [PATCH] INT-1500: add apply-sequence flag to splitter. --- .../config/SplitterFactoryBean.java | 8 ++++++ .../config/xml/SplitterParser.java | 9 +++++++ .../splitter/AbstractMessageSplitter.java | 13 ++++++++- .../config/xml/spring-integration-2.0.xsd | 9 +++++++ .../router/config/SplitterParserTests.java | 27 ++++++++++++++----- .../router/config/splitterParserTests.xml | 7 +++++ 6 files changed, 65 insertions(+), 8 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java index a158603081..92a4736230 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/SplitterFactoryBean.java @@ -29,6 +29,7 @@ import org.springframework.util.StringUtils; * Factory bean for creating a Message Splitter. * * @author Mark Fisher + * @author Iwein Fuld */ public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { @@ -36,6 +37,8 @@ public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { private volatile boolean requiresReply; + private volatile boolean applySequence = true; + public void setSendTimeout(Long sendTimeout) { this.sendTimeout = sendTimeout; @@ -49,6 +52,10 @@ public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { this.requiresReply = requiresReply; } + public void setApplySequence(boolean applySequence) { + this.applySequence = applySequence; + } + @Override MessageHandler createMethodInvokingHandler(Object targetObject, String targetMethodName) { Assert.notNull(targetObject, "targetObject must not be null"); @@ -89,6 +96,7 @@ public class SplitterFactoryBean extends AbstractMessageHandlerFactoryBean { splitter.setSendTimeout(sendTimeout); } splitter.setRequiresReply(requiresReply); + splitter.setApplySequence(applySequence); return splitter; } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/SplitterParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/SplitterParser.java index f28969c749..95e7530367 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/xml/SplitterParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/xml/SplitterParser.java @@ -16,10 +16,15 @@ package org.springframework.integration.config.xml; +import org.springframework.beans.factory.support.BeanDefinitionBuilder; +import org.springframework.beans.factory.xml.ParserContext; +import org.w3c.dom.Element; + /** * Parser for the <splitter/> element. * * @author Mark Fisher + * @author Iwein Fuld */ public class SplitterParser extends AbstractDelegatingConsumerEndpointParser { @@ -33,4 +38,8 @@ public class SplitterParser extends AbstractDelegatingConsumerEndpointParser { return true; } + @Override + void postProcess(BeanDefinitionBuilder builder, Element element, ParserContext parserContext) { + IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "apply-sequence"); + } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java b/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java index 90a216011f..ad2ca47016 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/splitter/AbstractMessageSplitter.java @@ -35,6 +35,15 @@ import org.springframework.util.ObjectUtils; */ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMessageHandler { + private boolean applySequence = true; + + /** + * Set the applySequence flag to the specified value. Defaults to true. + */ + public void setApplySequence(boolean applySequence) { + this.applySequence = applySequence; + } + @Override @SuppressWarnings("unchecked") protected final Object handleRequestMessage(Message message) { @@ -80,7 +89,9 @@ public abstract class AbstractMessageSplitter extends AbstractReplyProducingMess builder = MessageBuilder.withPayload(item); builder.copyHeaders(headers); } - builder.pushSequenceDetails(correlationId, sequenceNumber, sequenceSize); + if (this.applySequence) { + builder.pushSequenceDetails(correlationId, sequenceNumber, sequenceSize); + } return builder; } diff --git a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd index e86350c029..001c8bff53 100644 --- a/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd +++ b/spring-integration-core/src/main/resources/org/springframework/integration/config/xml/spring-integration-2.0.xsd @@ -2118,6 +2118,15 @@ Name of the header whose value to use. + + + + Set this flag to false to prevent adding sequence related headers in this splitter. This + can be convenient in cases where the set sequence numbers conflict with downstream custom + aggregations. + + + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java index 8e13d1a5d6..1bddf63cc8 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/config/SplitterParserTests.java @@ -16,25 +16,24 @@ package org.springframework.integration.router.config; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNull; - -import java.util.Collections; - import org.junit.Test; - import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; import org.springframework.integration.MessageHandlingException; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.core.PollableChannel; -import org.springframework.integration.endpoint.EventDrivenConsumer; import org.springframework.integration.message.GenericMessage; import org.springframework.integration.support.MessageBuilder; +import java.util.Collections; + +import static org.hamcrest.Matchers.is; +import static org.junit.Assert.*; + /** * @author Mark Fisher + * @author Iwein Fuld */ public class SplitterParserTests { @@ -104,4 +103,18 @@ public class SplitterParserTests { inputChannel.send(MessageBuilder.withPayload(Collections.emptyList()).build()); } + @Test + public void splitterParserTestApplySequenceFalse() { + ClassPathXmlApplicationContext context = new ClassPathXmlApplicationContext( + "splitterParserTests.xml", this.getClass()); + context.start(); + DirectChannel inputChannel = context.getBean("noSequenceInput", DirectChannel.class); + PollableChannel output = (PollableChannel) context.getBean("output"); + inputChannel.send(MessageBuilder.withPayload(Collections.emptyList()).build()); + Message message = output.receive(1000); + assertThat(message.getHeaders().getSequenceNumber(), is(0)); + assertThat(message.getHeaders().getSequenceSize(), is(0)); + } + + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml index aa96b3df2b..adfa3dfce2 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/router/config/splitterParserTests.xml @@ -32,6 +32,13 @@ output-channel="output" requires-reply="true"/> + +