From d59de750d0f5b54962436cbf8726a73f60867848 Mon Sep 17 00:00:00 2001 From: IrynaMatveieva Date: Fri, 14 Nov 2025 10:35:08 +0200 Subject: [PATCH] refactoring --- ...tractCalculatedFieldProcessingService.java | 89 ++++++++++--------- 1 file changed, 48 insertions(+), 41 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 6d16e11a17..39ab78f86a 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 @@ -33,10 +33,12 @@ 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.common.adaptor.JsonConverter; +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.OutputStrategy; import org.thingsboard.server.common.data.cf.configuration.OutputType; import org.thingsboard.server.common.data.cf.configuration.RelationPathQueryDynamicSourceConfiguration; import org.thingsboard.server.common.data.cf.configuration.TimeSeriesImmediateOutputStrategy; @@ -357,8 +359,8 @@ public abstract class AbstractCalculatedFieldProcessingService { SettableFuture future = SettableFuture.create(); switch (type) { - case ATTRIBUTES -> saveAttributes(tenantId, entityId, cfResult, cfIds, future); - case TIME_SERIES -> saveTimeSeries(tenantId, entityId, cfResult, cfIds, System.currentTimeMillis(), future); + case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfIds, future); + case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), future); } Futures.addCallback(future, new FutureCallback<>() { @@ -376,39 +378,37 @@ public abstract class AbstractCalculatedFieldProcessingService { }, MoreExecutors.directExecutor()); } - private void saveAttributes(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List cfIds, SettableFuture future) { - if (!(cfResult.getOutputStrategy() instanceof AttributeImmediateOutputStrategy outputStrategy)) { - future.setException(new IllegalArgumentException("Expected AttributeImmediateOutputStrategy")); - return; - } - JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); + private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List cfIds, SettableFuture future) { + if (!(outputStrategy instanceof AttributeImmediateOutputStrategy attOutputStrategy)) { + future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported.")); + } else { + AttributesSaveRequest.Strategy strategy = new Strategy(attOutputStrategy.isSaveAttribute(), attOutputStrategy.isSendWsUpdate(), attOutputStrategy.isProcessCfs()); + List newAttributes = JsonConverter.convertToAttributes(jsonResult); - AttributesSaveRequest.Strategy strategy = new Strategy(outputStrategy.isSaveAttribute(), outputStrategy.isSendWsUpdate(), outputStrategy.isProcessCfs()); - List newAttributes = JsonConverter.convertToAttributes(jsonResult); + if (!attOutputStrategy.isUpdateAttributesOnlyOnValueChange()) { + saveAttributesInternal(tenantId, entityId, scope, cfIds, newAttributes, strategy, future); + return; + } - if (!outputStrategy.isUpdateAttributesOnlyOnValueChange()) { - saveAttributesInternal(tenantId, entityId, cfResult, cfIds, newAttributes, strategy, future); - return; - } + List keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); + ListenableFuture> findFuture = attributesService.find(tenantId, entityId, scope, keys); - List keys = newAttributes.stream().map(KvEntry::getKey).collect(Collectors.toList()); - ListenableFuture> findFuture = attributesService.find(tenantId, entityId, cfResult.getScope(), keys); - - DonAsynchron.withCallback(findFuture, - existingAttributes -> { - List changed = filterChangedAttr(existingAttributes, newAttributes); - if (changed.isEmpty()) { - future.set(null); - return; - } - saveAttributesInternal(tenantId, entityId, cfResult, cfIds, changed, strategy, future); - }, - future::setException, - MoreExecutors.directExecutor()); + DonAsynchron.withCallback(findFuture, + existingAttributes -> { + List changed = filterChangedAttr(existingAttributes, newAttributes); + if (changed.isEmpty()) { + future.set(null); + return; + } + saveAttributesInternal(tenantId, entityId, scope, cfIds, changed, strategy, future); + }, + future::setException, + MoreExecutors.directExecutor()); + } } private void saveAttributesInternal(TenantId tenantId, EntityId entityId, - TelemetryCalculatedFieldResult cfResult, + AttributeScope scope, List cfIds, List entries, AttributesSaveRequest.Strategy strategy, @@ -416,7 +416,7 @@ public abstract class AbstractCalculatedFieldProcessingService { tsSubService.saveAttributesInternal(AttributesSaveRequest.builder() .tenantId(tenantId) .entityId(entityId) - .scope(cfResult.getScope()) + .scope(scope) .entries(entries) .strategy(strategy) .previousCalculatedFieldIds(cfIds) @@ -424,28 +424,35 @@ public abstract class AbstractCalculatedFieldProcessingService { .build()); } - private void saveTimeSeries(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List cfIds, long ts, SettableFuture future) { - if (!(cfResult.getOutputStrategy() instanceof TimeSeriesImmediateOutputStrategy outputStrategy)) { - future.setException(new IllegalArgumentException("Expected TimeSeriesImmediateOutputStrategy")); - return; + private void saveTimeSeries(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, List cfIds, long ts, SettableFuture future) { + if (!(outputStrategy instanceof TimeSeriesImmediateOutputStrategy tsOutputStrategy)) { + future.setException(new IllegalArgumentException("Only TimeSeriesImmediateOutputStrategy is supported.")); + } else { + TimeseriesSaveRequest.Strategy strategy = new TimeseriesSaveRequest.Strategy(tsOutputStrategy.isSaveTimeSeries(), tsOutputStrategy.isSaveLatest(), tsOutputStrategy.isSendWsUpdate(), tsOutputStrategy.isProcessCfs()); + saveTimeSeriesInternal(tenantId, entityId, jsonResult, tsOutputStrategy.getTtl(), cfIds, ts, strategy, future); } - JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); + } + + private void saveTimeSeriesInternal(TenantId tenantId, EntityId entityId, JsonElement jsonResult, Long ttl, List cfIds, long ts, TimeseriesSaveRequest.Strategy strategy, SettableFuture future) { Map> tsKvMap = JsonConverter.convertToTelemetry(jsonResult, ts); if (tsKvMap.isEmpty()) { future.set(null); return; } List tsEntries = toTsKvEntryList(tsKvMap); - TimeseriesSaveRequest.Strategy strategy = new TimeseriesSaveRequest.Strategy(outputStrategy.isSaveTimeSeries(), outputStrategy.isSaveLatest(), outputStrategy.isSendWsUpdate(), outputStrategy.isProcessCfs()); - tsSubService.saveTimeseriesInternal(TimeseriesSaveRequest.builder() + TimeseriesSaveRequest.Builder builder = TimeseriesSaveRequest.builder() .tenantId(tenantId) .entityId(entityId) .entries(tsEntries) - .ttl(outputStrategy.getTtl()) .strategy(strategy) - .previousCalculatedFieldIds(cfIds) - .future(future) - .build()); + .future(future); + if (ttl != null) { + builder.ttl(ttl); + } + if (cfIds != null && !cfIds.isEmpty()) { + builder.previousCalculatedFieldIds(cfIds); + } + tsSubService.saveTimeseriesInternal(builder.build()); } }