From fa4e056fb244477a0c2558386d6f14b6b06b00b5 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 5 Nov 2025 09:52:40 +0200 Subject: [PATCH] added min aggregation interval to tenant profile --- .../main/data/upgrade/basic/schema_update.sql | 8 ++++++ .../controller/SystemInfoController.java | 1 + .../cf/ctx/state/CalculatedFieldCtx.java | 27 ++++++++++--------- ...EntityAggregationCalculatedFieldState.java | 2 +- .../src/main/resources/thingsboard.yml | 4 +-- .../EntityAggregationCalculatedFieldTest.java | 6 ++--- .../server/common/data/SystemParams.java | 1 + ...gregationCalculatedFieldConfiguration.java | 5 ++-- .../single/interval/Watermark.java | 1 - .../DefaultTenantProfileConfiguration.java | 2 ++ .../CalculatedFieldDataValidator.java | 17 ++++++++++++ 11 files changed, 52 insertions(+), 22 deletions(-) diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index fe79fce3a2..3903bf1bf5 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -46,6 +46,12 @@ SET profile_data = jsonb_set( WHEN (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF' THEN NULL ELSE to_jsonb(3600) + END, + 'minAggregationIntervalInSecForCF', + CASE + WHEN (profile_data -> 'configuration') ? 'minAggregationIntervalInSecForCF' + THEN NULL + ELSE to_jsonb(60) END ) ), @@ -59,6 +65,8 @@ WHERE NOT ( (profile_data -> 'configuration') ? 'maxRelatedEntitiesToReturnPerCfArgument' AND (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF' + AND + (profile_data -> 'configuration') ? 'minAggregationIntervalInSecForCF' ); -- UPDATE TENANT PROFILE CONFIGURATION END diff --git a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java index 82807d0762..7b6d7b5768 100644 --- a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java +++ b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java @@ -165,6 +165,7 @@ public class SystemInfoController extends BaseController { systemParams.setMinAllowedScheduledUpdateIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedScheduledUpdateIntervalInSecForCF()); systemParams.setMaxRelationLevelPerCfArgument(tenantProfileConfiguration.getMaxRelationLevelPerCfArgument()); systemParams.setMinAllowedDeduplicationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedDeduplicationIntervalInSecForCF()); + systemParams.setMinAggregationIntervalInSecForCF(tenantProfileConfiguration.getMinAggregationIntervalInSecForCF()); systemParams.setTrendzSettings(trendzSettingsService.findTrendzSettings(currentUser.getTenantId())); } systemParams.setMobileQrEnabled(Optional.ofNullable(qrCodeSettingService.findQrCodeSettings(TenantId.SYS_TENANT_ID)) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java index 36c9436f9f..cdf055fe7a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java @@ -97,7 +97,7 @@ public class CalculatedFieldCtx implements Closeable { private boolean useLatestTs; private boolean requiresScheduledReevaluation; - private long lastReevaluationTs; + private long aggCheckInterval; private ActorSystemContext systemContext; private TbelInvokeService tbelInvokeService; @@ -199,6 +199,7 @@ public class CalculatedFieldCtx implements Closeable { if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig) { this.useLatestTs = aggConfig.isUseLatestTs(); } + this.aggCheckInterval = systemContext.getCfCheckInterval(); this.systemContext = systemContext; this.tbelInvokeService = systemContext.getTbelInvokeService(); this.relationService = systemContext.getRelationService(); @@ -213,20 +214,13 @@ public class CalculatedFieldCtx implements Closeable { public boolean isRequiresScheduledReevaluation() { if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration entityAggregationConfig) { long now = System.currentTimeMillis(); - long cfCheckIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getCfCheckInterval()); Watermark watermark = entityAggregationConfig.getWatermark(); - if (watermark != null) { - long checkIntervalMillis = TimeUnit.SECONDS.toMillis(watermark.getCheckInterval()); - if (cfCheckIntervalMillis == checkIntervalMillis) { - return true; - } - if (now + cfCheckIntervalMillis >= lastReevaluationTs + checkIntervalMillis) { - lastReevaluationTs = now; - return true; - } + if (watermark != null && watermark.getDuration() > 0) { + return true; } - long intervalEndTsMillis = entityAggregationConfig.getInterval().getCurrentIntervalEndTs(); - if (now + cfCheckIntervalMillis >= intervalEndTsMillis) { + long cfCheckIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getCfCheckInterval()); + long intervalEndTs = entityAggregationConfig.getInterval().getCurrentIntervalEndTs(); + if (now + cfCheckIntervalMillis >= intervalEndTs) { return true; } } @@ -635,6 +629,13 @@ public class CalculatedFieldCtx implements Closeable { && (thisConfig.getDeduplicationIntervalInSec() != otherConfig.getDeduplicationIntervalInSec() || !thisConfig.getMetrics().equals(otherConfig.getMetrics()))) { return true; } + if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig + && other.getCalculatedField().getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) { + boolean metricsChanged = thisConfig.getMetrics().equals(otherConfig.getMetrics()); + boolean intervalChanged = thisConfig.getInterval().equals(otherConfig.getInterval()); + boolean watermarkChanged = thisConfig.getWatermark().equals(otherConfig.getWatermark()); + return metricsChanged || intervalChanged || watermarkChanged; + } return false; } 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 ff93d4a924..4e7f073192 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 @@ -104,7 +104,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt intervalDuration = configuration.getInterval().getIntervalDurationMillis(); Watermark watermark = configuration.getWatermark(); watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); - checkInterval = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getCheckInterval()); + checkInterval = ctx.getAggCheckInterval(); interval = configuration.getInterval(); metrics = configuration.getMetrics(); } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index e49bcb9bf0..1c32d62f6c 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -529,8 +529,8 @@ actors: configuration: "${ACTORS_CALCULATED_FIELD_DEBUG_MODE_RATE_LIMITS_PER_TENANT_CONFIGURATION:50000:3600}" # Time in seconds to receive calculation result. calculation_timeout: "${ACTORS_CALCULATION_TIMEOUT_SEC:5}" - # Interval in seconds to re-evaluate calculated fields that have a time schedule. 2 minutes by default. - check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:30}" + # Interval in seconds to re-evaluate calculated fields that have a time schedule. 1 minute by default. + check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:60}" debug: settings: 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 f4123b31d6..fb92ebe3af 100644 --- a/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java +++ b/application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java @@ -109,7 +109,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); long interval = customInterval.getIntervalDurationMillis(); - Watermark watermark = new Watermark(60, 10); + Watermark watermark = new Watermark(60); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); await().alias("create CF and perform aggregation after interval end") @@ -140,7 +140,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); long interval = customInterval.getIntervalDurationMillis(); - Watermark watermark = new Watermark(60, 10); + Watermark watermark = new Watermark(60); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); await().alias("create CF and perform aggregation after interval end") @@ -155,7 +155,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":300}}", tsInInterval_1)); await().alias("create CF and perform aggregation after interval end") - .atMost(2 * watermark.getCheckInterval(), TimeUnit.SECONDS) + .atMost(2 * 10, TimeUnit.SECONDS) .pollInterval(POLL_INTERVAL, TimeUnit.SECONDS) .untilAsserted(() -> { ObjectNode result = getLatestTelemetry(device.getId(), "consumptionPerMin"); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java index 6a475daae3..c46b466e90 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/SystemParams.java @@ -41,5 +41,6 @@ public class SystemParams { int minAllowedScheduledUpdateIntervalInSecForCF; int maxRelationLevelPerCfArgument; long minAllowedDeduplicationIntervalInSecForCF; + long minAggregationIntervalInSecForCF; TrendzSettings trendzSettings; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/EntityAggregationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/EntityAggregationCalculatedFieldConfiguration.java index 15b85af7bc..78e74ee64d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/EntityAggregationCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/EntityAggregationCalculatedFieldConfiguration.java @@ -20,6 +20,7 @@ import jakarta.validation.constraints.NotEmpty; import lombok.Data; import org.thingsboard.server.common.data.cf.CalculatedFieldType; 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.ArgumentsBasedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.aggregation.AggMetric; @@ -49,8 +50,8 @@ public class EntityAggregationCalculatedFieldConfiguration implements ArgumentsB if (arguments.containsKey("ctx")) { throw new IllegalArgumentException("Argument name 'ctx' is reserved and cannot be used."); } - if (arguments.values().stream().anyMatch(Argument::hasTsRollingArgument)) { - throw new IllegalArgumentException("Calculated field with type: '" + getType() + "' doesn't support TS_ROLLING arguments."); + if (arguments.values().stream().anyMatch(argument -> !ArgumentType.TS_LATEST.equals(argument.getRefEntityKey().getType()))) { + throw new IllegalArgumentException("Calculated field with type: '" + getType() + "' support only TS_LATEST arguments."); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java index 20d03194ed..3b11681982 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/Watermark.java @@ -25,6 +25,5 @@ import lombok.NoArgsConstructor; public class Watermark { private long duration; - private long checkInterval; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java index 1cb1998a0e..9457b688f2 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java @@ -188,6 +188,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxSingleValueArgumentSizeInKBytes = 2; @Schema(example = "3600") private long minAllowedDeduplicationIntervalInSecForCF = 3600; + @Schema(example = "60") + private long minAggregationIntervalInSecForCF = 60; @Override public long getProfileThreshold(ApiUsageRecordKey key) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java index 3ccb837b3f..7813a244e1 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java +++ b/dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java @@ -22,6 +22,7 @@ import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalcula import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; +import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.dao.cf.CalculatedFieldDao; @@ -30,6 +31,7 @@ import org.thingsboard.server.dao.service.DataValidator; import org.thingsboard.server.dao.usagerecord.ApiLimitService; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; @Component @@ -48,6 +50,7 @@ public class CalculatedFieldDataValidator extends DataValidator validateSchedulingConfiguration(tenantId, calculatedField); validateRelationQuerySourceArguments(tenantId, calculatedField); validateAggregationConfiguration(tenantId, calculatedField); + validateEntityAggregationConfiguration(tenantId, calculatedField); } @Override @@ -123,6 +126,20 @@ public class CalculatedFieldDataValidator extends DataValidator } } + private void validateEntityAggregationConfiguration(TenantId tenantId, CalculatedField calculatedField) { + if (!(calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration aggConfiguration)) { + return; + } + long minAggregationIntervalInSec = apiLimitService.getLimit(tenantId, DefaultTenantProfileConfiguration::getMinAggregationIntervalInSecForCF); + if (minAggregationIntervalInSec <= 0) { + return; + } + if (aggConfiguration.getInterval().getIntervalDurationMillis() > TimeUnit.SECONDS.toMillis(minAggregationIntervalInSec)) { + throw new IllegalArgumentException("Aggregation interval duration is less than configured " + + "minimum allowed aggregation interval in tenant profile: " + minAggregationIntervalInSec + " sec."); + } + } + private static void wrapAsDataValidation(Runnable validation) { try { validation.run();