From f1ce07ef3f9522a5e9e36c33425c1c4d2e3896a5 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 31 Oct 2025 13:54:28 +0200 Subject: [PATCH] refactoring --- ...tractCalculatedFieldProcessingService.java | 34 +++++------------- .../cf/CalculatedFieldProcessingService.java | 3 +- ...faultCalculatedFieldProcessingService.java | 5 +-- ...EntityAggregationCalculatedFieldState.java | 36 +++++++++---------- .../utils/CalculatedFieldArgumentUtils.java | 11 ++++-- 5 files changed, 38 insertions(+), 51 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 4e2e81e9be..c58444397c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunction; -import org.thingsboard.server.common.data.cf.configuration.aggregation.AggKeyInput; import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; @@ -54,6 +53,7 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.SingleValueArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.aggregation.single.AggIntervalEntry; +import org.thingsboard.server.utils.CalculatedFieldArgumentUtils; import java.util.Collections; import java.util.HashMap; @@ -145,14 +145,14 @@ public abstract class AbstractCalculatedFieldProcessingService { )); } - protected ArgumentEntry resolveArgumentValue(String argName, ListenableFuture future) { + protected ArgumentEntry resolveArgumentValue(String key, ListenableFuture future) { try { return future.get(); } catch (ExecutionException e) { Throwable cause = e.getCause(); - throw new RuntimeException("Failed to fetch " + argName + ": " + cause.getMessage(), cause); + throw new RuntimeException("Failed to fetch " + key + ": " + cause.getMessage(), cause); } catch (InterruptedException e) { - throw new RuntimeException("Failed to fetch" + argName, e); + throw new RuntimeException("Failed to fetch" + key, e); } } @@ -298,30 +298,12 @@ public abstract class AbstractCalculatedFieldProcessingService { }; } - protected ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String metricName, CalculatedFieldCtx ctx) throws Exception { - var config = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); - - AggMetric metric = config.getMetrics().get(metricName); + protected ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval) { AggFunction function = metric.getFunction(); - - AggKeyInput input = (AggKeyInput) metric.getInput(); - String argName = input.getKey(); - Argument argument = ctx.getArguments().get(argName); - String key = argument.getRefEntityKey().getKey(); - long intervalMs = interval.getEndTs() - interval.getStartTs(); - BaseReadTsKvQuery query = new BaseReadTsKvQuery(key, interval.getStartTs(), interval.getEndTs(), intervalMs, 1, Aggregation.valueOf(function.name())); - log.trace("[{}][{}] Fetching timeseries for query {}", ctx.getTenantId(), entityId, query); - ListenableFuture> tsFuture = timeseriesService.findAll(ctx.getTenantId(), entityId, List.of(query)); - ListenableFuture argumentEntryFut = Futures.transform(tsFuture, timeSeries -> { - log.debug("[{}][{}] Fetched {} timeseries for query {}", ctx.getTenantId(), entityId, timeSeries == null ? 0 : timeSeries.size(), query); - if (timeSeries == null || timeSeries.isEmpty()) { - return new SingleValueArgumentEntry(); - } - return ArgumentEntry.createSingleValueArgument(timeSeries.get(0)); - }, calculatedFieldCallbackExecutor); - - return resolveArgumentValue(argName, argumentEntryFut); + BaseReadTsKvQuery query = new BaseReadTsKvQuery(argKey, interval.getStartTs(), interval.getEndTs(), intervalMs, 1, Aggregation.valueOf(function.name())); + ListenableFuture argumentEntryFut = fetchTimeSeriesInternal(tenantId, entityId, query, CalculatedFieldArgumentUtils::transformAggMetricArgument); + return resolveArgumentValue(argKey, argumentEntryFut); } private ListenableFuture fetchTimeSeries(TenantId tenantId, EntityId entityId, Argument argument, AggInterval interval, long queryEndTs) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java index e65df09f98..67bd2fec0c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.cf; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -38,7 +39,7 @@ public interface CalculatedFieldProcessingService { Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); - ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String argName, CalculatedFieldCtx ctx) throws Exception; + ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval); void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index 10ace6b8cf..c0114780b2 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -24,6 +24,7 @@ import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.configuration.Argument; +import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -113,8 +114,8 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF } @Override - public ArgumentEntry fetchMetricDuringInterval(EntityId entityId, AggIntervalEntry interval, String metricName, CalculatedFieldCtx ctx) throws Exception { - return super.fetchMetricDuringInterval(entityId, interval, metricName, ctx); + public ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval) { + return super.fetchMetricDuringInterval(tenantId, entityId, argKey, metric, interval); } @Override 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 93d83476cc..a5daf36b58 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 @@ -120,11 +120,9 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results = new HashMap<>(); List expiredIntervals = new ArrayList<>(); - for (Map.Entry> entry : intervals.entrySet()) { - AggIntervalEntry intervalEntry = entry.getKey(); - Map args = entry.getValue(); - processInterval(now, intervalEntry, args, expiredIntervals, results); - } + intervals.forEach((intervalEntry, argIntervalStatuses) -> { + processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); + }); expiredIntervals.forEach(intervals::remove); ArrayNode result = toResult(results); @@ -165,7 +163,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt AggIntervalEntry intervalEntry, Map args, List expiredIntervals, - Map> results) throws Exception { + Map> results) { long startTs = intervalEntry.getStartTs(); long endTs = intervalEntry.getEndTs(); @@ -179,37 +177,35 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt private void handleExpiredInterval(AggIntervalEntry intervalEntry, Map args, - Map> results) throws Exception { - for (Map.Entry argStatus : args.entrySet()) { - String argName = argStatus.getKey(); - AggIntervalEntryStatus argEntryIntervalStatus = argStatus.getValue(); + Map> results) { + args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { processMetric(intervalEntry, argName, results); } - } + }); } private void handleActiveInterval(AggIntervalEntry intervalEntry, Map args, - Map> results) throws Exception { - for (Map.Entry argStatus : args.entrySet()) { - String argName = argStatus.getKey(); - AggIntervalEntryStatus argEntryIntervalStatus = argStatus.getValue(); + Map> results) { + args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.shouldRecalculate(checkInterval)) { processMetric(intervalEntry, argName, results); ctx.scheduleReevaluation(checkInterval, actorCtx); } - } + }); } private void processMetric(AggIntervalEntry intervalEntry, String argName, - Map> results) throws Exception { + Map> results) { String metricName = findMetricName(argName); if (metricName != null) { - ArgumentEntry metric = cfProcessingService.fetchMetricDuringInterval(entityId, intervalEntry, metricName, ctx); - if (!metric.isEmpty()) { - results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metric); + AggMetric metric = metrics.get(metricName); + String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); + ArgumentEntry metricEntry = 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 0477669663..c6c64782a9 100644 --- a/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java +++ b/application/src/main/java/org/thingsboard/server/utils/CalculatedFieldArgumentUtils.java @@ -67,10 +67,17 @@ public class CalculatedFieldArgumentUtils { return ArgumentEntry.createTsRollingArgument(tsRolling, limit, argTimeWindow); } - public static ArgumentEntry transformAggregationArgument(List telemetry, long startIntervalTs, long endIntervalTs) { + public static ArgumentEntry transformAggMetricArgument(List timeSeries) { + if (timeSeries == null || timeSeries.isEmpty()) { + return new SingleValueArgumentEntry(); + } + return ArgumentEntry.createSingleValueArgument(timeSeries.get(0)); + } + + public static ArgumentEntry transformAggregationArgument(List timeSeries, long startIntervalTs, long endIntervalTs) { Map aggIntervals = new HashMap<>(); AggIntervalEntry aggIntervalEntry = new AggIntervalEntry(startIntervalTs, endIntervalTs); - if (telemetry == null || telemetry.isEmpty()) { + if (timeSeries == null || timeSeries.isEmpty()) { aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus()); } else { aggIntervals.put(aggIntervalEntry, new AggIntervalEntryStatus(System.currentTimeMillis()));