|
|
|
@ -271,7 +271,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
case ASSET_PROFILE, DEVICE_PROFILE -> { |
|
|
|
log.info("Initializing state for all entities in profile: tenantId=[{}], profileId=[{}]", tenantId, entityId); |
|
|
|
Map<String, Argument> commonArguments = calculatedFieldCtx.getArguments().entrySet().stream() |
|
|
|
.filter(entry -> !isProfileEntity(entry.getValue().getRefEntityId())) |
|
|
|
.filter(entry -> entry.getValue().getRefEntityId() != null) |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); |
|
|
|
fetchArguments(tenantId, entityId, commonArguments, commonArgs -> { |
|
|
|
calculatedFieldCache.getEntitiesByProfile(tenantId, entityId).forEach(targetEntityId -> { |
|
|
|
@ -375,9 +375,10 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private void processCalculatedFields(CalculatedFieldTelemetryUpdateRequest request, EntityId cfTargetEntityId) { |
|
|
|
if (cfTargetEntityId != null) { |
|
|
|
calculatedFieldCache.getCalculatedFieldCtxsByEntityId(cfTargetEntityId, tbelInvokeService).forEach(ctx -> { |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx); |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, cfTargetEntityId); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
executeTelemetryUpdate(ctx, request.getEntityId(), request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
EntityId targetEntityId = isProfileEntity(cfTargetEntityId) ? request.getEntityId() : cfTargetEntityId; |
|
|
|
executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
@ -406,9 +407,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private void processCalculatedFieldLink(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldCtx ctx, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> tpiStates) { |
|
|
|
TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, request.getTenantId(), targetEntity); |
|
|
|
if (targetEntityTpi.isMyPartition()) { |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx); |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
executeTelemetryUpdate(ctx, request.getEntityId(), request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
executeTelemetryUpdate(ctx, targetEntity, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
} else { |
|
|
|
List<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); |
|
|
|
@ -427,13 +428,13 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
proto.getLinksList().forEach(ctxIdProto -> { |
|
|
|
EntityId entityId = request.getEntityId(); |
|
|
|
CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(ctxIdProto.getCalculatedFieldIdMSB(), ctxIdProto.getCalculatedFieldIdLSB())); |
|
|
|
CalculatedFieldCtx ctx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); |
|
|
|
|
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx); |
|
|
|
Map<String, KvEntry> updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); |
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
executeTelemetryUpdate(ctx, entityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
EntityId targetEntityId = EntityIdFactory.getByTypeAndUuid(ctxIdProto.getEntityType(), new UUID(ctxIdProto.getEntityIdMSB(), ctxIdProto.getEntityIdLSB())); |
|
|
|
executeTelemetryUpdate(ctx, targetEntityId, request.getPreviousCalculatedFieldIds(), updatedTelemetry); |
|
|
|
} |
|
|
|
}); |
|
|
|
} catch (Exception e) { |
|
|
|
@ -654,7 +655,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
|
|
|
|
private ListenableFuture<ArgumentEntry> fetchArgumentValue(TenantId tenantId, EntityId targetEntityId, Argument argument) { |
|
|
|
EntityId argumentEntityId = argument.getRefEntityId(); |
|
|
|
EntityId entityId = isProfileEntity(argumentEntityId) |
|
|
|
EntityId entityId = (argumentEntityId == null || isProfileEntity(argumentEntityId)) |
|
|
|
? targetEntityId |
|
|
|
: argumentEntityId; |
|
|
|
return fetchKvEntry(tenantId, entityId, argument); |
|
|
|
|