Browse Source

Merge pull request #14405 from irynamatveieva/fix/entity-agg-cf

Fixed incorrect skipping of metric calculation when using the same argument
pull/14413/head
Viacheslav Klimov 10 months ago
committed by GitHub
parent
commit
41ef18e6bf
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 45
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  2. 101
      application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java

45
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<AggIntervalEntry, Map<String, ArgumentEntry>> 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<AggIntervalEntry, Map<String, ArgumentEntry>> results) {
Set<String> 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<String> 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<AggIntervalEntry, Map<String, ArgumentEntry>> results, Integer precision) {

101
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<String, Argument> 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<String, Argument> 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<String, AggMetric> 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,

Loading…
Cancel
Save