From 256a7292ef49150786a2a11e71aefba9b03ec344 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Wed, 20 Apr 2016 10:36:36 -0400 Subject: [PATCH] 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 --- .../AbstractCorrelatingMessageHandler.java | 28 ++++--- .../integration/handler/DelayHandler.java | 73 +++++++++++++++---- 2 files changed, 76 insertions(+), 25 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java index 493d42c31c..0508bb13ff 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/aggregator/AbstractCorrelatingMessageHandler.java @@ -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) { diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java index 64eeaa6b28..fad0f8717f 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/DelayHandler.java @@ -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 { - 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);