Browse Source

Merge pull request #14280 from irynamatveieva/related-entities-aggregation

Related Entities Aggregation Calculated Field
pull/14036/head
Viacheslav Klimov 11 months ago
committed by GitHub
parent
commit
3aef3a2c7f
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 2
      application/src/main/data/upgrade/basic/schema_update.sql
  2. 17
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  3. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  4. 11
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
  6. 15
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java
  8. 14
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  9. 87
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  10. 2
      application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java
  11. 10
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/RelatedEntitiesAggregationCalculatedFieldConfiguration.java
  12. 4
      common/data/src/main/java/org/thingsboard/server/common/data/tenant/profile/DefaultTenantProfileConfiguration.java
  13. 2
      ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts
  14. 1
      ui-ngx/src/app/shared/models/calculated-field.models.ts
  15. 2
      ui-ngx/src/app/shared/models/tenant.model.ts

2
application/src/main/data/upgrade/basic/schema_update.sql

@ -45,7 +45,7 @@ SET profile_data = jsonb_set(
CASE CASE
WHEN (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF' WHEN (profile_data -> 'configuration') ? 'minAllowedDeduplicationIntervalInSecForCF'
THEN NULL THEN NULL
ELSE to_jsonb(3600) ELSE to_jsonb(60)
END END
) )
), ),

17
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -411,6 +411,23 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} catch (Exception e) { } catch (Exception e) {
throw CalculatedFieldException.builder().ctx(ctx).eventEntity(entityId).cause(e).build(); 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<EntityId> relatedEntities = cfService.fetchRelatedEntities(ctx, entityId);
List<EntityId> missingEntities = relatedEntitiesState.checkRelatedEntities(relatedEntities);
if (!missingEntities.isEmpty()) {
missingEntities.forEach(missingEntityId -> {
Map<String, ArgumentEntry> 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()) { if (state.isSizeOk()) {
Map<String, ArgumentEntry> updatedArgs = state.update(newArgValues, ctx); Map<String, ArgumentEntry> updatedArgs = state.update(newArgValues, ctx);

3
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java

@ -349,9 +349,8 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) { if (cf.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) {
RelationPathLevel relation = config.getRelation(); RelationPathLevel relation = config.getRelation();
return direction.equals(relation.direction()) && relationType.equals(relation.relationType()); return direction.equals(relation.direction()) && relationType.equals(relation.relationType());
} else {
return false;
} }
return false;
}) })
.toList(); .toList();

11
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -174,18 +174,19 @@ public abstract class AbstractCalculatedFieldProcessingService {
} }
protected Map<String, ListenableFuture<ArgumentEntry>> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { protected Map<String, ListenableFuture<ArgumentEntry>> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) {
RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); if (!(ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config)) {
return Collections.emptyMap();
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, aggConfig.getRelation()); }
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, config.getRelation());
return aggConfig.getArguments().entrySet().stream() return config.getArguments().entrySet().stream()
.collect(Collectors.toMap( .collect(Collectors.toMap(
Map.Entry::getKey, Map.Entry::getKey,
entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchRelatedEntitiesArgumentEntry(ctx.getTenantId(), relatedEntities, entry.getValue(), ts), MoreExecutors.directExecutor()) entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchRelatedEntitiesArgumentEntry(ctx.getTenantId(), relatedEntities, entry.getValue(), ts), MoreExecutors.directExecutor())
)); ));
} }
private ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) { protected ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) {
Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo()); Predicate<EntityRelation> filter = entityRelation -> CalculatedField.isSupportedRefEntity(entityRelation.getFrom()) && CalculatedField.isSupportedRefEntity(entityRelation.getTo());
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findFilteredRelationsByPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation)), filter); ListenableFuture<List<EntityRelation>> relationsFut = relationService.findFilteredRelationsByPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation)), filter);

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java

@ -35,6 +35,8 @@ public interface CalculatedFieldProcessingService {
Map<String, ArgumentEntry> fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId); Map<String, ArgumentEntry> fetchDynamicArgsFromDb(CalculatedFieldCtx ctx, EntityId entityId);
List<EntityId> fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId);
Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments); Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments);
void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback); void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);

15
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -24,6 +24,7 @@ import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.DataConstants;
import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityType;
import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
@ -54,6 +55,7 @@ import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.UUID; 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.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT;
import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto; import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@ -97,6 +99,19 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
}; };
} }
@Override
public List<EntityId> fetchRelatedEntities(CalculatedFieldCtx ctx, EntityId entityId) {
try {
if (ctx.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration config) {
return resolveRelatedEntities(ctx.getTenantId(), entityId, config.getRelation()).get();
}
return Collections.emptyList();
} catch (ExecutionException | InterruptedException e) {
Throwable cause = e.getCause();
throw new RuntimeException("Failed to fetch related entities for entity [" + entityId + "]: " + cause.getMessage(), cause);
}
}
@Override @Override
public Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments) { public Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments) {
Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>(); Map<String, ListenableFuture<ArgumentEntry>> argFutures = new HashMap<>();

1
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldQueueService.java

@ -206,7 +206,6 @@ public class DefaultCalculatedFieldQueueService implements CalculatedFieldQueueS
return true; return true;
} }
} }
return false;
} }
} }
} }

