INT-4377: aggregator groupTimeout as Date (#3505)
* INT-4377: aggregator groupTimeout as Date JIRA: https://jira.spring.io/browse/INT-4377 Change the `groupTimeoutExpression` logic to let it to be evaluated to `Date` instance for some fine-grained scheduling use-case, e.g. to determine a scheduling moment from the group creation time (`timestamp`) instead of a current message arrival * Fix language in docs accoridng PR review Co-authored-by: Gary Russell <grussell@vmware.com> Co-authored-by: Gary Russell <grussell@vmware.com>
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2002-2020 the original author or authors.
|
||||
* Copyright 2002-2021 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.
|
||||
@@ -626,16 +626,24 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
}
|
||||
|
||||
private void scheduleGroupToForceComplete(MessageGroup messageGroup) {
|
||||
final Long groupTimeout = obtainGroupTimeout(messageGroup);
|
||||
Object groupTimeout = obtainGroupTimeout(messageGroup);
|
||||
/*
|
||||
* When 'groupTimeout' is evaluated to 'null' we do nothing.
|
||||
* The 'MessageGroupStoreReaper' can be used to 'forceComplete' message groups.
|
||||
*/
|
||||
if (groupTimeout != null) {
|
||||
if (groupTimeout > 0) {
|
||||
final Object groupId = messageGroup.getGroupId();
|
||||
final long timestamp = messageGroup.getTimestamp();
|
||||
final long lastModified = messageGroup.getLastModified();
|
||||
Date startTime = null;
|
||||
if (groupTimeout instanceof Date) {
|
||||
startTime = (Date) groupTimeout;
|
||||
}
|
||||
else if ((Long) groupTimeout > 0) {
|
||||
startTime = new Date(System.currentTimeMillis() + (Long) groupTimeout);
|
||||
}
|
||||
|
||||
if (startTime != null) {
|
||||
Object groupId = messageGroup.getGroupId();
|
||||
long timestamp = messageGroup.getTimestamp();
|
||||
long lastModified = messageGroup.getLastModified();
|
||||
ScheduledFuture<?> scheduledFuture =
|
||||
getTaskScheduler()
|
||||
.schedule(() -> {
|
||||
@@ -643,14 +651,12 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
processForceRelease(groupId, timestamp, lastModified);
|
||||
}
|
||||
catch (MessageDeliveryException ex) {
|
||||
if (AbstractCorrelatingMessageHandler.this.logger.isWarnEnabled()) {
|
||||
AbstractCorrelatingMessageHandler.this.logger.warn(ex,
|
||||
() -> "The MessageGroup [" + groupId
|
||||
+ "] is rescheduled by the reason of:");
|
||||
}
|
||||
logger.warn(ex, () ->
|
||||
"The MessageGroup [" + groupId +
|
||||
"] is rescheduled by the reason of: ");
|
||||
scheduleGroupToForceComplete(groupId);
|
||||
}
|
||||
}, new Date(System.currentTimeMillis() + groupTimeout));
|
||||
}, startTime);
|
||||
|
||||
this.logger.debug(() -> "Schedule MessageGroup [ " + messageGroup + "] to 'forceComplete'.");
|
||||
this.expireGroupScheduledFutures.put(UUIDConverter.getUUID(groupId), scheduledFuture);
|
||||
@@ -905,9 +911,22 @@ public abstract class AbstractCorrelatingMessageHandler extends AbstractMessageP
|
||||
"The expected collection of Messages contains non-Message element: " + commonElementType);
|
||||
}
|
||||
|
||||
protected Long obtainGroupTimeout(MessageGroup group) {
|
||||
return this.groupTimeoutExpression != null
|
||||
? this.groupTimeoutExpression.getValue(this.evaluationContext, group, Long.class) : null;
|
||||
protected Object obtainGroupTimeout(MessageGroup group) {
|
||||
if (this.groupTimeoutExpression != null) {
|
||||
Object timeout = this.groupTimeoutExpression.getValue(this.evaluationContext, group);
|
||||
if (timeout instanceof Date) {
|
||||
return timeout;
|
||||
}
|
||||
else if (timeout != null) {
|
||||
try {
|
||||
return Long.parseLong(timeout.toString());
|
||||
}
|
||||
catch (NumberFormatException ex) {
|
||||
throw new IllegalStateException("Error evaluating 'groupTimeoutExpression'", ex);
|
||||
}
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -107,6 +107,8 @@ public abstract class CorrelationHandlerSpec<S extends CorrelationHandlerSpec<S,
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a SpEL expression to evaluate the group timeout for scheduled expiration.
|
||||
* Must return {@link java.util.Date}, {@link java.lang.Long} or {@link String} as a long.
|
||||
* @param groupTimeoutExpression the group timeout expression string.
|
||||
* @return the handler spec.
|
||||
* @see AbstractCorrelatingMessageHandler#setGroupTimeoutExpression
|
||||
@@ -122,11 +124,12 @@ public abstract class CorrelationHandlerSpec<S extends CorrelationHandlerSpec<S,
|
||||
* based on the message group.
|
||||
* Usually used with a JDK8 lambda:
|
||||
* <p>{@code .groupTimeout(g -> g.size() * 2000L)}.
|
||||
* Must return {@link java.util.Date}, {@link java.lang.Long} or {@link String} a long.
|
||||
* @param groupTimeoutFunction a function invoked to resolve the group timeout in milliseconds.
|
||||
* @return the handler spec.
|
||||
* @see AbstractCorrelatingMessageHandler#setGroupTimeoutExpression
|
||||
*/
|
||||
public S groupTimeout(Function<MessageGroup, Long> groupTimeoutFunction) {
|
||||
public S groupTimeout(Function<MessageGroup, ?> groupTimeoutFunction) {
|
||||
this.handler.setGroupTimeoutExpression(new FunctionExpression<>(groupTimeoutFunction));
|
||||
return _this();
|
||||
}
|
||||
|
||||
@@ -3913,6 +3913,7 @@
|
||||
the MessageGroup won't be scheduled to be forced complete.
|
||||
The action taken when the group is forced complete depends on the
|
||||
'send-partial-result-on-expiry' attribute.
|
||||
Can be evaluated directly to 'java.util.Date' instance for a scheduled task.
|
||||
Mutually exclusive with the 'group-timeout' attribute.
|
||||
</xsd:documentation>
|
||||
</xsd:annotation>
|
||||
|
||||
Reference in New Issue
Block a user