diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index 94b1a8b878..9091b603fd 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -24,8 +24,9 @@ SET profile_data = jsonb_set( 'minAllowedScheduledUpdateIntervalInSecForCF', 60, 'maxRelationLevelPerCfArgument', 10, 'maxRelatedEntitiesToReturnPerCfArgument', 100, - 'minAllowedDeduplicationIntervalInSecForCF', 60, - 'minAllowedAggregationIntervalInSecForCF', 60 + 'minAllowedDeduplicationIntervalInSecForCF', 10, + 'minAllowedAggregationIntervalInSecForCF', 60, + 'minAllowedRealtimeAggregationIntervalInSecForCF', 300 ) || jsonb_strip_nulls(profile_data -> 'configuration') @@ -36,7 +37,8 @@ WHERE NOT ( 'maxRelationLevelPerCfArgument', 'maxRelatedEntitiesToReturnPerCfArgument', 'minAllowedDeduplicationIntervalInSecForCF', - 'minAllowedAggregationIntervalInSecForCF' + 'minAllowedAggregationIntervalInSecForCF', + 'minAllowedRealtimeAggregationIntervalInSecForCF' ] ); 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 9c04cb92bd..2c0ce1c938 100644 --- a/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java +++ b/application/src/main/java/org/thingsboard/server/controller/SystemInfoController.java @@ -166,6 +166,7 @@ public class SystemInfoController extends BaseController { systemParams.setMaxRelationLevelPerCfArgument(tenantProfileConfiguration.getMaxRelationLevelPerCfArgument()); systemParams.setMinAllowedDeduplicationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedDeduplicationIntervalInSecForCF()); systemParams.setMinAllowedAggregationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedAggregationIntervalInSecForCF()); + systemParams.setMinAllowedRealtimeAggregationIntervalInSecForCF(tenantProfileConfiguration.getMinAllowedRealtimeAggregationIntervalInSecForCF()); 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/controller/TenantProfileController.java b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java index 19cc2341ad..9c3985d3b0 100644 --- a/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java +++ b/application/src/main/java/org/thingsboard/server/controller/TenantProfileController.java @@ -166,9 +166,10 @@ public class TenantProfileController extends BaseController { " \"maxRelatedEntitiesToReturnPerCfArgument\": 100,\n" + " \"maxDataPointsPerRollingArg\": 1000,\n" + " \"maxStateSizeInKBytes\": 32,\n" + - " \"maxSingleValueArgumentSizeInKBytes\": 2" + - " \"minAllowedDeduplicationIntervalInSecForCF\": 60" + - " \"minAllowedAggregationIntervalInSecForCF\": 60" + + " \"maxSingleValueArgumentSizeInKBytes\": 2," + + " \"minAllowedDeduplicationIntervalInSecForCF\": 10," + + " \"minAllowedAggregationIntervalInSecForCF\": 60," + + " \"minAllowedRealtimeAggregationIntervalInSecForCF\": 300" + " }\n" + " },\n" + " \"default\": false\n" + 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 f04f3b109a..ec2d11ebf9 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 @@ -115,6 +115,7 @@ public class CalculatedFieldCtx implements Closeable { private long maxStateSize; private long maxSingleValueArgumentSize; + private long realtimeAggregationIntervalMillis; private boolean relationQueryDynamicArguments; private List mainEntityGeofencingArgumentNames; @@ -210,6 +211,7 @@ public class CalculatedFieldCtx implements Closeable { this.maxStateSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; this.maxSingleValueArgumentSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; + this.realtimeAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMinAllowedRealtimeAggregationIntervalInSecForCF)); } public boolean requiresScheduledReevaluation() { @@ -223,6 +225,12 @@ public class CalculatedFieldCtx implements Closeable { lastReevaluationTs = now; return true; } + if (entityAggregationConfig.isProduceIntermediateResult()) { + if (now - lastReevaluationTs >= realtimeAggregationIntervalMillis) { + lastReevaluationTs = now; + return true; + } + } ZonedDateTime lastReevaluationTime = TimeUtils.toZonedDateTime(lastReevaluationTs, entityAggregationConfig.getInterval().getZoneId()); long previousIntervalEndTs = entityAggregationConfig.getInterval().getDateTimeIntervalEndTs(lastReevaluationTime); if (now >= previousIntervalEndTs) { @@ -291,6 +299,7 @@ public class CalculatedFieldCtx implements Closeable { public void updateTenantProfileProperties() { this.maxStateSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxStateSizeInKBytes) * 1024; this.maxSingleValueArgumentSize = systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMaxSingleValueArgumentSizeInKBytes) * 1024; + this.realtimeAggregationIntervalMillis = TimeUnit.SECONDS.toMillis(systemContext.getApiLimitService().getLimit(tenantId, DefaultTenantProfileConfiguration::getMinAllowedRealtimeAggregationIntervalInSecForCF)); } public double evaluateSimpleExpression(Expression expression, CalculatedFieldState state) { 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 f3c3e8a1cc..a722f688dc 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 @@ -63,6 +63,8 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt private long checkInterval; private Map metrics; + private boolean produceIntermediateResult; + private EntityAggregationDebugArgumentsTracker debugTracker; private CalculatedFieldProcessingService cfProcessingService; @@ -81,6 +83,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt checkInterval = TimeUnit.SECONDS.toMillis(ctx.getSystemContext().getCfCheckInterval()); interval = configuration.getInterval(); metrics = configuration.getMetrics(); + produceIntermediateResult = configuration.isProduceIntermediateResult(); } @Override @@ -113,7 +116,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt Map> results = new HashMap<>(); List expiredIntervals = new ArrayList<>(); getIntervals().forEach((intervalEntry, argIntervalStatuses) -> { - processInterval(now, intervalEntry, argIntervalStatuses, expiredIntervals, results); + processInterval(now, ctx, intervalEntry, argIntervalStatuses, expiredIntervals, results); }); removeExpiredIntervals(expiredIntervals); @@ -193,6 +196,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt } private void processInterval(long now, + CalculatedFieldCtx ctx, AggIntervalEntry intervalEntry, Map args, List expiredIntervals, @@ -208,6 +212,10 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt if (watermarkDuration == 0) { expiredIntervals.add(intervalEntry); } + } else if (now - startTs < intervalEntry.getIntervalDuration()) { + if (produceIntermediateResult) { + handleCurrentInterval(ctx, intervalEntry, args, results); + } } } @@ -242,6 +250,25 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt }); } + private void handleCurrentInterval(CalculatedFieldCtx ctx, + AggIntervalEntry intervalEntry, + Map args, + Map> results) { + long realtimeAggregationInterval = ctx.getRealtimeAggregationIntervalMillis(); + args.forEach((argName, argEntryIntervalStatus) -> { + if (argEntryIntervalStatus.intervalPassed(realtimeAggregationInterval)) { + if (argEntryIntervalStatus.argsUpdated()) { + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); + argEntryIntervalStatus.setLastArgsRefreshTs(-1); + processArgument(intervalEntry, argName, false, results); + } else if (argEntryIntervalStatus.getLastMetricsEvalTs() == -1) {// TODO: should we return default value when the interval has not ended + argEntryIntervalStatus.setLastMetricsEvalTs(System.currentTimeMillis()); + processArgument(intervalEntry, argName, true, results); + } + } + }); + } + private void processArgument(AggIntervalEntry intervalEntry, String argName, boolean useDefault, diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index d3d04bcea0..8dab52f100 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -541,7 +541,7 @@ actors: # Interval in seconds to check calculated fields for re-evaluation interval. 1 minute by default. check_interval: "${ACTORS_CALCULATED_FIELDS_CHECK_INTERVAL_SEC:60}" alarms: - # Interval in seconds to re-evaluate Alarm rules that have a time schedule. 2 minutes by default. + # Interval in seconds to re-evaluate Alarm rules that have a time schedule. 1 minute by default. reevaluation_interval: "${ACTORS_ALARMS_REEVALUATION_INTERVAL_SEC:60}" debug: 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 0fa9b2dd78..d1b8bb0562 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 @@ -42,5 +42,6 @@ public class SystemParams { int maxRelationLevelPerCfArgument; long minAllowedDeduplicationIntervalInSecForCF; long minAllowedAggregationIntervalInSecForCF; + long minAllowedRealtimeAggregationIntervalInSecForCF; 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 f6095d41a7..488db86870 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 @@ -43,6 +43,7 @@ public class EntityAggregationCalculatedFieldConfiguration implements ArgumentsB private AggInterval interval; @Valid private Watermark watermark; + private boolean produceIntermediateResult; @Valid @NotNull private Output output; 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 87fa4a85da..3825c3c264 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 @@ -190,10 +190,12 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxStateSizeInKBytes = 32; @Schema(example = "2") private long maxSingleValueArgumentSizeInKBytes = 2; - @Schema(example = "60") - private long minAllowedDeduplicationIntervalInSecForCF = 60; + @Schema(example = "10") + private long minAllowedDeduplicationIntervalInSecForCF = 10; @Schema(example = "60") private long minAllowedAggregationIntervalInSecForCF = 60; + @Schema(example = "300") + private long minAllowedRealtimeAggregationIntervalInSecForCF = 300; @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 c10da4e6c6..6c333f36ad 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 @@ -50,7 +50,7 @@ public class CalculatedFieldDataValidator extends DataValidator validateCalculatedFieldConfiguration(calculatedField); validateSchedulingConfiguration(tenantId, calculatedField); validateRelationQuerySourceArguments(tenantId, calculatedField); - validateAggregationConfiguration(tenantId, calculatedField); + validateRelatedAggregationConfiguration(tenantId, calculatedField); validateEntityAggregationConfiguration(tenantId, calculatedField); } @@ -119,7 +119,7 @@ public class CalculatedFieldDataValidator extends DataValidator wrapAsDataValidation(() -> relationQueryDynamicSourceConfiguration.validateMaxRelationLevel(argumentName, maxRelationLevel))); } - private void validateAggregationConfiguration(TenantId tenantId, CalculatedField calculatedField) { + private void validateRelatedAggregationConfiguration(TenantId tenantId, CalculatedField calculatedField) { if (!(calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfiguration)) { return; } diff --git a/ui-ngx/src/app/shared/models/tenant.model.ts b/ui-ngx/src/app/shared/models/tenant.model.ts index 0cfa8df888..b8f04250ce 100644 --- a/ui-ngx/src/app/shared/models/tenant.model.ts +++ b/ui-ngx/src/app/shared/models/tenant.model.ts @@ -176,7 +176,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan maxArgumentsPerCF: 10, maxDataPointsPerRollingArg: 1000, maxRelationLevelPerCfArgument: 10, - minAllowedDeduplicationIntervalInSecForCF: 60, + minAllowedDeduplicationIntervalInSecForCF: 10, minAllowedAggregationIntervalInSecForCF: 60, maxRelatedEntitiesToReturnPerCfArgument: 100, minAllowedScheduledUpdateIntervalInSecForCF: 0,