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 520882fa75..6ead645322 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 @@ -45,7 +45,9 @@ import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultMetricArgumentEntry; @@ -190,10 +192,10 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt args.forEach((argName, argEntryIntervalStatus) -> { if (argEntryIntervalStatus.getLastArgsRefreshTs() > argEntryIntervalStatus.getLastMetricsEvalTs()) { argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); - processMetric(intervalEntry, argName, false, results); + processArgument(intervalEntry, argName, false, results); } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); - processMetric(intervalEntry, argName, true, results); + processArgument(intervalEntry, argName, true, results); } }); } @@ -206,38 +208,39 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt if (argEntryIntervalStatus.argsUpdated()) { argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); argEntryIntervalStatus.setLastArgsRefreshTs(-1); - processMetric(intervalEntry, argName, false, results); + processArgument(intervalEntry, argName, false, results); } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) { argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); - processMetric(intervalEntry, argName, true, results); + processArgument(intervalEntry, argName, true, results); } } }); } - private void processMetric(AggIntervalEntry intervalEntry, - String argName, - boolean useDefault, - Map> results) { - String metricName = findMetricName(argName); - if (metricName != null) { - AggMetric metric = metrics.get(metricName); - String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); - ArgumentEntry metricEntry = useDefault - ? createDefaultMetricArgumentEntry(argKey, metric) - : cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); - if (!metricEntry.isEmpty()) { - results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry); - } + private void processArgument(AggIntervalEntry intervalEntry, + String argName, + boolean useDefault, + Map> results) { + Set metrics = findMetrics(argName); + if (!metrics.isEmpty()) { + metrics.forEach(metricName -> { + AggMetric metric = this.metrics.get(metricName); + String argKey = ctx.getArguments().get(argName).getRefEntityKey().getKey(); + ArgumentEntry metricEntry = useDefault + ? createDefaultMetricArgumentEntry(argKey, metric) + : cfProcessingService.fetchMetricDuringInterval(ctx.getTenantId(), entityId, argKey, metric, intervalEntry); + if (!metricEntry.isEmpty()) { + results.computeIfAbsent(intervalEntry, i -> new HashMap<>()).put(metricName, metricEntry); + } + }); } } - private String findMetricName(String argName) { + private Set findMetrics(String argName) { return metrics.entrySet().stream() .filter(e -> ((AggKeyInput) e.getValue().getInput()).getKey().equals(argName)) .map(Map.Entry::getKey) - .findFirst() - .orElse(null); + .collect(Collectors.toSet()); } protected ArrayNode toResult(Map> results, Integer precision) { diff --git a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java index cf2c558618..3044525757 100644 --- a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -99,16 +99,17 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); long intervalEndTs = customInterval.getCurrentIntervalEndTs(); - CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, null); + CalculatedField consumptionCF = createConsumptionCF(device.getId(), customInterval, null); long interval = customInterval.getCurrentIntervalDurationMillis(); await().alias("create CF and no telemetry during interval -> save metric with default value") .atMost(2 * interval, TimeUnit.MILLISECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); assertThat(result).isNotNull(); assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("9999"); + assertThat(result.get("avgConsumption").get(0).get("value").isNull()).isTrue(); }); } @@ -130,15 +131,16 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); long interval = customInterval.getCurrentIntervalDurationMillis(); - CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, null); + CalculatedField consumptionCF = createConsumptionCF(device.getId(), customInterval, null); await().alias("create CF -> perform aggregation after interval end") .atMost(2 * interval, TimeUnit.MILLISECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); assertThat(result).isNotNull(); assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("133"); }); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":500}}", tsInInterval_1)); @@ -147,9 +149,10 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest .atMost(2 * interval, TimeUnit.MILLISECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); assertThat(result).isNotNull(); assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("133"); }); } @@ -172,15 +175,16 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest long interval = customInterval.getCurrentIntervalDurationMillis(); Watermark watermark = new Watermark(10); - CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); + CalculatedField consumptionCF = createConsumptionCF(device.getId(), customInterval, watermark); await().alias("create CF -> perform aggregation after interval end") .atMost(2 * interval, TimeUnit.MILLISECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); assertThat(result).isNotNull(); assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("133"); }); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":300}}", tsInInterval_1)); @@ -189,13 +193,14 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest .atMost(2 * 10, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { - ObjectNode result = getLatestTelemetry(device.getId(), "consumption"); + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgConsumption"); assertThat(result).isNotNull(); assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("600"); + assertThat(result.get("avgConsumption").get(0).get("value").asText()).isEqualTo("200"); }); } - private CalculatedField createTotalConsumptionCF(EntityId entityId, AggInterval aggInterval, Watermark watermark) { + private CalculatedField createConsumptionCF(EntityId entityId, AggInterval aggInterval, Watermark watermark) { Map arguments = new HashMap<>(); Argument argument = new Argument(); argument.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null)); @@ -209,10 +214,86 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest consumption.setDefaultValue(9999L); aggMetrics.put("consumption", consumption); + AggMetric avgEnergyConsumption = new AggMetric(); + avgEnergyConsumption.setFunction(AggFunction.AVG); + avgEnergyConsumption.setInput(new AggKeyInput("en")); + aggMetrics.put("avgConsumption", avgEnergyConsumption); + + TimeSeriesOutput output = new TimeSeriesOutput(); + output.setDecimalsByDefault(0); + + return createAggCf("Consumption", entityId, + aggInterval, + watermark, + arguments, + aggMetrics, + output); + } + + @Test + public void testCreateCfWith2Arguments_checkAggregation() throws Exception { + Device device = createDevice("Device", "1234567890111"); + + CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 0L, 5L); + long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); + long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); + + long tsBeforeInterval = currentIntervalStartTs - 1000; + long tsInInterval_1 = currentIntervalStartTs + 1000; + long tsInInterval_2 = currentIntervalStartTs + 500; + long tsInInterval_3 = currentIntervalStartTs + 200; + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsBeforeInterval)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":100}}", tsInInterval_1)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); + + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"temperature\":43}}", tsBeforeInterval)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"temperature\":39}}", tsInInterval_1)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"temperature\":27}}", tsInInterval_2)); + postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"temperature\":50}}", tsInInterval_3)); + + long interval = customInterval.getCurrentIntervalDurationMillis(); + + CalculatedField consumptionCF = createCFWith2Args(device.getId(), customInterval, null); + + await().alias("create CF -> perform aggregation after interval end") + .atMost(2 * interval, TimeUnit.MILLISECONDS) + .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) + .untilAsserted(() -> { + ObjectNode result = getLatestTelemetry(device.getId(), "consumption", "avgTemperature"); + assertThat(result).isNotNull(); + assertThat(result.get("consumption").get(0).get("value").asText()).isEqualTo("400"); + assertThat(result.get("avgTemperature").get(0).get("value").asText()).isEqualTo("39"); + }); + } + + private CalculatedField createCFWith2Args(EntityId entityId, AggInterval aggInterval, Watermark watermark) { + Map arguments = new HashMap<>(); + Argument energy = new Argument(); + energy.setRefEntityKey(new ReferencedEntityKey("energy", ArgumentType.TS_LATEST, null)); + arguments.put("en", energy); + + Argument temperature = new Argument(); + temperature.setRefEntityKey(new ReferencedEntityKey("temperature", ArgumentType.TS_LATEST, null)); + arguments.put("temp", temperature); + + Map aggMetrics = new HashMap<>(); + + AggMetric consumption = new AggMetric(); + consumption.setFunction(AggFunction.SUM); + consumption.setInput(new AggKeyInput("en")); + consumption.setDefaultValue(9999L); + aggMetrics.put("consumption", consumption); + + AggMetric avgTemperature = new AggMetric(); + avgTemperature.setFunction(AggFunction.AVG); + avgTemperature.setInput(new AggKeyInput("temp")); + aggMetrics.put("avgTemperature", avgTemperature); + TimeSeriesOutput output = new TimeSeriesOutput(); output.setDecimalsByDefault(0); - return createAggCf("Consumption per minute", entityId, + return createAggCf("CF with 2 args", entityId, aggInterval, watermark, arguments,