From 9ba386e9926b9c180b31038328caf98e338da96e Mon Sep 17 00:00:00 2001 From: dshvaika Date: Wed, 10 Dec 2025 19:43:40 +0200 Subject: [PATCH] refactoring --- ...tractCalculatedFieldProcessingService.java | 24 ++++++++++++++++--- .../service/cf/CalculatedFieldResult.java | 7 ++++++ ...faultCalculatedFieldProcessingService.java | 19 --------------- 3 files changed, 28 insertions(+), 22 deletions(-) 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 13ceb406bb..ec645085e6 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 @@ -22,17 +22,18 @@ import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; import com.google.gson.JsonElement; -import com.google.gson.JsonParser; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.Data; import lombok.extern.slf4j.Slf4j; +import org.apache.commons.lang3.function.TriConsumer; import org.thingsboard.common.util.DonAsynchron; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.rule.engine.api.AttributesSaveRequest; import org.thingsboard.rule.engine.api.AttributesSaveRequest.Strategy; import org.thingsboard.rule.engine.api.TimeseriesSaveRequest; +import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.adaptor.JsonConverter; import org.thingsboard.server.common.data.AttributeScope; @@ -80,7 +81,6 @@ import java.util.Collections; import java.util.HashMap; import java.util.List; import java.util.Map; -import java.util.Objects; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.function.Function; @@ -391,6 +391,24 @@ public abstract class AbstractCalculatedFieldProcessingService { return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE); } + protected void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, + TriConsumer telemetryResultHandler) { + List propagationEntityIds = propagationResult.getEntityIds(); + if (propagationEntityIds.isEmpty()) { + callback.onSuccess(); + return; + } + if (propagationEntityIds.size() == 1) { + EntityId propagationEntityId = propagationEntityIds.get(0); + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), callback); + return; + } + MultipleTbCallback multipleTbCallback = new MultipleTbCallback(propagationEntityIds.size(), callback); + for (var propagationEntityId : propagationEntityIds) { + telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), multipleTbCallback); + } + } + protected void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) { try { clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() { @@ -413,7 +431,7 @@ public abstract class AbstractCalculatedFieldProcessingService { protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, String cfName, TelemetryCalculatedFieldResult cfResult, List cfIds, TbCallback callback) { OutputType type = cfResult.getType(); - JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); + JsonElement jsonResult = cfResult.toJsonElement(); log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult); switch (type) { 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 c973cebc18..8b4c2a0101 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 @@ -15,11 +15,14 @@ */ package org.thingsboard.server.service.cf; +import com.google.gson.JsonElement; +import com.google.gson.JsonParser; import org.thingsboard.server.common.data.id.CalculatedFieldId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.msg.TbMsg; import java.util.List; +import java.util.Objects; public interface CalculatedFieldResult { @@ -29,4 +32,8 @@ public interface CalculatedFieldResult { boolean isEmpty(); + default JsonElement toJsonElement() { + return JsonParser.parseString(Objects.requireNonNull(stringValue())); + } + } 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 271fdb828d..927ba7d1f1 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 @@ -17,7 +17,6 @@ package org.thingsboard.server.service.cf; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.function.TriConsumer; import org.springframework.stereotype.Service; import org.thingsboard.server.actors.calculatedField.CalculatedFieldTelemetryMsg; import org.thingsboard.server.actors.calculatedField.MultipleTbCallback; @@ -169,24 +168,6 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds)); } - private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, - TriConsumer telemetryResultHandler) { - List propagationEntityIds = propagationResult.getEntityIds(); - if (propagationEntityIds.isEmpty()) { - callback.onSuccess(); - return; - } - if (propagationEntityIds.size() == 1) { - EntityId propagationEntityId = propagationEntityIds.get(0); - telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), callback); - return; - } - MultipleTbCallback multipleTbCallback = new MultipleTbCallback(propagationEntityIds.size(), callback); - for (var propagationEntityId : propagationEntityIds) { - telemetryResultHandler.accept(propagationEntityId, propagationResult.getResult(), multipleTbCallback); - } - } - @Override public void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List linkedCalculatedFields, TbCallback callback) { Map> unicasts = new HashMap<>();