INT-602 Resequencer returning spurious messages. Added get/setAttributes to MessageBarrier so that information can be held by it as part of its state.
This commit is contained in:
@@ -359,8 +359,23 @@ public abstract class AbstractMessageBarrierHandler<T extends Collection<? exten
|
||||
* @param barrier the {@link MessageBarrier} to be processed
|
||||
*/
|
||||
protected abstract void processBarrier(MessageBarrier<T> barrier);
|
||||
|
||||
/**
|
||||
* A method for discarding the content of the message barrier.
|
||||
* Can be overridden by subclasses.
|
||||
* @param entry
|
||||
* @param barrier
|
||||
*/
|
||||
protected void discardBarrier(MessageBarrier<T> barrier) {
|
||||
for (Message message : barrier.getMessages()) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Handling of Message group with correlation key '" + barrier.getCorrelationKey()+ "' has timed out.");
|
||||
}
|
||||
discardMessage(message);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
/**
|
||||
* A task that runs periodically, pruning the timed-out message barriers.
|
||||
*/
|
||||
private class PrunerTask implements Runnable {
|
||||
@@ -377,13 +392,7 @@ public abstract class AbstractMessageBarrierHandler<T extends Collection<? exten
|
||||
processBarrier(barrier);
|
||||
}
|
||||
else {
|
||||
for (Message message : barrier.getMessages()) {
|
||||
if (logger.isDebugEnabled()) {
|
||||
logger.debug("Handling of Message group with correlationId '" + entry.getKey()
|
||||
+ "' has timed out.");
|
||||
}
|
||||
discardMessage(message);
|
||||
}
|
||||
discardBarrier(barrier);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
|
||||
package org.springframework.integration.aggregator;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Collection;
|
||||
|
||||
@@ -30,6 +31,8 @@ import org.springframework.integration.core.Message;
|
||||
* been added to it). This is a parameterized type, allowing different different
|
||||
* client classes to use different types of Collections and their respective features.
|
||||
*
|
||||
* Can store/retrieve attributes through its setAttribute() and getAttribute() methods.
|
||||
*
|
||||
* This class is not thread-safe and will be synchronized by the calling code.
|
||||
*
|
||||
* @author Marius Bogoevici
|
||||
@@ -43,6 +46,8 @@ public class MessageBarrier<T extends Collection<? extends Message>> {
|
||||
private Object correlationKey;
|
||||
|
||||
private final long timestamp = System.currentTimeMillis();
|
||||
|
||||
private final Map<String, Object> attributes = new HashMap<String, Object>();
|
||||
|
||||
public MessageBarrier(T messages, Object correlationKey) {
|
||||
this.messages = messages;
|
||||
@@ -81,5 +86,22 @@ public class MessageBarrier<T extends Collection<? extends Message>> {
|
||||
public T getMessages() {
|
||||
return this.messages;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sets a the value of a given attribute on the MessageBarrier.
|
||||
* @param attributeName
|
||||
* @param value
|
||||
*/
|
||||
public void setAttribute(String attributeName, Object value) {
|
||||
this.attributes.put(attributeName, value);
|
||||
}
|
||||
|
||||
/**
|
||||
* Gets the value of a given attribute from the MessageBarrier.
|
||||
* @param attributeName
|
||||
*/
|
||||
public <V> V getAttribute(String attributeName) {
|
||||
return (V)this.attributes.get(attributeName);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -23,7 +23,6 @@ import java.util.SortedSet;
|
||||
import java.util.TreeSet;
|
||||
|
||||
import org.springframework.integration.core.Message;
|
||||
import org.springframework.integration.message.MessageBuilder;
|
||||
import org.springframework.util.CollectionUtils;
|
||||
|
||||
/**
|
||||
@@ -37,10 +36,16 @@ import org.springframework.util.CollectionUtils;
|
||||
* All considerations regarding <code>timeout</code> and grouping by
|
||||
* '<code>correlationId</code>' from {@link AbstractMessageBarrierHandler}
|
||||
* apply here as well.
|
||||
*
|
||||
* It is assumed that all messages have the same <code>sequence_size</code> header attribute
|
||||
* and that the sequence numbers of the messages are successive, starting with
|
||||
* 1 up to <code>sequenceSize</code>. Messages that do not satisfy this condition are
|
||||
* considered out-of-sequence and thus rejected.
|
||||
*
|
||||
*
|
||||
* Note: messages with the same sequence number will be treated as equivalent
|
||||
* by this class (i.e. after a message with a given sequence number is received,
|
||||
* further messages from withing the same group, that have the same sequence number,
|
||||
* further messages from within the same group, that have the same sequence number,
|
||||
* will be ignored.
|
||||
*
|
||||
* @author Marius Bogoevici
|
||||
@@ -48,7 +53,9 @@ import org.springframework.util.CollectionUtils;
|
||||
public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message<?>>> {
|
||||
|
||||
private volatile boolean releasePartialSequences = true;
|
||||
|
||||
|
||||
private static final String LAST_RELEASED_SEQUENCE_NUMBER = "last.released.sequence.number";
|
||||
|
||||
|
||||
public void setReleasePartialSequences(boolean releasePartialSequences) {
|
||||
this.releasePartialSequences = releasePartialSequences;
|
||||
@@ -58,13 +65,13 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
protected MessageBarrier<SortedSet<Message<?>>> createMessageBarrier(Object correlationKey) {
|
||||
MessageBarrier<SortedSet<Message<?>>> messageBarrier
|
||||
= new MessageBarrier<SortedSet<Message<?>>>(new TreeSet<Message<?>>(new MessageSequenceComparator()), correlationKey);
|
||||
messageBarrier.getMessages().add(createFlagMessage(0));
|
||||
messageBarrier.setAttribute(LAST_RELEASED_SEQUENCE_NUMBER, 0);
|
||||
return messageBarrier;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void processBarrier(MessageBarrier<SortedSet<Message<?>>> barrier) {
|
||||
if (hasReceivedAllMessages(barrier.getMessages())) {
|
||||
if (hasReceivedAllMessages(barrier)) {
|
||||
barrier.setComplete();
|
||||
}
|
||||
List<Message<?>> releasedMessages = releaseAvailableMessages(barrier);
|
||||
@@ -77,21 +84,18 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
}
|
||||
}
|
||||
|
||||
private boolean hasReceivedAllMessages(SortedSet<Message<?>> messages) {
|
||||
Message<?> firstMessage = messages.first();
|
||||
Message<?> lastMessage = messages.last();
|
||||
return (lastMessage.getHeaders().getSequenceNumber().equals(lastMessage.getHeaders().getSequenceSize())
|
||||
&& (lastMessage.getHeaders().getSequenceNumber() - firstMessage.getHeaders().getSequenceNumber() == messages.size() - 1));
|
||||
private boolean hasReceivedAllMessages(MessageBarrier<SortedSet<Message<?>>> barrier) {
|
||||
int sequenceSize = barrier.getMessages().first().getHeaders().getSequenceSize();
|
||||
int messagesCurrentlyInBarrier = barrier.getMessages().size();
|
||||
int lastReleasedSequenceNumber = barrier.getAttribute(LAST_RELEASED_SEQUENCE_NUMBER);
|
||||
return (lastReleasedSequenceNumber + messagesCurrentlyInBarrier == sequenceSize);
|
||||
}
|
||||
|
||||
private List<Message<?>> releaseAvailableMessages(MessageBarrier<SortedSet<Message<?>>> barrier) {
|
||||
if (this.releasePartialSequences || barrier.isComplete()) {
|
||||
ArrayList<Message<?>> releasedMessages = new ArrayList<Message<?>>();
|
||||
Iterator<Message<?>> it = barrier.getMessages().iterator();
|
||||
//remove the initial flag from the list
|
||||
Message<?> flag = it.next();
|
||||
it.remove();
|
||||
int lastReleasedSequenceNumber = flag.getHeaders().getSequenceNumber();
|
||||
int lastReleasedSequenceNumber = barrier.getAttribute(LAST_RELEASED_SEQUENCE_NUMBER);
|
||||
while (it.hasNext()) {
|
||||
Message<?> currentMessage = it.next();
|
||||
if (lastReleasedSequenceNumber == currentMessage.getHeaders().getSequenceNumber() - 1) {
|
||||
@@ -103,8 +107,7 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
break;
|
||||
}
|
||||
}
|
||||
//re-insert the flag so that we know where to start releasing next
|
||||
barrier.getMessages().add(createFlagMessage(lastReleasedSequenceNumber));
|
||||
barrier.setAttribute(LAST_RELEASED_SEQUENCE_NUMBER, lastReleasedSequenceNumber);
|
||||
return releasedMessages;
|
||||
}
|
||||
else {
|
||||
@@ -113,28 +116,24 @@ public class Resequencer extends AbstractMessageBarrierHandler<SortedSet<Message
|
||||
}
|
||||
|
||||
@Override
|
||||
protected boolean canAddMessage(Message<?> message,
|
||||
MessageBarrier<SortedSet<Message<?>>> barrier) {
|
||||
protected boolean canAddMessage(Message<?> message, MessageBarrier<SortedSet<Message<?>>> barrier) {
|
||||
if (!super.canAddMessage(message, barrier)) {
|
||||
return false;
|
||||
}
|
||||
Message<?> flagMessage = barrier.getMessages().first();
|
||||
int lastReleasedSequenceNumber = barrier.getAttribute(LAST_RELEASED_SEQUENCE_NUMBER);
|
||||
if (barrier.messages.contains(message)
|
||||
|| flagMessage.getHeaders().getSequenceNumber() >= message.getHeaders().getSequenceNumber()) {
|
||||
|| lastReleasedSequenceNumber >= message.getHeaders().getSequenceNumber()) {
|
||||
logger.debug("A message with the same sequence number has been already received: " + message);
|
||||
return false;
|
||||
}
|
||||
Message<?> lastMessage = barrier.getMessages().last();
|
||||
if (lastMessage != flagMessage
|
||||
&& lastMessage.getHeaders().getSequenceSize() < message.getHeaders().getSequenceNumber()) {
|
||||
// one can always add a message to the barrier if it's empty. Afterwards, assume that the complete sequence size
|
||||
//
|
||||
if (!barrier.getMessages().isEmpty() &&
|
||||
barrier.getMessages().first().getHeaders().getSequenceSize() < message.getHeaders().getSequenceNumber()) {
|
||||
logger.debug("The message has a sequence number which is larger than the sequence size: "+ message);
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
private static Message<Integer> createFlagMessage(int sequenceNumber) {
|
||||
return MessageBuilder.withPayload(sequenceNumber).setSequenceNumber(sequenceNumber).build();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -124,6 +124,37 @@ public class ResequencerTests {
|
||||
assertNotNull(reply4);
|
||||
assertEquals(new Integer(4), reply4.getHeaders().getSequenceNumber());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testResequencingWithDiscard() throws InterruptedException {
|
||||
QueueChannel discardChannel = new QueueChannel();
|
||||
Message<?> message1 = createMessage("123", "ABC", 4, 2, null);
|
||||
Message<?> message2 = createMessage("456", "ABC", 4, 1, null);
|
||||
Message<?> message3 = createMessage("789", "ABC", 4, 4, null);
|
||||
Message<?> message4 = createMessage("XYZ", "ABC", 4, 3, null);
|
||||
this.resequencer.setSendPartialResultOnTimeout(false);
|
||||
this.resequencer.setReleasePartialSequences(false);
|
||||
this.resequencer.setDiscardChannel(discardChannel);
|
||||
this.resequencer.setTimeout(90000);
|
||||
this.resequencer.handleMessage(message1);
|
||||
this.resequencer.handleMessage(message2);
|
||||
this.resequencer.discardBarrier(this.resequencer.barriers.get("ABC"));
|
||||
Message<?> reply1 = discardChannel.receive(0);
|
||||
Message<?> reply2 = discardChannel.receive(0);
|
||||
Message<?> reply3 = discardChannel.receive(0);
|
||||
// only messages 1 and 2 should have been received by now
|
||||
assertNotNull(reply1);
|
||||
assertEquals(new Integer(1), reply1.getHeaders().getSequenceNumber());
|
||||
assertNotNull(reply2);
|
||||
assertEquals(new Integer(2), reply2.getHeaders().getSequenceNumber());
|
||||
assertNull(reply3);
|
||||
// when sending the last message, the whole sequence must have been sent
|
||||
this.resequencer.handleMessage(message4);
|
||||
reply3 = discardChannel.receive(0);
|
||||
assertNull(reply3);
|
||||
Message<?> reply4 = discardChannel.receive(0);
|
||||
assertNull(reply4);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
|
||||
Reference in New Issue
Block a user