Browse Source

implemented update telemetry for entity by profile

pull/12092/head
IrynaMatveieva 2 years ago
parent
commit
5c8e425ab0
  1. 3
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  2. 4
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  3. 61
      application/src/main/java/org/thingsboard/server/service/telemetry/DefaultTelemetrySubscriptionService.java

3
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<String, KvEntry> updatedTelemetry);
void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map<String, KvEntry> updatedTelemetry);
void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback);

4
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<String, KvEntry> updatedTelemetry) {
public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map<String, KvEntry> 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<String, ArgumentEntry> 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);

61
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<? extends KvEntry> telemetry) {
if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) {
List<CalculatedFieldLink> 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<CalculatedFieldLink> 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<String, KvEntry> 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<String, String> attributes, Map<String, String> 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<TsKvEntry> ts) {
if (EntityType.DEVICE.equals(entityId.getEntityType()) || EntityType.ASSET.equals(entityId.getEntityType())) {
Futures.addCallback(this.tbEntityViewService.findEntityViewsByTenantIdAndEntityIdAsync(tenantId, entityId),

Loading…
Cancel
Save