|
|
@ -123,7 +123,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt |
|
|
intervals.forEach((intervalEntry, argIntervalStatuses) -> { |
|
|
intervals.forEach((intervalEntry, argIntervalStatuses) -> { |
|
|
processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); |
|
|
processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); |
|
|
}); |
|
|
}); |
|
|
expiredIntervals.forEach(intervals::remove); |
|
|
removeExpiredIntervals(expiredIntervals); |
|
|
|
|
|
|
|
|
ArrayNode result = toResult(results); |
|
|
ArrayNode result = toResult(results); |
|
|
if (result.isEmpty()) { |
|
|
if (result.isEmpty()) { |
|
|
@ -146,6 +146,15 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt |
|
|
}); |
|
|
}); |
|
|
} |
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
private void removeExpiredIntervals(List<AggIntervalEntry> expiredIntervals) { |
|
|
|
|
|
expiredIntervals.forEach(expiredInterval -> { |
|
|
|
|
|
arguments.values().stream() |
|
|
|
|
|
.map(EntityAggregationArgumentEntry.class::cast) |
|
|
|
|
|
.forEach(arg -> arg.getAggIntervals().remove(expiredInterval)); |
|
|
|
|
|
intervals.remove(expiredInterval); |
|
|
|
|
|
}); |
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
private void createIntervalIfNotExist() { |
|
|
private void createIntervalIfNotExist() { |
|
|
AggIntervalEntry currentInterval = new AggIntervalEntry(interval.getCurrentIntervalStartTs(), interval.getCurrentIntervalEndTs()); |
|
|
AggIntervalEntry currentInterval = new AggIntervalEntry(interval.getCurrentIntervalStartTs(), interval.getCurrentIntervalEndTs()); |
|
|
if (intervals.containsKey(currentInterval)) { |
|
|
if (intervals.containsKey(currentInterval)) { |
|
|
|