diff --git a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java index d878df218c..f2f0a34811 100644 --- a/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java +++ b/org.springframework.integration/src/main/java/org/springframework/integration/aggregator/Resequencer.java @@ -49,6 +49,7 @@ import org.springframework.util.CollectionUtils; * will be ignored. * * @author Marius Bogoevici + * @author Alex Peters */ public class Resequencer extends AbstractMessageBarrierHandler>> { @@ -77,7 +78,7 @@ public class Resequencer extends AbstractMessageBarrierHandler> releasedMessages = releaseAvailableMessages(barrier); if (!CollectionUtils.isEmpty(releasedMessages)) { Message lastMessage = releasedMessages.get(releasedMessages.size()-1); - if (lastMessage.getHeaders().getSequenceNumber().equals(lastMessage.getHeaders().getSequenceSize() - 1)) { + if (lastMessage.getHeaders().getSequenceNumber().equals(lastMessage.getHeaders().getSequenceSize())) { this.removeBarrier(barrier.getCorrelationKey()); } this.sendReplies(releasedMessages, this.resolveReplyChannelFromMessage(releasedMessages.get(0))); @@ -134,7 +135,7 @@ public class Resequencer extends AbstractMessageBarrierHandler message1 = createMessage("123", correlationId, 1, 1, + replyChannel); + resequencer.handleMessage(message1); + assertThat(resequencer.barriers.containsKey(correlationId), is(false)); + } private static Message createMessage(String payload, Object correlationId, int sequenceSize, int sequenceNumber, MessageChannel replyChannel) { - Message message = MessageBuilder.withPayload(payload) + return MessageBuilder.withPayload(payload) .setCorrelationId(correlationId) .setSequenceSize(sequenceSize) .setSequenceNumber(sequenceNumber) .setReplyChannel(replyChannel) .build(); - return message; } @After