From 19f6f323260347e5af7eb0317dbe168dd9cb8e52 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Mon, 16 Dec 2024 12:47:55 +0200 Subject: [PATCH] moved logic when telemetry update to main cf service --- .../cf/CalculatedFieldExecutionService.java | 5 +- .../service/cf/CalculatedFieldResult.java | 6 +- ...efaultCalculatedFieldExecutionService.java | 149 +++++++++++++----- .../cf/ctx/CalculatedFieldEntityCtx.java | 5 +- .../service/cf/ctx/state/ArgumentEntry.java | 2 + .../ctx/state/BaseCalculatedFieldState.java | 22 ++- .../cf/ctx/state/CalculatedFieldCtx.java | 2 +- .../cf/ctx/state/CalculatedFieldState.java | 7 +- .../ctx/state/ScriptCalculatedFieldState.java | 35 ++-- .../ctx/state/SimpleCalculatedFieldState.java | 37 ++--- .../ctx/state/SingleValueArgumentEntry.java | 9 +- .../cf/ctx/state/TsRollingArgumentEntry.java | 4 + .../DefaultTelemetrySubscriptionService.java | 78 +-------- 13 files changed, 186 insertions(+), 175 deletions(-) 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 5a85529f6b..302e4b2511 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 @@ -15,20 +15,19 @@ */ 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; import org.thingsboard.server.gen.transport.TransportProtos; -import java.util.Map; +import java.util.List; public interface CalculatedFieldExecutionService { void onCalculatedFieldMsg(TransportProtos.CalculatedFieldMsgProto proto, TbCallback callback); - void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry); + void onTelemetryUpdate(TenantId tenantId, EntityId entityId, List telemetry); void onEntityProfileChanged(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java index 1f8a06c8fa..e8ea318bf6 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java @@ -23,9 +23,9 @@ import java.util.Map; @Data public final class CalculatedFieldResult { - private String type; - private AttributeScope scope; - private Map resultMap; + private final String type; + private final AttributeScope scope; + private final Map resultMap; public CalculatedFieldResult(String type, AttributeScope scope, Map resultMap) { this.type = type; 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 050c2d7e73..c1a86ac36f 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 @@ -36,16 +36,20 @@ import org.thingsboard.script.api.tbel.TbelInvokeService; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.cf.CalculatedField; +import org.thingsboard.server.common.data.cf.CalculatedFieldLink; import org.thingsboard.server.common.data.cf.CalculatedFieldType; import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; +import org.thingsboard.server.common.data.id.AssetId; import org.thingsboard.server.common.data.id.AssetProfileId; import org.thingsboard.server.common.data.id.CalculatedFieldId; +import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; @@ -76,6 +80,8 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; +import org.thingsboard.server.service.profile.TbAssetProfileCache; +import org.thingsboard.server.service.profile.TbDeviceProfileCache; import java.util.ArrayList; import java.util.Collections; @@ -102,6 +108,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private final CalculatedFieldService calculatedFieldService; private final AssetService assetService; private final DeviceService deviceService; + private final TbAssetProfileCache assetProfileCache; + private final TbDeviceProfileCache deviceProfileCache; private final AttributesService attributesService; private final TimeseriesService timeseriesService; private final RocksDBService rocksDBService; @@ -112,6 +120,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private ListeningExecutorService calculatedFieldCallbackExecutor; private final ConcurrentMap calculatedFields = new ConcurrentHashMap<>(); + private final ConcurrentMap> calculatedFieldLinks = new ConcurrentHashMap<>(); private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); private final ConcurrentMap states = new ConcurrentHashMap<>(); @@ -130,6 +139,16 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); + scheduledExecutor.submit(this::fetchCalculatedFields); + } + + private void fetchCalculatedFields() { + PageDataIterable cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); + cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf)); + PageDataIterable cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); + cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); + rocksDBService.getAll().forEach((ctxId, ctx) -> states.put(JacksonUtil.fromString(ctxId, CalculatedFieldEntityCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldEntityCtx.class))); + states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId()))); } @PreDestroy @@ -216,30 +235,77 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } @Override - public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, List telemetry) { try { - log.info("Received telemetry update msg: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); - CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); - CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); - Map argumentValues = updatedTelemetry.entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); - - EntityId cfEntityId = calculatedField.getEntityId(); - switch (cfEntityId.getEntityType()) { - case ASSET_PROFILE, DEVICE_PROFILE -> { - boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); - if (isCommonEntity) { - getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); - } else { - updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); - } + EntityType entityType = entityId.getEntityType(); + if (EntityType.DEVICE.equals(entityType) || EntityType.ASSET.equals(entityType) || EntityType.CUSTOMER.equals(entityType) || EntityType.TENANT.equals(entityType)) { + EntityId profileId = null; + if (EntityType.ASSET.equals(entityType)) { + profileId = assetProfileCache.get(tenantId, (AssetId) entityId).getId(); + } else if (EntityType.DEVICE.equals(entityType)) { + profileId = deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); } - default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); + List cfLinks = new ArrayList<>(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId)); + Optional.ofNullable(profileId).ifPresent(id -> cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, id))); + cfLinks.forEach(link -> { + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + Map attributes = link.getConfiguration().getAttributes(); + Map timeSeries = link.getConfiguration().getTimeSeries(); + Map updatedTelemetry = telemetry.stream() + .filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) + .collect(Collectors.toMap( + entry -> getMappedKey(entry, attributes, timeSeries), + entry -> entry, + (v1, v2) -> v1 + )); + + if (!updatedTelemetry.isEmpty()) { + executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, updatedTelemetry); + } + }); } - log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); } catch (Exception e) { - log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); + log.trace("Failed to update telemetry entityId: [{}]", entityId, e); + } + } + + private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map updatedTelemetry) { + log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", tenantId, entityId, calculatedFieldId); + CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); + CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); + Map argumentValues = updatedTelemetry.entrySet().stream() + .collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); + + EntityId cfEntityId = calculatedField.getEntityId(); + switch (cfEntityId.getEntityType()) { + case ASSET_PROFILE, DEVICE_PROFILE -> { + boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); + if (isCommonEntity) { + getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); + } else { + updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); + } + } + default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); + } + log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); + } + + 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(); } @Override @@ -493,26 +559,29 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas if (state == null) { state = createStateByType(calculatedFieldCtx.getCfType()); } - state.initState(argumentValues); - calculatedFieldEntityCtx.setState(state); - states.put(entityCtxId, calculatedFieldEntityCtx); - rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); - - ListenableFuture resultFuture = state.performCalculation(calculatedFieldCtx); - Futures.addCallback(resultFuture, new FutureCallback<>() { - @Override - public void onSuccess(CalculatedFieldResult result) { - if (result != null) { - pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); - } - } + if (state.updateState(argumentValues)) { + calculatedFieldEntityCtx.setState(state); + states.put(entityCtxId, calculatedFieldEntityCtx); + rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); + + boolean allArgsPresent = calculatedFieldCtx.getArguments().keySet().containsAll(state.getArguments().keySet()); + if (allArgsPresent) { + ListenableFuture resultFuture = state.performCalculation(calculatedFieldCtx); + Futures.addCallback(resultFuture, new FutureCallback<>() { + @Override + public void onSuccess(CalculatedFieldResult result) { + if (result != null) { + pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); + } + } - @Override - public void onFailure(Throwable t) { - log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); + @Override + public void onFailure(Throwable t) { + log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); + } + }, MoreExecutors.directExecutor()); } - }, MoreExecutors.directExecutor()); - + } } private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId) { @@ -520,14 +589,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas if (stateStr == null) { return new CalculatedFieldEntityCtx(entityCtxId, null); } - return JacksonUtil.fromString(rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)), CalculatedFieldEntityCtx.class); + return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); } private void pushMsgToRuleEngine(TenantId tenantId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult) { try { String type = calculatedFieldResult.getType(); - TbMsgType msgType = "ATTRIBUTES".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; - TbMsgMetaData md = "ATTRIBUTES".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; + TbMsgType msgType = "ATTRIBUTE".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; + TbMsgMetaData md = "ATTRIBUTE".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; ObjectNode payload = createJsonPayload(calculatedFieldResult); TbMsg msg = TbMsg.newMsg(msgType, originatorId, md, JacksonUtil.writeValueAsString(payload)); clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java index 7a8384b6bf..e6dc021951 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/CalculatedFieldEntityCtx.java @@ -16,17 +16,16 @@ package org.thingsboard.server.service.cf.ctx; import lombok.Data; +import lombok.NoArgsConstructor; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @Data +@NoArgsConstructor public class CalculatedFieldEntityCtx { private CalculatedFieldEntityCtxId id; private CalculatedFieldState state; - public CalculatedFieldEntityCtx() { - } - public CalculatedFieldEntityCtx(CalculatedFieldEntityCtxId id, CalculatedFieldState state) { this.id = id; this.state = state; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java index f70d614123..78222244c9 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ArgumentEntry.java @@ -41,6 +41,8 @@ public interface ArgumentEntry { Object getValue(); + boolean hasUpdatedValue(ArgumentEntry entry); + static ArgumentEntry createSingleValueArgument(KvEntry kvEntry) { return new SingleValueArgumentEntry(kvEntry); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java index 59b007a420..ae6fc9033a 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/BaseCalculatedFieldState.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.cf.ctx.state; import java.util.HashMap; import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; public abstract class BaseCalculatedFieldState implements CalculatedFieldState { @@ -31,26 +32,39 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { } @Override - public void initState(Map argumentValues) { + public boolean updateState(Map argumentValues) { if (arguments == null) { arguments = new HashMap<>(); } + AtomicBoolean stateUpdated = new AtomicBoolean(false); argumentValues.forEach((key, argumentEntry) -> { ArgumentEntry existingArgumentEntry = arguments.get(key); if (existingArgumentEntry != null) { if (existingArgumentEntry instanceof SingleValueArgumentEntry) { - arguments.put(key, argumentEntry); + if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { + arguments.put(key, argumentEntry); + stateUpdated.set(true); + } } else if (existingArgumentEntry instanceof TsRollingArgumentEntry existingTsRollingArgumentEntry) { if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - existingTsRollingArgumentEntry.getTsRecords().putAll(tsRollingArgumentEntry.getTsRecords()); + if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { + existingTsRollingArgumentEntry.getTsRecords().putAll(tsRollingArgumentEntry.getTsRecords()); + stateUpdated.set(true); + } } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - existingTsRollingArgumentEntry.getTsRecords().put(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); + if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { + existingTsRollingArgumentEntry.getTsRecords().put(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); + stateUpdated.set(true); + } + } } } else { arguments.put(key, argumentEntry); + stateUpdated.set(true); } }); + return stateUpdated.get(); } } 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 b436e0421e..2cd5c68144 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 @@ -35,7 +35,7 @@ public class CalculatedFieldCtx { private TenantId tenantId; private EntityId entityId; private CalculatedFieldType cfType; - private Map arguments; + private final Map arguments; private Output output; private String expression; private TbelInvokeService tbelInvokeService; diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java index a5ac6b2c47..3c4a680df1 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldState.java @@ -20,7 +20,6 @@ import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.google.common.util.concurrent.ListenableFuture; import org.thingsboard.server.common.data.cf.CalculatedFieldType; -import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.service.cf.CalculatedFieldResult; import java.util.Map; @@ -41,11 +40,7 @@ public interface CalculatedFieldState { Map getArguments(); - default boolean isValid(Map arguments) { - return getArguments().keySet().containsAll(arguments.keySet()); - } - - void initState(Map argumentValues); + boolean updateState(Map argumentValues); ListenableFuture performCalculation(CalculatedFieldCtx ctx); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java index b9b98f9c5e..87429050de 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java @@ -39,26 +39,23 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { - if (isValid(ctx.getArguments())) { - arguments.forEach((key, argumentEntry) -> { - if (argumentEntry instanceof TsRollingArgumentEntry) { - Argument argument = ctx.getArguments().get(key); - TreeMap tsRecords = ((TsRollingArgumentEntry) argumentEntry).getTsRecords(); - if (tsRecords.size() > argument.getLimit()) { - tsRecords.pollFirstEntry(); - } - tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); + arguments.forEach((key, argumentEntry) -> { + if (argumentEntry instanceof TsRollingArgumentEntry) { + Argument argument = ctx.getArguments().get(key); + TreeMap tsRecords = ((TsRollingArgumentEntry) argumentEntry).getTsRecords(); + if (tsRecords.size() > argument.getLimit()) { + tsRecords.pollFirstEntry(); } - }); - Object[] args = arguments.values().stream().map(ArgumentEntry::getValue).toArray(); - ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); - Output output = ctx.getOutput(); - return Futures.transform(resultFuture, - result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), - MoreExecutors.directExecutor() - ); - } - return Futures.immediateFuture(null); + tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); + } + }); + Object[] args = arguments.values().stream().map(ArgumentEntry::getValue).toArray(); + ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); + Output output = ctx.getOutput(); + return Futures.transform(resultFuture, + result -> new CalculatedFieldResult(output.getType(), output.getScope(), result), + MoreExecutors.directExecutor() + ); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java index 491419b40a..e16d310b3e 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java @@ -37,27 +37,24 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState { @Override public ListenableFuture performCalculation(CalculatedFieldCtx ctx) { - if (isValid(ctx.getArguments())) { - String expression = ctx.getExpression(); - ThreadLocal customExpression = new ThreadLocal<>(); - var expr = customExpression.get(); - if (expr == null) { - expr = new ExpressionBuilder(expression) - .implicitMultiplication(true) - .variables(this.arguments.keySet()) - .build(); - customExpression.set(expr); - } - Map variables = new HashMap<>(); - this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); - expr.setVariables(variables); - - double expressionResult = expr.evaluate(); - - Output output = ctx.getOutput(); - return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); + String expression = ctx.getExpression(); + ThreadLocal customExpression = new ThreadLocal<>(); + var expr = customExpression.get(); + if (expr == null) { + expr = new ExpressionBuilder(expression) + .implicitMultiplication(true) + .variables(this.arguments.keySet()) + .build(); + customExpression.set(expr); } - return Futures.immediateFuture(null); + Map variables = new HashMap<>(); + this.arguments.forEach((k, v) -> variables.put(k, Double.parseDouble(v.getValue().toString()))); + expr.setVariables(variables); + + double expressionResult = expr.evaluate(); + + Output output = ctx.getOutput(); + return Futures.immediateFuture(new CalculatedFieldResult(output.getType(), output.getScope(), Map.of(output.getName(), expressionResult))); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java index e0db8c50fb..e6cb24b970 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SingleValueArgumentEntry.java @@ -16,19 +16,18 @@ package org.thingsboard.server.service.cf.ctx.state; import lombok.Data; +import lombok.NoArgsConstructor; import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @Data +@NoArgsConstructor public class SingleValueArgumentEntry implements ArgumentEntry { private long ts; private Object value; - public SingleValueArgumentEntry() { - } - public SingleValueArgumentEntry(KvEntry entry) { if (entry instanceof TsKvEntry) { this.ts = ((TsKvEntry) entry).getTs(); @@ -48,4 +47,8 @@ public class SingleValueArgumentEntry implements ArgumentEntry { return value; } + @Override + public boolean hasUpdatedValue(ArgumentEntry entry) { + return this.ts != ((SingleValueArgumentEntry) entry).getTs(); + } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java index 1166da113d..104d0ae90c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/ctx/state/TsRollingArgumentEntry.java @@ -40,4 +40,8 @@ public class TsRollingArgumentEntry implements ArgumentEntry { return tsRecords; } + @Override + public boolean hasUpdatedValue(ArgumentEntry entry) { + return !tsRecords.containsKey(((SingleValueArgumentEntry) entry).getTs()); + } } 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 ad64204e0f..b94b319e71 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 @@ -32,11 +32,7 @@ import org.thingsboard.server.common.data.ApiUsageRecordKey; 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; @@ -44,7 +40,6 @@ import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DoubleDataEntry; -import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; @@ -52,14 +47,11 @@ import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.msg.queue.TbCallback; import org.thingsboard.server.common.stats.TbApiUsageReportClient; import org.thingsboard.server.dao.attributes.AttributesService; -import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.timeseries.TimeseriesService; 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; @@ -73,7 +65,6 @@ import java.util.Objects; import java.util.Optional; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.stream.Collectors; /** * Created by ashvayka on 27.03.18. @@ -87,11 +78,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer private final TbEntityViewService tbEntityViewService; private final TbApiUsageReportClient apiUsageClient; private final TbApiUsageStateService apiUsageStateService; - private final CalculatedFieldService calculatedFieldService; private final CalculatedFieldExecutionService calculatedFieldExecutionService; - private final TbAssetProfileCache assetProfileCache; - private final TbDeviceProfileCache deviceProfileCache; - private ExecutorService tsCallBackExecutor; @Value("${sql.ts.value_no_xss_validation:false}") @@ -102,19 +89,13 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer @Lazy TbEntityViewService tbEntityViewService, TbApiUsageReportClient apiUsageClient, TbApiUsageStateService apiUsageStateService, - CalculatedFieldService calculatedFieldService, - CalculatedFieldExecutionService calculatedFieldExecutionService, - TbAssetProfileCache assetProfileCache, - TbDeviceProfileCache deviceProfileCache) { + CalculatedFieldExecutionService calculatedFieldExecutionService) { this.attrService = attrService; this.tsService = tsService; this.tbEntityViewService = tbEntityViewService; this.apiUsageClient = apiUsageClient; this.apiUsageStateService = apiUsageStateService; - this.calculatedFieldService = calculatedFieldService; this.calculatedFieldExecutionService = calculatedFieldExecutionService; - this.assetProfileCache = assetProfileCache; - this.deviceProfileCache = deviceProfileCache; } @PostConstruct @@ -201,7 +182,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer addMainCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); addEntityViewCallback(tenantId, entityId, ts); - updateTelemetryInCalculatedFields(tenantId, entityId, ts); + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, ts); } private void saveWithoutLatestAndNotifyInternal(TenantId tenantId, EntityId entityId, List ts, long ttl, FutureCallback callback) { @@ -210,55 +191,6 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); } - private void updateTelemetryInCalculatedFields(TenantId tenantId, EntityId entityId, List telemetry) { - EntityType entityType = entityId.getEntityType(); - if (EntityType.DEVICE.equals(entityType) || EntityType.ASSET.equals(entityType) || EntityType.CUSTOMER.equals(entityType) || EntityType.TENANT.equals(entityType)) { - EntityId profileId = null; - if (EntityType.ASSET.equals(entityType)) { - profileId = assetProfileCache.get(tenantId, (AssetId) entityId).getId(); - } else if (EntityType.DEVICE.equals(entityType)) { - profileId = deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); - } - List cfLinks = new ArrayList<>(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId)); - Optional.ofNullable(profileId).ifPresent(id -> cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, id))); - if (!cfLinks.isEmpty()) { - cfLinks.forEach(link -> { - CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); - Map attributes = link.getConfiguration().getAttributes(); - Map timeSeries = link.getConfiguration().getTimeSeries(); - Map updatedTelemetry = telemetry.stream() - .filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) - .collect(Collectors.toMap( - entry -> getMappedKey(entry, attributes, timeSeries), - entry -> entry, - (v1, v2) -> v1 - )); - - if (!updatedTelemetry.isEmpty()) { - 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), @@ -335,7 +267,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = attrService.save(tenantId, entityId, scope, attributes); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope, attributes, notifyDevice)); - updateTelemetryInCalculatedFields(tenantId, entityId, attributes); + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, attributes); } @Override @@ -343,7 +275,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = attrService.save(tenantId, entityId, scope, attributes); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onAttributesUpdate(tenantId, entityId, scope.name(), attributes, notifyDevice)); - updateTelemetryInCalculatedFields(tenantId, entityId, attributes); + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, attributes); } @Override @@ -357,7 +289,7 @@ public class DefaultTelemetrySubscriptionService extends AbstractSubscriptionSer ListenableFuture> saveFuture = tsService.saveLatest(tenantId, entityId, ts); addVoidCallback(saveFuture, callback); addWsCallback(saveFuture, success -> onTimeSeriesUpdate(tenantId, entityId, ts)); - updateTelemetryInCalculatedFields(tenantId, entityId, ts); + calculatedFieldExecutionService.onTelemetryUpdate(tenantId, entityId, ts); } @Override