From 9ef68584c9aecb19659496210cf914a7eb5f81b5 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Wed, 22 Jan 2025 09:40:12 +0200 Subject: [PATCH] updated getMappedTelemetry method --- ...efaultCalculatedFieldExecutionService.java | 19 ++++++++++--------- .../cf/ctx/state/CalculatedFieldCtx.java | 2 +- ...CalculatedFieldAttributeUpdateRequest.java | 15 ++------------- ...CalculatedFieldTelemetryUpdateRequest.java | 5 +---- ...alculatedFieldTimeSeriesUpdateRequest.java | 13 +++---------- 5 files changed, 17 insertions(+), 37 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java index 7c17d54649..84a21ab315 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java @@ -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 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 updatedTelemetry = request.getMappedTelemetry(ctx); + Map 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> tpiStates) { TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, request.getTenantId(), targetEntity); if (targetEntityTpi.isMyPartition()) { - Map updatedTelemetry = request.getMappedTelemetry(ctx); + Map updatedTelemetry = request.getMappedTelemetry(ctx, request.getEntityId()); if (!updatedTelemetry.isEmpty()) { - executeTelemetryUpdate(ctx, request.getEntityId(), request.getPreviousCalculatedFieldIds(), updatedTelemetry); + executeTelemetryUpdate(ctx, targetEntity, request.getPreviousCalculatedFieldIds(), updatedTelemetry); } } else { List 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 updatedTelemetry = request.getMappedTelemetry(ctx); + Map 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 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); 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 e17a1a61f9..cb4052b7df 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 @@ -57,7 +57,7 @@ public class CalculatedFieldCtx { this.arguments = configuration.getArguments(); this.referencedEntityKeys = arguments.entrySet().stream() .collect(Collectors.toMap( - entry -> new TbPair<>(entry.getValue().getRefEntityId(), entry.getValue().getRefEntityKey()), + entry -> new TbPair<>(entry.getValue().getRefEntityId() == null ? entityId : entry.getValue().getRefEntityId(), entry.getValue().getRefEntityKey()), Map.Entry::getKey )); this.argNames = new ArrayList<>(arguments.keySet()); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java index 6050370fd2..d2eb31cd6d 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldAttributeUpdateRequest.java @@ -19,7 +19,6 @@ import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.server.common.data.AttributeScope; -import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.id.CalculatedFieldId; @@ -52,16 +51,7 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel } @Override - public Map getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration) { - return switch (scope) { - case CLIENT_SCOPE -> linkConfiguration.getClientAttributes(); - case SERVER_SCOPE -> linkConfiguration.getServerAttributes(); - case SHARED_SCOPE -> linkConfiguration.getSharedAttributes(); - }; - } - - @Override - public Map getMappedTelemetry(CalculatedFieldCtx ctx) { + public Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId) { Map mappedKvEntries = new HashMap<>(); Map, String> referencedKeys = ctx.getReferencedEntityKeys(); @@ -70,7 +60,7 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel ReferencedEntityKey referencedEntityKey = new ReferencedEntityKey(key, ArgumentType.ATTRIBUTE, scope); - String argName = referencedKeys.get(new TbPair<>(entityId, referencedEntityKey)); + String argName = referencedKeys.get(new TbPair<>(referencedEntityId, referencedEntityKey)); if (argName != null) { mappedKvEntries.put(argName, entry); @@ -79,5 +69,4 @@ public class CalculatedFieldAttributeUpdateRequest implements CalculatedFieldTel return mappedKvEntries; } - } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java index f85117dc41..3f7250f4ef 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTelemetryUpdateRequest.java @@ -15,7 +15,6 @@ */ package org.thingsboard.server.service.cf.telemetry; -import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @@ -35,8 +34,6 @@ public interface CalculatedFieldTelemetryUpdateRequest { List getPreviousCalculatedFieldIds(); - Map getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration); - - Map getMappedTelemetry(CalculatedFieldCtx ctx); + Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java index 646145a46e..a5637c8cfd 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/telemetry/CalculatedFieldTimeSeriesUpdateRequest.java @@ -18,7 +18,6 @@ package org.thingsboard.server.service.cf.telemetry; import lombok.AllArgsConstructor; import lombok.Data; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; -import org.thingsboard.server.common.data.cf.CalculatedFieldLinkConfiguration; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.id.CalculatedFieldId; @@ -49,12 +48,7 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe } @Override - public Map getTelemetryKeysFromLink(CalculatedFieldLinkConfiguration linkConfiguration) { - return linkConfiguration.getTimeSeries(); - } - - @Override - public Map getMappedTelemetry(CalculatedFieldCtx ctx) { + public Map getMappedTelemetry(CalculatedFieldCtx ctx, EntityId referencedEntityId) { Map mappedKvEntries = new HashMap<>(); Map, String> referencedKeys = ctx.getReferencedEntityKeys(); @@ -62,13 +56,13 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe String key = entry.getKey(); ReferencedEntityKey tsLatestKey = new ReferencedEntityKey(key, ArgumentType.TS_LATEST, null); - String argTsLatestName = referencedKeys.get(new TbPair<>(entityId, tsLatestKey)); + String argTsLatestName = referencedKeys.get(new TbPair<>(referencedEntityId, tsLatestKey)); if (argTsLatestName != null) { mappedKvEntries.put(argTsLatestName, entry); } else { ReferencedEntityKey tsRollingKey = new ReferencedEntityKey(key, ArgumentType.TS_ROLLING, null); - String argTsRollingName = referencedKeys.get(new TbPair<>(entityId, tsRollingKey)); + String argTsRollingName = referencedKeys.get(new TbPair<>(referencedEntityId, tsRollingKey)); if (argTsRollingName != null) { mappedKvEntries.put(argTsRollingName, entry); @@ -78,5 +72,4 @@ public class CalculatedFieldTimeSeriesUpdateRequest implements CalculatedFieldTe return mappedKvEntries; } - }