14
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.dao.relation.RelationService;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; 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.cf.ctx.state.geofencing.GeofencingCalculatedFieldState;
import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; 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 @Override
public void close() { public void close() {
try { try {

87
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java

@ -38,9 +38,12 @@ 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.CalculatedFieldCtx;
import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry; import org.thingsboard.server.service.cf.ctx.state.aggregation.function.AggEntry;
import java.util.ArrayList;
import java.util.HashMap; import java.util.HashMap;
import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Map.Entry; import java.util.Map.Entry;
import java.util.concurrent.ScheduledFuture;
import static java.util.concurrent.TimeUnit.SECONDS; import static java.util.concurrent.TimeUnit.SECONDS;
@ -52,9 +55,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
private long lastArgsRefreshTs = -1; private long lastArgsRefreshTs = -1;
@Setter @Setter
private long lastMetricsEvalTs = -1; private long lastMetricsEvalTs = -1;
@Setter
private long lastRelatedEntitiesRefreshTs = -1;
private long deduplicationIntervalMs = -1; private long deduplicationIntervalMs = -1;
private Map<String, AggMetric> metrics; private Map<String, AggMetric> metrics;
private ScheduledFuture<?> reevaluationFuture;
public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) { public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) {
super(entityId); super(entityId);
} }
@ -67,8 +74,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec()); deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec());
} }
public void scheduleReevaluation() { @Override
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx); public void close() {
super.close();
if (reevaluationFuture != null) {
reevaluationFuture.cancel(true);
reevaluationFuture = null;
}
} }
@Override @Override
@ -76,9 +88,14 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
super.reset(); super.reset();
lastArgsRefreshTs = -1; lastArgsRefreshTs = -1;
lastMetricsEvalTs = -1; lastMetricsEvalTs = -1;
lastRelatedEntitiesRefreshTs = -1;
metrics = null; metrics = null;
} }
public void updateLastRelatedEntitiesRefreshTs() {
lastRelatedEntitiesRefreshTs = System.currentTimeMillis();
}
@Override @Override
public CalculatedFieldType getType() { public CalculatedFieldType getType() {
return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
@ -90,6 +107,56 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
return super.update(argumentValues, ctx); return super.update(argumentValues, ctx);
} }
public List<EntityId> checkRelatedEntities(List<EntityId> relatedEntities) {
Map<EntityId, Map<String, ArgumentEntry>> entityInputs = prepareInputs();
findOutdatedEntities(entityInputs, relatedEntities).forEach(this::cleanupEntityData);
updateLastRelatedEntitiesRefreshTs();
return findMissingEntities(entityInputs, relatedEntities);
}
private List<EntityId> findMissingEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) {
List<EntityId> 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<EntityId> findOutdatedEntities(Map<EntityId, Map<String, ArgumentEntry>> entityInputs, List<EntityId> relatedEntities) {
List<EntityId> 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<String, ArgumentEntry> updateEntityData(Map<String, ArgumentEntry> 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();
}
public void scheduleReevaluation() {
ScheduledFuture<?> future = ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
if (future != null) {
reevaluationFuture = future;
}
}
@Override @Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception { public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception {
boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty(); boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty();
@ -97,7 +164,7 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
Output output = ctx.getOutput(); Output output = ctx.getOutput();
ObjectNode aggResult = aggregateMetrics(output); ObjectNode aggResult = aggregateMetrics(output);
lastMetricsEvalTs = System.currentTimeMillis(); lastMetricsEvalTs = System.currentTimeMillis();
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx); scheduleReevaluation();
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())
@ -108,20 +175,6 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
} }
} }
public Map<String, ArgumentEntry> updateEntityData(Map<String, ArgumentEntry> 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() { private boolean shouldRecalculate() {
boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationIntervalMs; boolean intervalPassed = lastMetricsEvalTs <= System.currentTimeMillis() - deduplicationIntervalMs;
boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs; boolean argsUpdatedDuringInterval = lastArgsRefreshTs > lastMetricsEvalTs;

2
application/src/test/java/org/thingsboard/server/cf/RelatedEntitiesAggregationCalculatedFieldTest.java

@ -88,6 +88,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
updateDefaultTenantProfileConfig(tenantProfileConfig -> { updateDefaultTenantProfileConfig(tenantProfileConfig -> {
tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1); tenantProfileConfig.setMinAllowedDeduplicationIntervalInSecForCF(1);
tenantProfileConfig.setMinAllowedScheduledUpdateIntervalInSecForCF(1);
}); });
Tenant tenant = new Tenant(); Tenant tenant = new Tenant();
@ -801,6 +802,7 @@ public class RelatedEntitiesAggregationCalculatedFieldTest extends AbstractContr
configuration.setRelation(relation); configuration.setRelation(relation);
configuration.setArguments(inputs); configuration.setArguments(inputs);
configuration.setDeduplicationIntervalInSec(deduplicationInterval); configuration.setDeduplicationIntervalInSec(deduplicationInterval);
configuration.setScheduledUpdateInterval(10);
configuration.setMetrics(metrics); configuration.setMetrics(metrics);
configuration.setOutput(output); configuration.setOutput(output);

10
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.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration; 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.Output;
import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.relation.RelationPathLevel; import org.thingsboard.server.common.data.relation.RelationPathLevel;
import java.util.Map; import java.util.Map;
@Data @Data
public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration { public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements ArgumentsBasedCalculatedFieldConfiguration, ScheduledUpdateSupportedCalculatedFieldConfiguration {
@NotNull @NotNull
private RelationPathLevel relation; private RelationPathLevel relation;
@ -40,11 +41,18 @@ public class RelatedEntitiesAggregationCalculatedFieldConfiguration implements A
private Output output; private Output output;
private boolean useLatestTs; private boolean useLatestTs;
private int scheduledUpdateInterval;
@Override @Override
public CalculatedFieldType getType() { public CalculatedFieldType getType() {
return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; return CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
} }
@Override
public boolean isScheduledUpdateEnabled() {
return true;
}
@Override @Override
public void validate() { public void validate() {
relation.validate(); relation.validate();

4
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; private long maxStateSizeInKBytes = 32;
@Schema(example = "2") @Schema(example = "2")
private long maxSingleValueArgumentSizeInKBytes = 2; private long maxSingleValueArgumentSizeInKBytes = 2;
@Schema(example = "3600") @Schema(example = "60")
private long minAllowedDeduplicationIntervalInSecForCF = 3600; private long minAllowedDeduplicationIntervalInSecForCF = 60;
@Override @Override
public long getProfileThreshold(ApiUsageRecordKey key) { public long getProfileThreshold(ApiUsageRecordKey key) {

2
ui-ngx/src/app/modules/home/components/calculated-fields/components/related-entities-aggregation-configuration/related-entities-aggregation-component.component.ts

@ -105,6 +105,7 @@ export class RelatedEntitiesAggregationComponentComponent implements ControlValu
map(argumentsObj => getCalculatedFieldArgumentsHighlights(argumentsObj)) map(argumentsObj => getCalculatedFieldArgumentsHighlights(argumentsObj))
); );
private readonly minAllowedScheduledUpdateIntervalInSecForCF = getCurrentAuthState(this.store).minAllowedScheduledUpdateIntervalInSecForCF;
private propagateChange: (config: CalculatedFieldRelatedAggregationConfiguration) => void = () => { }; private propagateChange: (config: CalculatedFieldRelatedAggregationConfiguration) => void = () => { };
constructor(private fb: FormBuilder, constructor(private fb: FormBuilder,
@ -149,6 +150,7 @@ export class RelatedEntitiesAggregationComponentComponent implements ControlValu
private updatedModel(value: CalculatedFieldRelatedAggregationConfiguration): void { private updatedModel(value: CalculatedFieldRelatedAggregationConfiguration): void {
value.type = CalculatedFieldType.RELATED_ENTITIES_AGGREGATION; value.type = CalculatedFieldType.RELATED_ENTITIES_AGGREGATION;
value.scheduledUpdateInterval = this.minAllowedScheduledUpdateIntervalInSecForCF;
this.propagateChange(value); this.propagateChange(value);
} }
} }

1
ui-ngx/src/app/shared/models/calculated-field.models.ts

@ -142,6 +142,7 @@ export interface CalculatedFieldRelatedAggregationConfiguration {
arguments: Record<string, CalculatedFieldArgument>; arguments: Record<string, CalculatedFieldArgument>;
metrics: Record<string, CalculatedFieldAggMetric>; metrics: Record<string, CalculatedFieldAggMetric>;
deduplicationIntervalInSec: number; deduplicationIntervalInSec: number;
scheduledUpdateInterval?: number;
useLatestTs: boolean; useLatestTs: boolean;
output: CalculatedFieldOutput & { decimalsByDefault?: number; }; output: CalculatedFieldOutput & { decimalsByDefault?: number; };
} }

2
ui-ngx/src/app/shared/models/tenant.model.ts

@ -175,7 +175,7 @@ export function createTenantProfileConfiguration(type: TenantProfileType): Tenan
maxArgumentsPerCF: 10, maxArgumentsPerCF: 10,
maxDataPointsPerRollingArg: 1000, maxDataPointsPerRollingArg: 1000,
maxRelationLevelPerCfArgument: 10, maxRelationLevelPerCfArgument: 10,
minAllowedDeduplicationIntervalInSecForCF: 3600, minAllowedDeduplicationIntervalInSecForCF: 60,
maxRelatedEntitiesToReturnPerCfArgument: 100, maxRelatedEntitiesToReturnPerCfArgument: 100,
minAllowedScheduledUpdateIntervalInSecForCF: 0, minAllowedScheduledUpdateIntervalInSecForCF: 0,
maxStateSizeInKBytes: 32, maxStateSizeInKBytes: 32,

Loading…
Cancel
Save