diff --git a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java index 39ab78f86a..d44d06407c 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java @@ -37,7 +37,7 @@ import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.cf.CalculatedField; import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.ArgumentType; -import org.thingsboard.server.common.data.cf.configuration.AttributeImmediateOutputStrategy; +import org.thingsboard.server.common.data.cf.configuration.AttributesImmediateOutputStrategy; import org.thingsboard.server.common.data.cf.configuration.OutputStrategy; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; @@ -79,12 +79,12 @@ import java.util.concurrent.ExecutionException; import java.util.function.Predicate; import java.util.stream.Collectors; -import static org.thingsboard.rule.engine.util.TelemetryUtil.filterChangedAttr; -import static org.thingsboard.rule.engine.util.TelemetryUtil.toTsKvEntryList; import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION; import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT; import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY; import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY; +import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr; +import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultKvEntry; import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.transformSingleValueArgument; @@ -379,7 +379,7 @@ public abstract class AbstractCalculatedFieldProcessingService { } private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List cfIds, SettableFuture future) { - if (!(outputStrategy instanceof AttributeImmediateOutputStrategy attOutputStrategy)) { + if (!(outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy)) { future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); } else { AttributesSaveRequest.Strategy strategy = new Strategy(attOutputStrategy.isSaveAttribute(), attOutputStrategy.isSendWsUpdate(), attOutputStrategy.isProcessCfs()); @@ -413,7 +413,7 @@ public abstract class AbstractCalculatedFieldProcessingService { List entries, AttributesSaveRequest.Strategy strategy, SettableFuture future) { - tsSubService.saveAttributesInternal(AttributesSaveRequest.builder() + tsSubService.saveAttributes(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(entityId) .scope(scope) @@ -452,7 +452,7 @@ public abstract class AbstractCalculatedFieldProcessingService { if (cfIds != null && !cfIds.isEmpty()) { builder.previousCalculatedFieldIds(cfIds); } - tsSubService.saveTimeseriesInternal(builder.build()); + tsSubService.saveTimeseries(builder.build()); } } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java index 4a6fc08b57..858f9eb2f3 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java @@ -39,12 +39,8 @@ public interface CalculatedFieldProcessingService { Map fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map arguments); - void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); - void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); - void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback); - void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List linkedCalculatedFields, TbCallback callback); } diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java index 302f51ecb1..851717326f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java @@ -143,8 +143,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF } } - @Override - public void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + private void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { if (result instanceof TelemetryCalculatedFieldResult telemetryResult) { saveTelemetryResult(tenantId, entityId, telemetryResult, cfIds, callback); return; @@ -157,8 +156,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF callback.onSuccess(); } - @Override - public void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { + private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List cfIds, TbCallback callback) { if (result instanceof PropagationCalculatedFieldResult propagationResult) { handlePropagationResults(propagationResult, callback, (entity, res, cb) -> sendMsgToRuleEngine(tenantId, entityId, cb, res.toTbMsg(entity, cfIds))); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeImmediateOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java similarity index 92% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeImmediateOutputStrategy.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java index 737c4fc64e..73bc65274d 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeImmediateOutputStrategy.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java @@ -22,7 +22,7 @@ import lombok.NoArgsConstructor; @Data @AllArgsConstructor @NoArgsConstructor -public class AttributeImmediateOutputStrategy implements AttributeOutputStrategy { +public class AttributesImmediateOutputStrategy implements AttributesOutputStrategy { private boolean updateAttributesOnlyOnValueChange; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java index 61195fc5f4..578af5c6ea 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutput.java @@ -25,10 +25,10 @@ public class AttributesOutput implements Output { private AttributeScope scope; private Integer decimalsByDefault; - private AttributeOutputStrategy strategy; + private AttributesOutputStrategy strategy; public AttributesOutput() { - this.strategy = new AttributeRuleChainOutputStrategy(); + this.strategy = new AttributesRuleChainOutputStrategy(); } @Override diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutputStrategy.java similarity index 79% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeOutputStrategy.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutputStrategy.java index f47bc27579..057fb7d8d7 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeOutputStrategy.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesOutputStrategy.java @@ -26,8 +26,8 @@ import com.fasterxml.jackson.annotation.JsonTypeInfo; property = "type" ) @JsonSubTypes({ - @JsonSubTypes.Type(value = AttributeImmediateOutputStrategy.class, name = "IMMEDIATE"), - @JsonSubTypes.Type(value = AttributeRuleChainOutputStrategy.class, name = "RULE_CHAIN"), + @JsonSubTypes.Type(value = AttributesImmediateOutputStrategy.class, name = "IMMEDIATE"), + @JsonSubTypes.Type(value = AttributesRuleChainOutputStrategy.class, name = "RULE_CHAIN"), }) -public interface AttributeOutputStrategy extends OutputStrategy { +public interface AttributesOutputStrategy extends OutputStrategy { } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeRuleChainOutputStrategy.java b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesRuleChainOutputStrategy.java similarity index 91% rename from common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeRuleChainOutputStrategy.java rename to common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesRuleChainOutputStrategy.java index ce03aeb750..1a3348ce74 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributeRuleChainOutputStrategy.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesRuleChainOutputStrategy.java @@ -20,7 +20,7 @@ import lombok.NoArgsConstructor; @Data @NoArgsConstructor -public class AttributeRuleChainOutputStrategy implements AttributeOutputStrategy { +public class AttributesRuleChainOutputStrategy implements AttributesOutputStrategy { @Override public OutputStrategyType getType() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java b/dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java index 8b95ddcb57..92f65c6e6d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java +++ b/dao/src/main/java/org/thingsboard/server/dao/util/KvUtils.java @@ -19,13 +19,21 @@ import com.fasterxml.jackson.databind.JsonNode; import com.github.benmanes.caffeine.cache.Cache; import com.github.benmanes.caffeine.cache.Caffeine; import org.thingsboard.server.common.data.StringUtils; +import org.thingsboard.server.common.data.kv.AttributeKvEntry; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.service.NoXssValidator; +import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.Objects; import java.util.concurrent.TimeUnit; +import java.util.function.Function; +import java.util.stream.Collectors; public class KvUtils { @@ -74,4 +82,33 @@ public class KvUtils { } } } + + public static List toTsKvEntryList(Map> tsKvMap) { + List tsKvEntryList = new ArrayList<>(); + for (Map.Entry> tsKvEntry : tsKvMap.entrySet()) { + for (KvEntry kvEntry : tsKvEntry.getValue()) { + tsKvEntryList.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry)); + } + } + return tsKvEntryList; + } + + public static List filterChangedAttr(List currentAttributes, List newAttributes) { + if (currentAttributes == null || currentAttributes.isEmpty()) { + return newAttributes; + } + + Map currentAttrMap = currentAttributes.stream() + .collect(Collectors.toMap(AttributeKvEntry::getKey, Function.identity(), (existing, replacement) -> existing)); + + return newAttributes.stream() + .filter(item -> { + AttributeKvEntry cacheAttr = currentAttrMap.get(item.getKey()); + return cacheAttr == null + || !Objects.equals(item.getValue(), cacheAttr.getValue()) //JSON and String can be equals by value, but different by type + || !Objects.equals(item.getDataType(), cacheAttr.getDataType()); + }) + .collect(Collectors.toList()); + } + } diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java index 20aa7993a1..0c73efa1b9 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgAttributesNode.java @@ -48,10 +48,10 @@ import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessin import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.Deduplicate; import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.OnEveryMessage; import static org.thingsboard.rule.engine.telemetry.settings.AttributesProcessingSettings.WebSocketsOnly; -import static org.thingsboard.rule.engine.util.TelemetryUtil.filterChangedAttr; import static org.thingsboard.server.common.data.DataConstants.NOTIFY_DEVICE_METADATA_KEY; import static org.thingsboard.server.common.data.DataConstants.SCOPE; import static org.thingsboard.server.common.data.msg.TbMsgType.POST_ATTRIBUTES_REQUEST; +import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr; @RuleNode( type = ComponentType.ACTION, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java index 32f06b1e00..80e964e893 100644 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java +++ b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/telemetry/TbMsgTimeseriesNode.java @@ -31,7 +31,6 @@ import org.thingsboard.rule.engine.telemetry.strategy.ProcessingStrategy; import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.StringUtils; import org.thingsboard.server.common.data.TenantProfile; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.plugin.ComponentType; @@ -39,7 +38,6 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon import org.thingsboard.server.common.data.util.TbPair; import org.thingsboard.server.common.msg.TbMsg; -import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.UUID; @@ -49,8 +47,8 @@ import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessin import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.Deduplicate; import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.OnEveryMessage; import static org.thingsboard.rule.engine.telemetry.settings.TimeseriesProcessingSettings.WebSocketsOnly; -import static org.thingsboard.rule.engine.util.TelemetryUtil.toTsKvEntryList; import static org.thingsboard.server.common.data.msg.TbMsgType.POST_TELEMETRY_REQUEST; +import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList; @RuleNode( type = ComponentType.ACTION, diff --git a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TelemetryUtil.java b/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TelemetryUtil.java deleted file mode 100644 index 41d6f1ce1d..0000000000 --- a/rule-engine/rule-engine-components/src/main/java/org/thingsboard/rule/engine/util/TelemetryUtil.java +++ /dev/null @@ -1,60 +0,0 @@ -/** - * Copyright © 2016-2025 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.rule.engine.util; - -import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.KvEntry; -import org.thingsboard.server.common.data.kv.TsKvEntry; - -import java.util.ArrayList; -import java.util.List; -import java.util.Map; -import java.util.Objects; -import java.util.function.Function; -import java.util.stream.Collectors; - -public class TelemetryUtil { - - public static List toTsKvEntryList(Map> tsKvMap) { - List tsKvEntryList = new ArrayList<>(); - for (Map.Entry> tsKvEntry : tsKvMap.entrySet()) { - for (KvEntry kvEntry : tsKvEntry.getValue()) { - tsKvEntryList.add(new BasicTsKvEntry(tsKvEntry.getKey(), kvEntry)); - } - } - return tsKvEntryList; - } - - public static List filterChangedAttr(List currentAttributes, List newAttributes) { - if (currentAttributes == null || currentAttributes.isEmpty()) { - return newAttributes; - } - - Map currentAttrMap = currentAttributes.stream() - .collect(Collectors.toMap(AttributeKvEntry::getKey, Function.identity(), (existing, replacement) -> existing)); - - return newAttributes.stream() - .filter(item -> { - AttributeKvEntry cacheAttr = currentAttrMap.get(item.getKey()); - return cacheAttr == null - || !Objects.equals(item.getValue(), cacheAttr.getValue()) //JSON and String can be equals by value, but different by type - || !Objects.equals(item.getDataType(), cacheAttr.getDataType()); - }) - .collect(Collectors.toList()); - } - -}