INT-2832/2833 Aggregator Fix and Documentation

Expire empty groups.

Due to indentation changes, the code changes look more extensive
than they are. In effect the if (group.size() > 0) test is moved
to a narrower scope and the remove(group) is now performed if the
group is empty.

Document expire-groups-upon-completion.

INT-2832 Doc Polishing

PR Review + punctuation.

INT-2833 Add Delay For Expiring Empty Groups

minimumTimeoutForEmptyGroups
This commit is contained in:
Gary Russell
2012-11-26 20:24:58 -05:00
committed by Mark Fisher
parent 84f4bea5e7
commit 3ea05576b3
3 changed files with 193 additions and 49 deletions

View File

@@ -92,6 +92,8 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
private boolean lockRegistrySet = false;
private volatile long minimumTimeoutForEmptyGroups;
public AbstractCorrelatingMessageHandler(MessageGroupProcessor processor, MessageGroupStore store,
CorrelationStrategy correlationStrategy, ReleaseStrategy releaseStrategy) {
Assert.notNull(processor);
@@ -172,6 +174,21 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
this.sendPartialResultOnExpiry = sendPartialResultOnExpiry;
}
/**
* By default, when a MessageGroupStoreReaper is configured to expire partial
* groups, empty groups are also removed. Empty groups exist after a group
* is released normally. This is to enable the detection and discarding of
* late-arriving messages. If you wish to run empty group deletion on a longer
* schedule than expiring partial groups, set this property. Empty groups will
* then not be removed from the MessageStore until they have not been modified
* for at least this number of milliseconds.
*
* @param minimumTimeoutForEmptyGroups The minimum timeout.
*/
public void setMinimumTimeoutForEmptyGroups(long minimumTimeoutForEmptyGroups) {
this.minimumTimeoutForEmptyGroups = minimumTimeoutForEmptyGroups;
}
public void setReleasePartialSequences(boolean releasePartialSequences){
Assert.isInstanceOf(SequenceSizeReleaseStrategy.class, this.releaseStrategy,
"Release strategy of type [" + this.releaseStrategy.getClass().getSimpleName()
@@ -241,7 +258,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
*/
protected abstract void afterRelease(MessageGroup group, Collection<Message<?>> completedMessages);
private final boolean forceComplete(MessageGroup group) {
private void forceComplete(MessageGroup group) {
Object correlationKey = group.getGroupId();
// UUIDConverter is no-op if already converted
@@ -250,41 +267,47 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
try {
lock.lockInterruptibly();
try {
if (group.size() > 0) {
try {
/*
* Need to verify the group hasn't changed while we were waiting on
* its lock. We have to re-fetch the group for this. A possible
* future improvement would be to add MessageGroupStore.getLastModified(groupId).
*/
MessageGroup messageGroupNow = this.messageStore.getMessageGroup(
group.getGroupId());
long lastModifiedNow = messageGroupNow.getLastModified();
if (group.getLastModified() == lastModifiedNow) {
if (releaseStrategy.canRelease(group)) {
this.completeGroup(correlationKey, group);
}
else {
this.expireGroup(correlationKey, group);
}
/*
* Need to verify the group hasn't changed while we were waiting on
* its lock. We have to re-fetch the group for this. A possible
* future improvement would be to add MessageGroupStore.getLastModified(groupId).
*/
MessageGroup messageGroupNow = this.messageStore.getMessageGroup(
group.getGroupId());
long lastModifiedNow = messageGroupNow.getLastModified();
if (group.getLastModified() == lastModifiedNow) {
if (group.size() > 0) {
if (releaseStrategy.canRelease(group)) {
this.completeGroup(correlationKey, group);
}
else {
removeGroup = false;
if (logger.isDebugEnabled()) {
logger.debug("Group expiry candidate (" + group.getGroupId() +
") has changed - it may be reconsidered for a future expiration");
}
this.expireGroup(correlationKey, group);
}
}
finally {
if (removeGroup) {
this.remove(group);
else {
/*
* By default empty groups are removed on the same schedule as non-empty
* groups. A longer timeout for empty groups can be enabled by
* setting minimumTimeoutForEmptyGroups.
*/
removeGroup = lastModifiedNow < (System.currentTimeMillis() - this.minimumTimeoutForEmptyGroups);
if (removeGroup && logger.isDebugEnabled()) {
logger.debug("Removing empty group: " + group.getGroupId());
}
}
return true;
}
else {
removeGroup = false;
if (logger.isDebugEnabled()) {
logger.debug("Group expiry candidate (" + group.getGroupId() +
") has changed - it may be reconsidered for a future expiration");
}
}
}
finally {
if (removeGroup) {
this.remove(group);
}
lock.unlock();
}
}
@@ -292,7 +315,6 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageH
Thread.currentThread().interrupt();
throw new MessagingException("Thread was interrupted while trying to obtain lock");
}
return false;
}
void remove(MessageGroup group) {

View File

@@ -22,6 +22,7 @@ import static org.junit.Assert.assertTrue;
import java.util.ArrayList;
import java.util.Collection;
import java.util.List;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
@@ -34,6 +35,7 @@ import org.springframework.integration.store.MessageGroup;
import org.springframework.integration.store.MessageGroupStore;
import org.springframework.integration.store.SimpleMessageStore;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.test.util.TestUtils;
/**
* @author Gary Russell
@@ -150,4 +152,101 @@ public class AbstractCorrelatingMessageHandlerTests {
assertNull(discards.receive(0));
}
@Test // INT-2833
public void testReaperReapsAnEmptyGroup() throws Exception {
final MessageGroupStore groupStore = new SimpleMessageStore();
AggregatingMessageHandler handler = new AggregatingMessageHandler(
new MessageGroupProcessor() {
public Object processMessageGroup(MessageGroup group) {
return group;
}
}, groupStore) {
};
final List<Message<?>> outputMessages = new ArrayList<Message<?>>();
handler.setOutputChannel(new MessageChannel() {
/*
* Executes when group 'bar' completes normally
*/
public boolean send(Message<?> message, long timeout) {
outputMessages.add(message);
return true;
}
public boolean send(Message<?> message) {
return this.send(message, 0);
}
});
handler.setReleaseStrategy(new ReleaseStrategy() {
public boolean canRelease(MessageGroup group) {
return group.size() == 1;
}
});
Message<String> message = MessageBuilder.withPayload("foo")
.setCorrelationId("bar")
.build();
handler.handleMessage(message);
assertEquals(1, outputMessages.size());
assertEquals(1, TestUtils.getPropertyValue(handler, "messageStore.groupIdToMessageGroup", Map.class).size());
groupStore.expireMessageGroups(0);
assertEquals(0, TestUtils.getPropertyValue(handler, "messageStore.groupIdToMessageGroup", Map.class).size());
}
@Test // INT-2833
public void testReaperReapsAnEmptyGroupAfterConfiguredDelay() throws Exception {
final MessageGroupStore groupStore = new SimpleMessageStore();
AggregatingMessageHandler handler = new AggregatingMessageHandler(
new MessageGroupProcessor() {
public Object processMessageGroup(MessageGroup group) {
return group;
}
}, groupStore) {
};
final List<Message<?>> outputMessages = new ArrayList<Message<?>>();
handler.setOutputChannel(new MessageChannel() {
/*
* Executes when group 'bar' completes normally
*/
public boolean send(Message<?> message, long timeout) {
outputMessages.add(message);
return true;
}
public boolean send(Message<?> message) {
return this.send(message, 0);
}
});
handler.setReleaseStrategy(new ReleaseStrategy() {
public boolean canRelease(MessageGroup group) {
return group.size() == 1;
}
});
handler.setMinimumTimeoutForEmptyGroups(1000);
Message<String> message = MessageBuilder.withPayload("foo")
.setCorrelationId("bar")
.build();
handler.handleMessage(message);
assertEquals(1, outputMessages.size());
assertEquals(1, TestUtils.getPropertyValue(handler, "messageStore.groupIdToMessageGroup", Map.class).size());
groupStore.expireMessageGroups(0);
assertEquals(1, TestUtils.getPropertyValue(handler, "messageStore.groupIdToMessageGroup", Map.class).size());
Thread.sleep(1010);
groupStore.expireMessageGroups(0);
assertEquals(0, TestUtils.getPropertyValue(handler, "messageStore.groupIdToMessageGroup", Map.class).size());
}
}