INT-3999: Avoid "hard" References from Futures

JIRA: https://jira.spring.io/browse/INT-3999

Since the scheduled tasks may live for a long time it can finish
with the `OutOfMemory` if we use the direct reference to big objects, like `Message<?>`.

* Fix `AbstractCorrelatingMessageHandler` to deal only with the `groupId`
from the `scheduleGroupToForceComplete()` when we `schedule` `Runnable` for the `forceRelease` logic.
* Fix `DelayHandler` to deal only with `messageId` in the `releaseMessageAfterDelay()`, when we
`schedule` `Runnable` for the `releaseMessageAfterDelay`.
* Since the logic hasn't been changed for those components, there is no any new test.
There is just enough to be sure that all existing tests are fine.

**Cherry-pick to 4.0.x, 4.1.x, 4.2.x**

Optimise the release task for the `SimpleMessageStore` case
This commit is contained in:
Artem Bilan
2016-04-20 10:36:36 -04:00
committed by Gary Russell
parent 882f4c017e
commit 256a7292ef
2 changed files with 76 additions and 25 deletions

View File

@@ -426,38 +426,38 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
}
}
private void scheduleGroupToForceComplete(final MessageGroup messageGroup) {
final Long groupTimeout = this.obtainGroupTimeout(messageGroup);
private void scheduleGroupToForceComplete(MessageGroup messageGroup) {
final Long groupTimeout = obtainGroupTimeout(messageGroup);
/*
* When 'groupTimeout' is evaluated to 'null' we do nothing.
* The 'MessageGroupStoreReaper' can be used to 'forceComplete' message groups.
*/
if (groupTimeout != null && groupTimeout >= 0) {
if (groupTimeout > 0) {
final MessageGroupProcessor forceReleaseProcessor =
AbstractCorrelatingMessageHandler.this.forceReleaseProcessor;
ScheduledFuture<?> scheduledFuture = this.getTaskScheduler()
final Object groupId = messageGroup.getGroupId();
ScheduledFuture<?> scheduledFuture = getTaskScheduler()
.schedule(new Runnable() {
@Override
public void run() {
try {
forceReleaseProcessor.processMessageGroup(messageGroup);
processForceRelease(groupId);
}
catch (MessageDeliveryException e) {
if (logger.isDebugEnabled()) {
logger.debug("The MessageGroup [ " + messageGroup +
logger.debug("The MessageGroup [ " + groupId +
"] is rescheduled by the reason: " + e.getMessage());
}
scheduleGroupToForceComplete(messageGroup);
scheduleGroupToForceComplete(groupId);
}
}
}, new Date(System.currentTimeMillis() + groupTimeout));
if (logger.isDebugEnabled()) {
logger.debug("Schedule MessageGroup [ " + messageGroup + "] to 'forceComplete'.");
}
this.expireGroupScheduledFutures.put(UUIDConverter.getUUID(messageGroup.getGroupId()), scheduledFuture);
this.expireGroupScheduledFutures.put(UUIDConverter.getUUID(groupId), scheduledFuture);
}
else {
this.forceReleaseProcessor.processMessageGroup(messageGroup);
@@ -465,6 +465,16 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
}
}
private void scheduleGroupToForceComplete(Object groupId) {
MessageGroup messageGroup = this.messageStore.getMessageGroup(groupId);
scheduleGroupToForceComplete(messageGroup);
}
private void processForceRelease(Object groupId) {
MessageGroup messageGroup = this.messageStore.getMessageGroup(groupId);
this.forceReleaseProcessor.processMessageGroup(messageGroup);
}
private void discardMessage(Message<?> message) {
if (this.discardChannelName != null) {
synchronized (this) {

View File

@@ -20,6 +20,7 @@ import java.io.Serializable;
import java.util.Collection;
import java.util.Date;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.atomic.AtomicBoolean;
import org.aopalliance.aop.Advice;
@@ -30,8 +31,6 @@ import org.springframework.context.event.ContextRefreshedEvent;
import org.springframework.expression.EvaluationContext;
import org.springframework.expression.EvaluationException;
import org.springframework.expression.Expression;
import org.springframework.expression.ExpressionParser;
import org.springframework.expression.spel.standard.SpelExpressionParser;
import org.springframework.integration.context.IntegrationObjectSupport;
import org.springframework.integration.expression.ExpressionUtils;
import org.springframework.integration.store.MessageGroup;
@@ -84,8 +83,6 @@ import org.springframework.util.CollectionUtils;
public class DelayHandler extends AbstractReplyProducingMessageHandler implements DelayHandlerManagement,
ApplicationListener<ContextRefreshedEvent> {
private static final ExpressionParser expressionParser = new SpelExpressionParser();
private final String messageGroupId;
private volatile long defaultDelay;
@@ -287,12 +284,14 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
if (delayValueException != null) {
if (this.ignoreExpressionFailures) {
if (logger.isDebugEnabled()) {
logger.debug("Failed to get delay value from 'delayExpression': " + delayValueException.getMessage() +
logger.debug("Failed to get delay value from 'delayExpression': " +
delayValueException.getMessage() +
". Will fall back to default delay: " + this.defaultDelay);
}
}
else {
throw new MessageHandlingException(message, "Error occurred during 'delay' value determination", delayValueException);
throw new MessageHandlingException(message, "Error occurred during 'delay' value determination",
delayValueException);
}
}
@@ -309,20 +308,60 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
}
else {
messageWrapper = new DelayedMessageWrapper(message, System.currentTimeMillis());
delayedMessage = this.getMessageBuilderFactory().withPayload(messageWrapper).copyHeaders(message.getHeaders()).build();
delayedMessage = getMessageBuilderFactory()
.withPayload(messageWrapper)
.copyHeaders(message.getHeaders())
.build();
this.messageStore.addMessageToGroup(this.messageGroupId, delayedMessage);
}
final Message<?> messageToSchedule = delayedMessage;
this.getTaskScheduler().schedule(new Runnable() {
Runnable releaseTask;
@Override
public void run() {
releaseMessage(messageToSchedule);
if (this.messageStore instanceof SimpleMessageStore) {
final Message<?> messageToSchedule = delayedMessage;
releaseTask = new Runnable() {
@Override
public void run() {
releaseMessage(messageToSchedule);
}
};
}
else {
final UUID messageId = delayedMessage.getHeaders().getId();
releaseTask = new Runnable() {
@Override
public void run() {
Message<?> messageToRelease = getMessageById(messageId);
if (messageToRelease != null) {
releaseMessage(messageToRelease);
}
}
};
}
getTaskScheduler().schedule(releaseTask, new Date(messageWrapper.getRequestDate() + delay));
}
private Message<?> getMessageById(UUID messageId) {
Message<?> theMessage = ((MessageStore) this.messageStore).getMessage(messageId);
if (theMessage == null) {
if (logger.isDebugEnabled()) {
logger.debug("No message in the Message Store for id: " + messageId +
". Likely another instance has already released it.");
}
}, new Date(messageWrapper.getRequestDate() + delay));
return null;
}
else {
return theMessage;
}
}
private void releaseMessage(Message<?> message) {
@@ -371,17 +410,19 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement
* Used for reading persisted Messages in the 'messageStore'
* to reschedule them e.g. upon application restart.
* The logic is based on iteration over {@code messageGroup.getMessages()}
* and schedules task about 'delay' logic.
* and schedules task for 'delay' logic.
* This behavior is dictated by the avoidance of invocation thread overload.
*/
@Override
public synchronized void reschedulePersistedMessages() {
MessageGroup messageGroup = this.messageStore.getMessageGroup(this.messageGroupId);
for (final Message<?> message : messageGroup.getMessages()) {
this.getTaskScheduler().schedule(new Runnable() {
getTaskScheduler().schedule(new Runnable() {
@Override
public void run() {
// This is fine to keep the reference to the message,
// because the scheduled task is performed immediately.
long delay = determineDelayForMessage(message);
if (delay > 0) {
releaseMessageAfterDelay(message, delay);