Browse Source

Merge pull request #14580 from irynamatveieva/entity-aggregation-cf/update-ts

Time series aggregation CF: updated result ts to align with existing TB aggregation logic
pull/14630/head
Viacheslav Klimov 10 months ago
committed by GitHub
parent
commit
cb622f33a7
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java

11
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -127,8 +127,8 @@ public class CalculatedFieldCtx implements Closeable {
private List<String> relatedEntityArgumentNames; private List<String> relatedEntityArgumentNames;
private long scheduledUpdateIntervalMillis; private long scheduledUpdateIntervalMillis;
private long cfCheckReevaluationInterval; private long cfCheckReevaluationIntervalMillis;
private long alarmReevaluationInterval; private long alarmReevaluationIntervalMillis;
private Argument propagationArgument; private Argument propagationArgument;
private boolean applyExpressionForResolvedArguments; private boolean applyExpressionForResolvedArguments;
@ -246,8 +246,7 @@ public class CalculatedFieldCtx implements Closeable {
boolean requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation(); boolean requiresScheduledReevaluation = calculatedField.getConfiguration().requiresScheduledReevaluation();
if (calculatedField.getConfiguration() instanceof AlarmCalculatedFieldConfiguration) { if (calculatedField.getConfiguration() instanceof AlarmCalculatedFieldConfiguration) {
if (requiresScheduledReevaluation) { if (requiresScheduledReevaluation) {
long reevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(alarmReevaluationInterval); if (now - lastReevaluationTs >= alarmReevaluationIntervalMillis) {
if (now - lastReevaluationTs >= reevaluationIntervalMillis) {
lastReevaluationTs = now; lastReevaluationTs = now;
return true; return true;
} }
@ -306,8 +305,8 @@ public class CalculatedFieldCtx implements Closeable {
this.maxStateSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; this.maxStateSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024;
this.maxSingleValueArgumentSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; this.maxSingleValueArgumentSize = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024;
this.intermediateAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getIntermediateAggregationIntervalInSecForCF)); this.intermediateAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getIntermediateAggregationIntervalInSecForCF));
this.cfCheckReevaluationInterval = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getCfReevaluationCheckInterval); this.cfCheckReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getCfReevaluationCheckInterval));
this.alarmReevaluationInterval = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getAlarmsReevaluationInterval); this.alarmReevaluationIntervalMillis = TimeUnit.SECONDS.toMillis(apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getAlarmsReevaluationInterval));
} }
public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) { public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) {

4
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java

@ -205,7 +205,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
handleExpiredInterval(intervalEntry, args, results); handleExpiredInterval(intervalEntry, args, results);
expiredIntervals.add(intervalEntry); expiredIntervals.add(intervalEntry);
} else if (now - startTs >= intervalEntry.getIntervalDuration()) { } else if (now - startTs >= intervalEntry.getIntervalDuration()) {
handleActiveInterval(ctx.getCfCheckReevaluationInterval(), intervalEntry, args, results); handleActiveInterval(ctx.getCfCheckReevaluationIntervalMillis(), intervalEntry, args, results);
if (watermarkDuration == 0) { if (watermarkDuration == 0) {
expiredIntervals.add(intervalEntry); expiredIntervals.add(intervalEntry);
} }
@ -288,7 +288,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
} }
if (!metricsNode.isEmpty()) { if (!metricsNode.isEmpty()) {
ObjectNode resultNode = JacksonUtil.newObjectNode(); ObjectNode resultNode = JacksonUtil.newObjectNode();
resultNode.put("ts", interval.getEndTs() - 1); resultNode.put("ts", interval.getStartTs());
resultNode.set("values", metricsNode); resultNode.set("values", metricsNode);
result.add(resultNode); result.add(resultNode);

Loading…
Cancel
Save