diff --git a/application/src/main/data/upgrade/basic/schema_update.sql b/application/src/main/data/upgrade/basic/schema_update.sql index fe79fce3a2..e5fdba571b 100644 --- a/application/src/main/data/upgrade/basic/schema_update.sql +++ b/application/src/main/data/upgrade/basic/schema_update.sql @@ -45,7 +45,7 @@ SET profile_data = jsonb_set( CASE WHEN (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF' THEN NULL - ELSE to_jsonb(3600) + ELSE to_jsonb(60) END ) ), diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java index 4966b8132e..673db74863 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java @@ -411,6 +411,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM } catch (Exception e) { throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); } + } else if (ctx.shouldFetchEntityRelations(state)) { + log.debug("[{}][{}] Going to update related entities for CF.", entityId, ctx.getCfId()); + try { + if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesState) { + List relatedEntities = cfService.fetchRelatedEntities(ctx, entityId); + List missingEntities = relatedEntitiesState.checkRelatedEntities(relatedEntities); + if (!missingEntities.isEmpty()) { + missingEntities.forEach(missingEntityId -> { + Map fetchedArgs = cfService.fetchArgsFromDb(tenantId, missingEntityId, ctx.getArguments()); + relatedEntitiesState.updateEntityData(setEntityIdToSingleEntityArguments(missingEntityId, fetchedArgs)); + }); + justRestored = true; + } + } + } catch (Exception e) { + throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); + } } if (state.isSizeOk()) { Map updatedArgs = state.update(newArgValues, ctx); diff --git a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java index 4f0e323b99..e628d82b2f 100644 --- a/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java +++ b/application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java @@ -339,9 +339,11 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware List matchingCfs = cfsByEntityIdAndProfile.stream() .filter(cf -> { - var config = (RelatedEntitiesAggregationCalculatedFieldConfiguration) cf.getCalculatedField().getConfiguration(); - RelationPathLevel relation = config.getRelation(); - return direction.equals(relation.direction()) && relationType.equals(relation.relationType()); + if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) { + RelationPathLevel relation = config.getRelation(); + return direction.equals(relation.direction()) && relationType.equals(relation.relationType()); + } + return false; }) .toList(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 945792ebcd..5de44dc3f4 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -39,7 +39,6 @@ import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.relation.EntityRelation; import org.thingsboard.server.common.data.relation.EntityRelationPathQuery; -import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.relation.RelationService; @@ -172,26 +171,31 @@ public abstract class AbstractCalculatedFieldProcessingService { } protected Map> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { - RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); - - ListenableFuture> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, aggConfig.getRelation()); + if (!(ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config)) { + return Collections.emptyMap(); + } + ListenableFuture> relatedEntitiesFut = resolveRelatedEntities(entityId, ctx); - return aggConfig.getArguments().entrySet().stream() + return config.getArguments().entrySet().stream() .collect(Collectors.toMap( Map.Entry::getKey, entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchRelatedEntitiesArgumentEntry(ctx.getTenantId(), relatedEntities, entry.getValue(), ts), MoreExecutors.directExecutor()) )); } - private ListenableFuture> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) { - ListenableFuture> relationsFut = relationService.findByRelationPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation))); + protected ListenableFuture> resolveRelatedEntities(EntityId entityId, CalculatedFieldCtx ctx) { + if (!(ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config)) { + return Futures.immediateFuture(Collections.emptyList()); + } + + ListenableFuture> relationsFut = relationService.findByRelationPathQueryAsync(ctx.getTenantId(), new EntityRelationPathQuery(entityId, List.of(config.getRelation()))); return Futures.transform(relationsFut, relations -> { if (relations == null) { return Collections.emptyList(); } - return switch (relation.direction()) { + return switch (config.getRelation().direction()) { case FROM -> relations.stream() .map(EntityRelation::getTo) .toList(); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java index a9139572b8..15988c4b7b 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java @@ -35,6 +35,8 @@ public interface CalculatedFieldProcessingService { Map fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId); + List fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId); + Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index 52393d0ffe..28e5e4abe5 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -54,6 +54,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.concurrent.ExecutionException; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; @@ -97,6 +98,16 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF }; } + @Override + public List fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId) { + try { + return resolveRelatedEntities(entityId, ctx).get(); + } catch (ExecutionException | InterruptedException e) { + Throwable cause = e.getCause(); + throw new RuntimeException("Failed to fetch related entities for entity [" + entityId + "]: " + cause.getMessage(), cause); + } + } + @Override public Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments) { Map> argFutures = new HashMap<>(); 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 927787eae1..0df8d5ebe8 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 @@ -58,6 +58,7 @@ import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.state.aggregation.RelatedEntitiesAggregationCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculatedFieldState; import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; @@ -672,6 +673,19 @@ public class CalculatedFieldCtx implements Closeable { }; } + public boolean shouldFetchEntityRelations(CalculatedFieldState state) { + if (!(state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState)) { + return false; + } + if (!isScheduledUpdateEnabled()) { + return false; + } + if (relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() == -1L) { + return true; + } + return relatedEntitiesAggState.getLastRelatedEntitiesRefreshTs() < System.currentTimeMillis() - scheduledUpdateIntervalMillis; + } + @Override public void close() { try { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java index 89ca14c94b..23a4381ddc 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java @@ -38,7 +38,9 @@ import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry; +import java.util.ArrayList; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Map.Entry; @@ -52,6 +54,8 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat private long lastArgsRefreshTs = -1; @Setter private long lastMetricsEvalTs = -1; + @Setter + private long lastRelatedEntitiesRefreshTs = -1; private long deduplicationIntervalMs = -1; private Map metrics; @@ -76,9 +80,14 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat super.reset(); lastArgsRefreshTs = -1; lastMetricsEvalTs = -1; + lastRelatedEntitiesRefreshTs = -1; metrics = null; } + public void updateLastRelatedEntitiesRefreshTs() { + lastRelatedEntitiesRefreshTs = System.currentTimeMillis(); + } + @Override public CalculatedFieldType getType() { return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; @@ -90,6 +99,49 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat return super.update(argumentValues, ctx); } + public List checkRelatedEntities(List relatedEntities) { + Map> entityInputs = prepareInputs(); + findOutdatedEntities(entityInputs, relatedEntities).forEach(this::cleanupEntityData); + updateLastRelatedEntitiesRefreshTs(); + return findMissingEntities(entityInputs, relatedEntities); + } + + private List findMissingEntities(Map> entityInputs, List relatedEntities) { + List missing = new ArrayList<>(); + relatedEntities.forEach(entityId -> { + if (!entityInputs.containsKey(entityId)) { + missing.add(entityId); + log.warn("[{}] Missing related entity inputs for {}", ctx.getCfId(), entityId); + } + }); + return missing; + } + + private List findOutdatedEntities(Map> entityInputs, List relatedEntities) { + List outdated = new ArrayList<>(); + entityInputs.keySet().forEach(entityId -> { + if (!relatedEntities.contains(entityId)) { + outdated.add(entityId); + log.warn("[{}] CF state keeps outdated related entity {}", ctx.getCfId(), entityId); + } + }); + return outdated; + } + + public Map updateEntityData(Map fetchedArgs) { + lastMetricsEvalTs = -1; + return update(fetchedArgs, ctx); + } + + public void cleanupEntityData(EntityId relatedEntityId) { + arguments.values().forEach(argEntry -> { + RelatedEntitiesArgumentEntry aggEntry = (RelatedEntitiesArgumentEntry) argEntry; + aggEntry.getEntityInputs().remove(relatedEntityId); + }); + lastMetricsEvalTs = -1; + lastArgsRefreshTs = System.currentTimeMillis(); + } + @Override public ListenableFuture performCalculation(Map updatedArgs, CalculatedFieldCtx ctx) throws Exception { boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty(); @@ -108,20 +160,6 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat } } - public Map updateEntityData(Map fetchedArgs) { - lastMetricsEvalTs = -1; - return update(fetchedArgs, ctx); - } - - public void cleanupEntityData(EntityId relatedEntityId) { - arguments.values().forEach(argEntry -> { - RelatedEntitiesArgumentEntry aggEntry = (RelatedEntitiesArgumentEntry) argEntry; - aggEntry.getEntityInputs().remove(relatedEntityId); - }); - lastMetricsEvalTs = -1; - lastArgsRefreshTs = System.currentTimeMillis(); - } - private boolean shouldRecalculate() { boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationIntervalMs; boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java index 9d4c7bdaf6..b6eba53aaa 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java @@ -23,12 +23,13 @@ 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.ArgumentsBasedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.Output; +import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.relation.RelationPathLevel; import java.util.Map; @Data -public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration { +public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration, ScheduledUpdateSupportedCalculatedFieldConfiguration { @NotNull private RelationPathLevel relation; @@ -40,6 +41,9 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A private Output output; private boolean useLatestTs; + private boolean scheduledUpdateEnabled; + private int scheduledUpdateInterval; + @Override public CalculatedFieldType getType() { return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; 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..0e246f9268 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 @@ -186,8 +186,8 @@ public class DefaultTenantProfileConfiguration implements TenantProfileConfigura private long maxStateSizeInKBytes = 32; @Schema(example = "2") private long maxSingleValueArgumentSizeInKBytes = 2; - @Schema(example = "3600") - private long minAllowedDeduplicationIntervalInSecForCF = 3600; + @Schema(example = "60") + private long minAllowedDeduplicationIntervalInSecForCF = 60; @Override public long getProfileThreshold(ApiUsageRecordKey key) { diff --git a/ui-ngx/src/app/shared/models/tenant.model.ts b/ui-ngx/src/app/shared/models/tenant.model.ts index ae7a0ae8b8..34023f939b 100644 --- a/ui-ngx/src/app/shared/models/tenant.model.ts +++ b/ui-ngx/src/app/shared/models/tenant.model.ts @@ -175,7 +175,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan maxArgumentsPerCF: 10, maxDataPointsPerRollingArg: 1000, maxRelationLevelPerCfArgument: 10, - minAllowedDeduplicationIntervalInSecForCF: 3600, + minAllowedDeduplicationIntervalInSecForCF: 60, maxRelatedEntitiesToReturnPerCfArgument: 100, minAllowedScheduledUpdateIntervalInSecForCF: 0, maxStateSizeInKBytes: 32,