diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index a1202ec74b..b605ee909e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -16,6 +16,7 @@ package org.thingsboard.server.service.cf; import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.msg.queue.TbCallback; @@ -27,7 +28,7 @@ public interface CalculatedFieldExecutionService { void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); - void onTelemetryUpdate(TenantId tenantId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); 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 3bb9be065d..d821f632fe 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 @@ -225,7 +225,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } @Override - public void onTelemetryUpdate(TenantId tenantId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { try { log.info("Received telemetry update msg: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> { @@ -234,7 +234,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas }); Map argumentValues = updatedTelemetry.entrySet().stream() .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); - updateOrInitializeState(calculatedFieldCtx, calculatedFieldCtx.getEntityId(), argumentValues); + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); log.info("Successfully updated time series for calculatedFieldId: [{}]", calculatedFieldId); } catch (Exception e) { log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java index 76c1383f6a..f6c19eb078 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java @@ -33,8 +33,10 @@ import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; import org.thingsboard.server.common.data.cf.CalculatedFieldLink; +import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.CustomerId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; @@ -56,6 +58,8 @@ import org.thingsboard.server.dao.util.KvUtils; import org.thingsboard.server.service.apiusage.TbApiUsageStateService; import org.thingsboard.server.service.cf.CalculatedFieldExecutionService; import org.thingsboard.server.service.entitiy.entityview.TbEntityViewService; +import org.thingsboard.server.service.profile.TbAssetProfileCache; +import org.thingsboard.server.service.profile.TbDeviceProfileCache; import org.thingsboard.server.service.subscription.TbSubscriptionUtils; import java.util.ArrayList; @@ -85,6 +89,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer private final TbApiUsageStateService apiUsageStateService; private final CalculatedFieldService calculatedFieldService; private final CalculatedFieldExecutionService calculatedFieldExecutionService; + private final TbAssetProfileCache assetProfileCache; + private final TbDeviceProfileCache deviceProfileCache; private ExecutorService tsCallBackExecutor; @@ -97,7 +103,9 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer TbApiUsageReportClient apiUsageClient, TbApiUsageStateService apiUsageStateService, CalculatedFieldService calculatedFieldService, - CalculatedFieldExecutionService calculatedFieldExecutionService) { + CalculatedFieldExecutionService calculatedFieldExecutionService, + TbAssetProfileCache assetProfileCache, + TbDeviceProfileCache deviceProfileCache) { this.attrService = attrService; this.tsService = tsService; this.tbEntityViewService = tbEntityViewService; @@ -105,6 +113,8 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer this.apiUsageStateService = apiUsageStateService; this.calculatedFieldService = calculatedFieldService; this.calculatedFieldExecutionService = calculatedFieldExecutionService; + this.assetProfileCache = assetProfileCache; + this.deviceProfileCache = deviceProfileCache; } @PostConstruct @@ -201,8 +211,17 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer } private void updateTelemetryInCalculatedFields(TenantId tenantId, EntityId entityId, List telemetry) { - if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) { - List cfLinks = calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId); + EntityType entityType = entityId.getEntityType(); + if (EntityType.DEVICE.equals(entityType) || EntityType.ASSET.equals(entityType)) { + EntityId profileId; + if (EntityType.ASSET.equals(entityType)) { + profileId = assetProfileCache.get(tenantId, (AssetId) entityId).getId(); + } else { + profileId = deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); + } + List cfLinks = new ArrayList<>(); + cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId)); + cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, profileId)); if (!cfLinks.isEmpty()) { cfLinks.forEach(link -> { CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); @@ -211,34 +230,36 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer Map updatedTelemetry = telemetry.stream() .filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) .collect(Collectors.toMap( - entry -> { - if (entry instanceof AttributeKvEntry) { - return attributes.entrySet().stream() - .filter(attr -> attr.getValue().equals(entry.getKey())) - .map(Map.Entry::getKey) - .findFirst() - .orElse(entry.getKey()); - } else if (entry instanceof TsKvEntry) { - return timeSeries.entrySet().stream() - .filter(ts -> ts.getValue().equals(entry.getKey())) - .map(Map.Entry::getKey) - .findFirst() - .orElse(entry.getKey()); - } - return entry.getKey(); - }, + entry -> getMappedKey(entry, attributes, timeSeries), entry -> entry, (v1, v2) -> v1 )); if (!updatedTelemetry.isEmpty()) { - calculatedFieldExecutionService.onTelemetryUpdate(tenantId, calculatedFieldId, updatedTelemetry); + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, calculatedFieldId, updatedTelemetry); } }); } } } + private String getMappedKey(KvEntry entry, Map attributes, Map timeSeries) { + if (entry instanceof AttributeKvEntry) { + return attributes.entrySet().stream() + .filter(attr -> attr.getValue().equals(entry.getKey())) + .map(Map.Entry::getKey) + .findFirst() + .orElse(entry.getKey()); + } else if (entry instanceof TsKvEntry) { + return timeSeries.entrySet().stream() + .filter(ts -> ts.getValue().equals(entry.getKey())) + .map(Map.Entry::getKey) + .findFirst() + .orElse(entry.getKey()); + } + return entry.getKey(); + } + private void addEntityViewCallback(TenantId tenantId, EntityId entityId, List ts) { if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) { Futures.addCallback(this.tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId),