From 9a553178f5572c573e5a31885e2b45fd0abd75db Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 5 Nov 2025 11:20:45 +0200 Subject: [PATCH] return default value if no args updated --- .../single/AggIntervalEntryStatus.java | 25 ++--- ...EntityAggregationCalculatedFieldState.java | 93 ++++++++++--------- .../utils/CalculatedFieldArgumentUtils.java | 2 +- 3 files changed, 59 insertions(+), 61 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java index 2d350946e2..fbf344e5d3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntryStatus.java @@ -15,42 +15,31 @@ */ package org.thingsboard.server.service.cf.ctx.state.aggregation.single; +import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; -import lombok.Setter; @Data @NoArgsConstructor @AllArgsConstructor public class AggIntervalEntryStatus { - @Setter private long lastArgsRefreshTs = -1; - @Setter + private long lastMetricsEvalTs = -1; public AggIntervalEntryStatus(long lastArgsRefreshTs) { this.lastArgsRefreshTs = lastArgsRefreshTs; } - public boolean shouldRecalculate(long checkInterval) { - boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - checkInterval; - boolean argsUpdatedDuringInterval = lastArgsRefreshTs > -1; - if (intervalPassed && argsUpdatedDuringInterval) { - setLastMetricsEvalTs(System.currentTimeMillis()); - setLastArgsRefreshTs(-1); - return true; - } - return false; + public boolean intervalPassed(long checkInterval) { + return lastMetricsEvalTs <= System.currentTimeMillis() - checkInterval; } - public boolean intervalPassed(long checkInterval) { - boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - checkInterval; - if (intervalPassed) { - lastMetricsEvalTs = System.currentTimeMillis(); - } - return intervalPassed; + @JsonIgnore + public boolean argsUpdated() { + return lastArgsRefreshTs > -1; } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java index 4e7f073192..5e45b97f9a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java @@ -59,43 +59,6 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt super(entityId); } - public void fillMissingIntervals() { - long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); - long intervalDuration = interval.getIntervalDurationMillis(); - Map> intervals = getIntervals(); - AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); - if (lastIntervalEntry == null) { - return; - } - - long nextStartTs = lastIntervalEntry.getEndTs(); - long nextEndTs = nextStartTs + intervalDuration; - - while (nextEndTs <= currentIntervalEndTs) { - AggIntervalEntry missingAggIntervalEntry = new AggIntervalEntry(nextStartTs, nextEndTs); - - arguments.forEach((argName, argumentEntry) -> { - var entityAggEntry = (EntityAggregationArgumentEntry) argumentEntry; - AggIntervalEntryStatus intervalEntryStatus = new AggIntervalEntryStatus(System.currentTimeMillis()); - entityAggEntry.getAggIntervals().computeIfAbsent(missingAggIntervalEntry, missingInterval -> intervalEntryStatus); - }); - - nextStartTs = nextEndTs; - nextEndTs += intervalDuration; - } - } - - private Map> getIntervals() { - Map> intervals = new HashMap<>(); - arguments.forEach((argName, entry) -> { - var argEntry = (EntityAggregationArgumentEntry) entry; - argEntry.getAggIntervals().forEach((intervalEntry, status) -> - intervals.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(argName, status) - ); - }); - return intervals; - } - @Override public void setCtx(CalculatedFieldCtx ctx, TbActorRef actorCtx) { super.setCtx(ctx, actorCtx); @@ -104,7 +67,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt intervalDuration = configuration.getInterval().getIntervalDurationMillis(); Watermark watermark = configuration.getWatermark(); watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); - checkInterval = ctx.getAggCheckInterval(); + checkInterval = TimeUnit.SECONDS.toMillis(ctx.getAggCheckInterval()); interval = configuration.getInterval(); metrics = configuration.getMetrics(); } @@ -121,8 +84,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results = new HashMap<>(); List expiredIntervals = new ArrayList<>(); - Map> intervals = getIntervals(); - intervals.forEach((intervalEntry, argIntervalStatuses) -> { + getIntervals().forEach((intervalEntry, argIntervalStatuses) -> { processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); }); removeExpiredIntervals(expiredIntervals); @@ -157,6 +119,43 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt }); } + public void fillMissingIntervals() { + long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); + long intervalDuration = interval.getIntervalDurationMillis(); + Map> intervals = getIntervals(); + AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); + if (lastIntervalEntry == null) { + return; + } + + long nextStartTs = lastIntervalEntry.getEndTs(); + long nextEndTs = nextStartTs + intervalDuration; + + while (nextEndTs <= currentIntervalEndTs) { + AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs); + + arguments.forEach((argName, argumentEntry) -> { + var entityAggEntry = (EntityAggregationArgumentEntry) argumentEntry; + AggIntervalEntryStatus intervalEntryStatus = new AggIntervalEntryStatus(System.currentTimeMillis()); + entityAggEntry.getAggIntervals().computeIfAbsent(missing, missingInterval -> intervalEntryStatus); + }); + + nextStartTs = nextEndTs; + nextEndTs += intervalDuration; + } + } + + private Map> getIntervals() { + Map> intervals = new HashMap<>(); + arguments.forEach((argName, entry) -> { + var argEntry = (EntityAggregationArgumentEntry) entry; + argEntry.getAggIntervals().forEach((intervalEntry, status) -> + intervals.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(argName, status) + ); + }); + return intervals; + } + private void processInterval(long now, AggIntervalEntry intervalEntry, Map args, @@ -179,6 +178,9 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { processMetric(intervalEntry, argName, false, results); + } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); + processMetric(intervalEntry, argName, true, results); } }); } @@ -187,8 +189,15 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map args, Map> results) { args.forEach((argName, argEntryIntervalStatus) -> { - if (argEntryIntervalStatus.shouldRecalculate(checkInterval)) { - processMetric(intervalEntry, argName, false, results); + if (argEntryIntervalStatus.intervalPassed(checkInterval)) { + if (argEntryIntervalStatus.argsUpdated()) { + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); + argEntryIntervalStatus.setLastArgsRefreshTs(-1); + processMetric(intervalEntry, argName, false, results); + } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); + processMetric(intervalEntry, argName, true, results); + } } }); } diff --git a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java index 107b528957..97d53ac252 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -78,7 +78,7 @@ public class CalculatedFieldArgumentUtils { public static ArgumentEntry createDefaultMetricArgumentEntry(String argKey, AggMetric metric) { Long defaultValue = metric.getDefaultValue(); if (defaultValue != null) { - ArgumentEntry.createSingleValueArgument(new DoubleDataEntry(argKey, defaultValue.doubleValue())); + return ArgumentEntry.createSingleValueArgument(new DoubleDataEntry(argKey, defaultValue.doubleValue())); } return new SingleValueArgumentEntry(); }