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 22f7c49556..fa6184a7f3 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -579,7 +579,7 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP afterRelease(messageGroup, completedMessages); } if (!isExpireGroupsUponCompletion() && this.minimumTimeoutForEmptyGroups > 0) { - removeEmptyGroupAfterTimeout(messageGroup, this.minimumTimeoutForEmptyGroups); + removeEmptyGroupAfterTimeout(groupIdUuid, this.minimumTimeoutForEmptyGroups); } } else { @@ -616,29 +616,27 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP return false; } - private void removeEmptyGroupAfterTimeout(MessageGroup messageGroup, long timeout) { - Object groupId = messageGroup.getGroupId(); - UUID groupUuid = UUIDConverter.getUUID(groupId); + private void removeEmptyGroupAfterTimeout(UUID groupId, long timeout) { ScheduledFuture scheduledFuture = getTaskScheduler() .schedule(() -> { - Lock lock = this.lockRegistry.obtain(groupUuid.toString()); + Lock lock = this.lockRegistry.obtain(groupId.toString()); try { lock.lockInterruptibly(); try { - this.expireGroupScheduledFutures.remove(groupUuid); + this.expireGroupScheduledFutures.remove(groupId); /* * Obtain a fresh state for group from the MessageStore, * since it could be changed while we have waited for lock. */ - MessageGroup groupNow = this.messageStore.getMessageGroup(groupUuid); + MessageGroup groupNow = this.messageStore.getMessageGroup(groupId); boolean removeGroup = groupNow.size() == 0 && groupNow.getLastModified() <= (System.currentTimeMillis() - this.minimumTimeoutForEmptyGroups); if (removeGroup) { - this.logger.debug(() -> "Removing empty group: " + groupUuid); - remove(messageGroup); + this.logger.debug(() -> "Removing empty group: " + groupId); + remove(groupNow); } } finally { @@ -649,13 +647,13 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP Thread.currentThread().interrupt(); this.logger.debug(() -> "Thread was interrupted while trying to obtain lock." + "Rescheduling empty MessageGroup [ " + groupId + "] for removal."); - removeEmptyGroupAfterTimeout(messageGroup, timeout); + removeEmptyGroupAfterTimeout(groupId, timeout); } }, new Date(System.currentTimeMillis() + timeout)); this.logger.debug(() -> "Schedule empty MessageGroup [ " + groupId + "] for removal."); - this.expireGroupScheduledFutures.put(groupUuid, scheduledFuture); + this.expireGroupScheduledFutures.put(groupId, scheduledFuture); } private void scheduleGroupToForceComplete(MessageGroup messageGroup) { 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 1f0fd45f44..1b72b2d6d8 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 @@ -1,5 +1,5 @@ /* - * Copyright 2002-2021 the original author or authors. + * Copyright 2002-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -425,24 +425,7 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement this.messageStore.addMessageToGroup(this.messageGroupId, delayedMessage); } - Runnable releaseTask; - - if (this.messageStore instanceof SimpleMessageStore) { - final Message messageToSchedule = delayedMessage; - - releaseTask = () -> releaseMessage(messageToSchedule); - } - else { - final UUID messageId = delayedMessage.getHeaders().getId(); - - releaseTask = () -> { - Message messageToRelease = getMessageById(messageId); - if (messageToRelease != null) { - releaseMessage(messageToRelease); - } - }; - } - + Runnable releaseTask = releaseTaskForMessage(delayedMessage); Date startTime = new Date(messageWrapper.getRequestDate() + delay); if (TransactionSynchronizationManager.isSynchronizationActive() && @@ -463,6 +446,21 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement } } + private Runnable releaseTaskForMessage(Message delayedMessage) { + if (this.messageStore instanceof SimpleMessageStore) { + return () -> releaseMessage(delayedMessage); + } + else { + UUID messageId = delayedMessage.getHeaders().getId(); + return () -> { + Message messageToRelease = getMessageById(messageId); + if (messageToRelease != null) { + releaseMessage(messageToRelease); + } + }; + } + } + private Message getMessageById(UUID messageId) { Message theMessage = ((MessageStore) this.messageStore).getMessage(messageId); @@ -531,9 +529,9 @@ public class DelayHandler extends AbstractReplyProducingMessageHandler implement rescheduleAt(message, new Date()); } - protected void rescheduleAt(final Message message, Date startTime) { - getTaskScheduler() - .schedule(() -> releaseMessage(message), startTime); + protected void rescheduleAt(Message message, Date startTime) { + Runnable releaseTask = releaseTaskForMessage(message); + getTaskScheduler().schedule(releaseTask, startTime); } private void doReleaseMessage(Message message) {