diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/AggregatorParser.java b/spring-integration-core/src/main/java/org/springframework/integration/config/AggregatorParser.java index ce2d317cd2..9b6b884992 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/AggregatorParser.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/AggregatorParser.java @@ -16,6 +16,8 @@ package org.springframework.integration.config; +import org.w3c.dom.Element; + import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.RuntimeBeanReference; import org.springframework.beans.factory.parsing.BeanComponentDefinition; @@ -25,7 +27,6 @@ import org.springframework.beans.factory.xml.ParserContext; import org.springframework.integration.MessagingConfigurationException; import org.springframework.integration.router.AggregatingMessageHandler; import org.springframework.util.StringUtils; -import org.w3c.dom.Element; /** * Parser for the aggregator element of the integration namespace. @@ -69,6 +70,7 @@ public class AggregatorParser implements BeanDefinitionParser { public static final String TRACKED_CORRELATION_ID_CAPACITY_PROPERTY = "trackedCorrelationIdCapacity"; + public BeanDefinition parse(Element element, ParserContext parserContext) { return parseAggregatorElement(element, parserContext, true); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceUtils.java b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceUtils.java index 38ca60f781..2ddd6dc7d7 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceUtils.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/config/IntegrationNamespaceUtils.java @@ -39,6 +39,7 @@ public abstract class IntegrationNamespaceUtils { private static final String KEEP_ALIVE_ATTRIBUTE = "keep-alive"; + public static ConcurrencyPolicy parseConcurrencyPolicy(Element element) { ConcurrencyPolicy policy = new ConcurrencyPolicy(); String coreSize = element.getAttribute(CORE_SIZE_ATTRIBUTE); diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java index 7733d0ffd0..52537aebb7 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/AggregatorParserTests.java @@ -23,6 +23,7 @@ import java.util.List; import org.junit.Assert; import org.junit.Before; import org.junit.Test; + import org.springframework.context.ApplicationContext; import org.springframework.context.support.ClassPathXmlApplicationContext; import org.springframework.integration.channel.MessageChannel; @@ -40,6 +41,7 @@ public class AggregatorParserTests { private ApplicationContext context; + @Before public void setUp() { context = new ClassPathXmlApplicationContext("aggregatorParserTests.xml", this.getClass()); @@ -98,6 +100,7 @@ public class AggregatorParserTests { 99, getPropertyValue(completeAggregatingMessageHandler, "trackedCorrelationIdCapacity", int.class)); } + private static Message createMessage(String payload, Object correlationId, int sequenceSize, int sequenceNumber, MessageChannel replyChannel) { StringMessage message = new StringMessage(payload); @@ -120,4 +123,5 @@ public class AggregatorParserTests { ReflectionUtils.makeAccessible(field); return field.get(beanUnderTest); } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/TestAggregator.java b/spring-integration-core/src/test/java/org/springframework/integration/config/TestAggregator.java index 6caef00fdf..b4daefb50a 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/TestAggregator.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/TestAggregator.java @@ -14,13 +14,13 @@ * limitations under the License. */ - package org.springframework.integration.config; -import java.util.List; import java.util.ArrayList; import java.util.Collections; +import java.util.List; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import org.springframework.integration.message.Message; import org.springframework.integration.message.StringMessage; @@ -32,7 +32,8 @@ import org.springframework.integration.router.MessageSequenceComparator; */ public class TestAggregator implements Aggregator { - ConcurrentHashMap> aggregatedMessages = new ConcurrentHashMap>(); + private final ConcurrentMap> aggregatedMessages = new ConcurrentHashMap>(); + public Message aggregate(List> messages) { List> sortableList = new ArrayList>(messages); @@ -46,10 +47,8 @@ public class TestAggregator implements Aggregator { return returnedMessage; } - public ConcurrentHashMap> getAggregatedMessages() { + public ConcurrentMap> getAggregatedMessages() { return aggregatedMessages; } - + } - - diff --git a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml b/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml index 999bc131b2..adef845aa4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/config/aggregatorParserTests.xml @@ -9,8 +9,8 @@ - - + + + + +