INT-661: Applied patch uploaded in Jira
This commit is contained in:
@@ -49,6 +49,7 @@ import org.springframework.util.CollectionUtils;
|
||||
* will be ignored.
|
||||
*
|
||||
* @author Marius Bogoevici
|
||||
* @author Alex Peters
|
||||
*/
|
||||
public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message<?>>> {
|
||||
|
||||
@@ -77,7 +78,7 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
List<Message<?>> 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<SortedSet<Message
|
||||
return false;
|
||||
}
|
||||
if (!barrier.getMessages().isEmpty() &&
|
||||
message.getHeaders().getSequenceSize() != barrier.getMessages().first().getHeaders().getSequenceSize()) {
|
||||
! message.getHeaders().getSequenceSize().equals(barrier.getMessages().first().getHeaders().getSequenceSize())) {
|
||||
logger.debug("The message has a sequence size which is different from other messages handled so far: " + message
|
||||
+ ", expected value is " + barrier.getMessages().first().getHeaders().getSequenceNumber());
|
||||
return false;
|
||||
|
||||
@@ -16,14 +16,15 @@
|
||||
|
||||
package org.springframework.integration.aggregator;
|
||||
|
||||
import static org.hamcrest.CoreMatchers.is;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertThat;
|
||||
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.core.MessageChannel;
|
||||
@@ -33,6 +34,7 @@ import org.springframework.integration.util.TestUtils;
|
||||
|
||||
/**
|
||||
* @author Marius Bogoevici
|
||||
* @author Alex Peters
|
||||
*/
|
||||
public class ResequencerTests {
|
||||
|
||||
@@ -223,17 +225,26 @@ public class ResequencerTests {
|
||||
assertNotNull(reply4);
|
||||
assertEquals(new Integer(4), reply4.getHeaders().getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRemovalOfBarrierWhenLastMessageOfSequenceArrives() {
|
||||
QueueChannel replyChannel = new QueueChannel();
|
||||
String correlationId = "ABC";
|
||||
Message<?> 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<String> message = MessageBuilder.withPayload(payload)
|
||||
return MessageBuilder.withPayload(payload)
|
||||
.setCorrelationId(correlationId)
|
||||
.setSequenceSize(sequenceSize)
|
||||
.setSequenceNumber(sequenceNumber)
|
||||
.setReplyChannel(replyChannel)
|
||||
.build();
|
||||
return message;
|
||||
}
|
||||
|
||||
@After
|
||||
|
||||
Reference in New Issue
Block a user