From cb278ad93fae7c74fe2bec01a5137418a6688da1 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 3 Nov 2025 15:50:44 +0200 Subject: [PATCH] minor fixes --- .../single/EntityAggregationArgumentEntry.java | 1 - .../EntityAggregationCalculatedFieldState.java | 11 +++++++---- .../utils/CalculatedFieldArgumentUtils.java | 16 +++++++++++----- .../cf/configuration/aggregation/AggMetric.java | 2 +- 4 files changed, 19 insertions(+), 11 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java index c1d1b5d807..2723b35cc9 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationArgumentEntry.java @@ -53,7 +53,6 @@ public class EntityAggregationArgumentEntry implements ArgumentEntry { for (Map.Entry aggIntervalEntry : aggIntervals.entrySet()) { if (aggIntervalEntry.getKey().belongsToInterval(entryTs)) { aggIntervalEntry.getValue().setLastArgsRefreshTs(System.currentTimeMillis()); - aggIntervals.put(aggIntervalEntry.getKey(), aggIntervalEntry.getValue()); return true; } } 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 5510d63f24..0f64140cab 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 @@ -27,6 +27,7 @@ import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInp import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.AggInterval; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval.Watermark; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldResult; @@ -40,8 +41,9 @@ import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; -import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; +import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultMetricArgumentEntry; public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { @@ -104,8 +106,9 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt this.cfProcessingService = ctx.getCfProcessingService(); var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); intervalDuration = configuration.getInterval().getIntervalDurationMillis(); - watermarkDuration = configuration.getWatermark().getDuration(); - checkInterval = configuration.getWatermark().getCheckInterval(); + Watermark watermark = configuration.getWatermark(); + watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); + checkInterval = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getCheckInterval()); interval = configuration.getInterval(); metrics = configuration.getMetrics(); } @@ -221,7 +224,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt AggMetric metric = metrics.get(metricName); String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); ArgumentEntry metricEntry = useDefault - ? ArgumentEntry.createSingleValueArgument(createDefaultKvEntry(argKey, metric.getDefaultValue())) + ? createDefaultMetricArgumentEntry(argKey, metric) : cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); if (!metricEntry.isEmpty()) { results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry); 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 676903e239..38a204d882 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -70,11 +70,19 @@ public class CalculatedFieldArgumentUtils { public static ArgumentEntry transformAggMetricArgument(List timeSeries, String argKey, AggMetric aggMetric) { if (timeSeries == null || timeSeries.isEmpty()) { - return ArgumentEntry.createSingleValueArgument(createDefaultKvEntry(argKey, aggMetric.getDefaultValue())); + return createDefaultMetricArgumentEntry(argKey, aggMetric); } return ArgumentEntry.createSingleValueArgument(timeSeries.get(0)); } + 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 StringDataEntry(argKey, null)); + } + public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs) { Map aggIntervals = new HashMap<>(); AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(startIntervalTs, endIntervalTs); @@ -87,10 +95,8 @@ public class CalculatedFieldArgumentUtils { } public static KvEntry createDefaultKvEntry(Argument argument) { - return createDefaultKvEntry(argument.getRefEntityKey().getKey(), argument.getDefaultValue()); - } - - public static KvEntry createDefaultKvEntry(String key, String defaultValue) { + String key = argument.getRefEntityKey().getKey(); + String defaultValue = argument.getDefaultValue(); if (StringUtils.isBlank(defaultValue)) { return new StringDataEntry(key, null); } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java index 8bf2818bd2..355ca2c72d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/AggMetric.java @@ -27,6 +27,6 @@ public class AggMetric { private AggFunction function; private String filter; private AggInput input; - private String defaultValue; + private Long defaultValue; }