diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index 9f95c14183..0dbf89cfd3 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -13,9 +13,11 @@ package org.springframework.integration.aggregator; -import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.Comparator; +import java.util.Date; +import java.util.HashMap; import java.util.List; import java.util.concurrent.locks.Lock; @@ -63,6 +65,7 @@ import org.springframework.util.CollectionUtils; * @author Dave Syer * @author Oleg Zhurakousky * @author Gary Russell + * @author Enrique Rodríguez * @since 2.0 */ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageHandler implements MessageProducer { @@ -71,6 +74,8 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH public static final long DEFAULT_SEND_TIMEOUT = 1000L; + private final Comparator> sequenceNumberComparator = new SequenceNumberComparator(); + protected volatile MessageGroupStore messageStore; private final MessageGroupProcessor outputProcessor; @@ -346,11 +351,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH } protected int findLastReleasedSequenceNumber(Object groupId, Collection> partialSequence){ - List> sorted = new ArrayList>(partialSequence); - Collections.sort(sorted, new SequenceNumberComparator()); - - Message lastReleasedMessage = sorted.get(partialSequence.size()-1); - + Message lastReleasedMessage = Collections.max(partialSequence, this.sequenceNumberComparator); return lastReleasedMessage.getHeaders().getSequenceNumber(); } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/SequenceSizeReleaseStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/SequenceSizeReleaseStrategy.java index e36d632559..6d6796e700 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/SequenceSizeReleaseStrategy.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/SequenceSizeReleaseStrategy.java @@ -16,11 +16,9 @@ package org.springframework.integration.aggregator; -import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.Comparator; -import java.util.List; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -36,12 +34,13 @@ import org.springframework.integration.store.MessageGroup; * @author Dave Syer * @author Iwein Fuld * @author Oleg Zhurakousky + * @author Enrique Rodríguez */ public class SequenceSizeReleaseStrategy implements ReleaseStrategy { private static final Log logger = LogFactory.getLog(SequenceSizeReleaseStrategy.class); - private volatile Comparator> comparator = new SequenceNumberComparator(); + private final Comparator> comparator = new SequenceNumberComparator(); private volatile boolean releasePartialSequences; @@ -74,10 +73,8 @@ public class SequenceSizeReleaseStrategy implements ReleaseStrategy { if (logger.isTraceEnabled()) { logger.trace("Considering partial release of group [" + messageGroup + "]"); } - List> sorted = new ArrayList>(messages); - Collections.sort(sorted, comparator); - - int nextSequenceNumber = sorted.get(0).getHeaders().getSequenceNumber(); + Message minMessage = Collections.min(messages, this.comparator); + int nextSequenceNumber = minMessage.getHeaders().getSequenceNumber(); int lastReleasedMessageSequence = messageGroup.getLastReleasedMessageSequenceNumber(); if (nextSequenceNumber - lastReleasedMessageSequence == 1){ diff --git a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java index a43aca54fc..62c73d8a2c 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/channel/AbstractMessageChannel.java @@ -17,6 +17,7 @@ package org.springframework.integration.channel; import java.util.Collections; +import java.util.Comparator; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; @@ -49,12 +50,14 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport im protected final Log logger = LogFactory.getLog(this.getClass()); + private final ChannelInterceptorList interceptors = new ChannelInterceptorList(); + + private final Comparator orderComparator = new OrderComparator(); + private volatile boolean shouldTrack = false; private volatile Class[] datatypes = new Class[] { Object.class }; - private final ChannelInterceptorList interceptors = new ChannelInterceptorList(); - private volatile String fullChannelName; @@ -87,7 +90,7 @@ public abstract class AbstractMessageChannel extends IntegrationObjectSupport im * interceptors. */ public void setInterceptors(List interceptors) { - Collections.sort(interceptors, new OrderComparator()); + Collections.sort(interceptors, this.orderComparator); this.interceptors.set(interceptors); }