|
|
|
@ -64,7 +64,6 @@ import static org.thingsboard.server.common.data.cf.configuration.geofencing.Ent |
|
|
|
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; |
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; |
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; |
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformAggSingleArgument; |
|
|
|
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; |
|
|
|
|
|
|
|
@Data |
|
|
|
@ -98,7 +97,7 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
Map<String, ListenableFuture<ArgumentEntry>> argFutures = switch (ctx.getCfType()) { |
|
|
|
case GEOFENCING -> fetchGeofencingCalculatedFieldArguments(ctx, entityId, false, ts); |
|
|
|
case SIMPLE, SCRIPT, ALARM, PROPAGATION -> getBaseCalculatedFieldArguments(ctx, entityId, ts); |
|
|
|
case RELATED_ENTITIES_AGGREGATION -> fetchAggArguments(ctx, entityId, ts); |
|
|
|
case RELATED_ENTITIES_AGGREGATION -> fetchRelatedEntitiesAggArguments(ctx, entityId, ts); |
|
|
|
}; |
|
|
|
if (ctx.getCfType() == PROPAGATION) { |
|
|
|
argFutures.put(PROPAGATION_CONFIG_ARGUMENT, fetchPropagationCalculatedFieldArgument(ctx, entityId)); |
|
|
|
@ -128,23 +127,6 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
return resolveOwnerArgument(tenantId, entityId); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<List<EntityId>> resolveRelatedEntities(TenantId tenantId, EntityId entityId, RelationPathLevel relation) { |
|
|
|
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findByRelationPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation))); |
|
|
|
|
|
|
|
return Futures.transform(relationsFut, relations -> { |
|
|
|
if (relations == null) { |
|
|
|
return new ArrayList<>(); |
|
|
|
} |
|
|
|
|
|
|
|
return switch (relation.direction()) { |
|
|
|
case FROM -> relations.stream() |
|
|
|
.map(EntityRelation::getTo) |
|
|
|
.toList(); |
|
|
|
case TO -> relations.isEmpty() ? List.of() : List.of(relations.get(0).getFrom()); |
|
|
|
}; |
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
protected Map<String, ArgumentEntry> resolveArgumentFutures(Map<String, ListenableFuture<ArgumentEntry>> argFutures) { |
|
|
|
return argFutures.entrySet().stream() |
|
|
|
.collect(Collectors.toMap( |
|
|
|
@ -189,7 +171,7 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
return argFutures; |
|
|
|
} |
|
|
|
|
|
|
|
protected Map<String, ListenableFuture<ArgumentEntry>> fetchAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { |
|
|
|
protected Map<String, ListenableFuture<ArgumentEntry>> fetchRelatedEntitiesAggArguments(CalculatedFieldCtx ctx, EntityId entityId, long ts) { |
|
|
|
RelatedEntitiesAggregationCalculatedFieldConfiguration aggConfig = (RelatedEntitiesAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); |
|
|
|
|
|
|
|
ListenableFuture<List<EntityId>> relatedEntitiesFut = resolveRelatedEntities(ctx.getTenantId(), entityId, aggConfig.getRelation()); |
|
|
|
@ -197,10 +179,27 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
return aggConfig.getArguments().entrySet().stream() |
|
|
|
.collect(Collectors.toMap( |
|
|
|
Map.Entry::getKey, |
|
|
|
entry -> Futures.transformAsync(relatedEntitiesFut, relatedEntities -> fetchAggArgumentEntry(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) { |
|
|
|
ListenableFuture<List<EntityRelation>> relationsFut = relationService.findByRelationPathQueryAsync(tenantId, new EntityRelationPathQuery(entityId, List.of(relation))); |
|
|
|
|
|
|
|
return Futures.transform(relationsFut, relations -> { |
|
|
|
if (relations == null) { |
|
|
|
return new ArrayList<>(); |
|
|
|
} |
|
|
|
|
|
|
|
return switch (relation.direction()) { |
|
|
|
case FROM -> relations.stream() |
|
|
|
.map(EntityRelation::getTo) |
|
|
|
.toList(); |
|
|
|
case TO -> relations.isEmpty() ? List.of() : List.of(relations.get(0).getFrom()); |
|
|
|
}; |
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<List<EntityId>> resolveGeofencingEntityIds(TenantId tenantId, EntityId entityId, Map.Entry<String, Argument> entry) { |
|
|
|
Argument value = entry.getValue(); |
|
|
|
if (value.getRefEntityId() != null) { |
|
|
|
@ -252,11 +251,11 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue))), MoreExecutors.directExecutor()); |
|
|
|
} |
|
|
|
|
|
|
|
public ListenableFuture<ArgumentEntry> fetchAggArgumentEntry(TenantId tenantId, List<EntityId> aggEntities, Argument argument, long startTs) { |
|
|
|
public ListenableFuture<ArgumentEntry> fetchRelatedEntitiesArgumentEntry(TenantId tenantId, List<EntityId> aggEntities, Argument argument, long startTs) { |
|
|
|
List<ListenableFuture<Map.Entry<EntityId, ArgumentEntry>>> futures = aggEntities.stream() |
|
|
|
.map(entityId -> { |
|
|
|
ListenableFuture<ArgumentEntry> singleAggEntryFut = fetchSingleAggArgumentEntry(tenantId, entityId, argument, startTs); |
|
|
|
return Futures.transform(singleAggEntryFut, singleAggEntry -> Map.entry(entityId, singleAggEntry), MoreExecutors.directExecutor()); |
|
|
|
ListenableFuture<ArgumentEntry> argumentEntryFut = fetchArgumentValue(tenantId, entityId, argument, startTs); |
|
|
|
return Futures.transform(argumentEntryFut, argumentEntry -> Map.entry(entityId, ArgumentEntry.createAggSingleArgument(entityId, argumentEntry)), MoreExecutors.directExecutor()); |
|
|
|
}) |
|
|
|
.toList(); |
|
|
|
|
|
|
|
@ -321,34 +320,4 @@ public abstract class AbstractCalculatedFieldProcessingService { |
|
|
|
return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchSingleAggArgumentEntry(TenantId tenantId, EntityId entityId, Argument argument, long startTs) { |
|
|
|
return switch (argument.getRefEntityKey().getType()) { |
|
|
|
case TS_ROLLING -> throw new IllegalStateException("TS_ROLLING is not supported for aggregation"); |
|
|
|
case ATTRIBUTE -> fetchAttributeAggEntry(tenantId, entityId, argument, startTs); |
|
|
|
case TS_LATEST -> fetchTsLatestAggEntry(tenantId, entityId, argument, startTs); |
|
|
|
}; |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchAttributeAggEntry(TenantId tenantId, EntityId entityId, Argument argument, long defaultLastUpdateTs) { |
|
|
|
log.trace("[{}][{}] Fetching attribute for key {}", tenantId, entityId, argument.getRefEntityKey()); |
|
|
|
var attributeOptFuture = attributesService.find(tenantId, entityId, argument.getRefEntityKey().getScope(), argument.getRefEntityKey().getKey()); |
|
|
|
return Futures.transform(attributeOptFuture, attrOpt -> { |
|
|
|
log.debug("[{}][{}] Fetched attribute for key {}: {}", tenantId, entityId, argument.getRefEntityKey(), attrOpt); |
|
|
|
AttributeKvEntry attributeKvEntry = attrOpt.orElseGet(() -> new BaseAttributeKvEntry(createDefaultKvEntry(argument), defaultLastUpdateTs, SingleValueArgumentEntry.DEFAULT_VERSION)); |
|
|
|
return transformAggSingleArgument(entityId, Optional.of(attributeKvEntry)); |
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchTsLatestAggEntry(TenantId tenantId, EntityId entityId, Argument argument, long defaultTs) { |
|
|
|
String key = argument.getRefEntityKey().getKey(); |
|
|
|
log.trace("[{}][{}] Fetching latest timeseries {}", tenantId, entityId, key); |
|
|
|
return Futures.transform( |
|
|
|
timeseriesService.findLatest(tenantId, entityId, key), |
|
|
|
result -> { |
|
|
|
log.debug("[{}][{}] Fetched latest timeseries {}: {}", tenantId, entityId, key, result); |
|
|
|
Optional<TsKvEntry> tsKvEntry = result.or(() -> Optional.of(new BasicTsKvEntry(defaultTs, createDefaultKvEntry(argument), SingleValueArgumentEntry.DEFAULT_VERSION))); |
|
|
|
return transformAggSingleArgument(entityId, tsKvEntry); |
|
|
|
}, calculatedFieldCallbackExecutor); |
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|