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 2147e0a1b7..732d050a22 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 @@ -42,6 +42,7 @@ import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfig 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.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.DeviceProfileId; import org.thingsboard.server.common.data.id.EntityId; @@ -80,6 +81,7 @@ 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 java.util.ArrayList; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -115,6 +117,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private final ConcurrentMap calculatedFieldsCtx = new ConcurrentHashMap<>(); private final ConcurrentMap states = new ConcurrentHashMap<>(); + private final ConcurrentMap> tenantStorage = new ConcurrentHashMap<>(); + private final ConcurrentMap> customerStorage = new ConcurrentHashMap<>(); + private static final int MAX_LAST_RECORDS_VALUE = 1024; @Value("${calculatedField.initFetchPackSize:50000}") @@ -366,6 +371,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private ListenableFuture fetchArgumentValue(CalculatedFieldCtx calculatedFieldCtx, EntityId targetEntityId, Argument argument) { TenantId tenantId = calculatedFieldCtx.getTenantId(); EntityId argumentEntityId = argument.getEntityId(); + if (EntityType.TENANT.equals(argumentEntityId.getEntityType()) || EntityType.CUSTOMER.equals(argumentEntityId.getEntityType())) { + return fetchFromStorage(tenantId, argumentEntityId, argument); + } EntityId entityId = EntityType.DEVICE_PROFILE.equals(argumentEntityId.getEntityType()) || EntityType.ASSET_PROFILE.equals(argumentEntityId.getEntityType()) ? targetEntityId : argumentEntityId; @@ -375,6 +383,21 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return fetchKvEntry(tenantId, entityId, argument); } + private ListenableFuture fetchFromStorage(TenantId tenantId, EntityId entityId, Argument argument) { + List kvEntries; + if (EntityType.TENANT.equals(entityId.getEntityType())) { + kvEntries = tenantStorage.computeIfAbsent(tenantId, id -> new ArrayList<>()); + } else { + kvEntries = customerStorage.computeIfAbsent((CustomerId) entityId, id -> new ArrayList<>()); + } + for (KvEntry kvEntry : kvEntries) { + if (kvEntry.getKey().equals(argument.getKey())) { + return Futures.immediateFuture(ArgumentEntry.createSingleValueArgument(kvEntry)); + } + } + return fetchKvEntry(tenantId, entityId, argument); + } + private ListenableFuture fetchLastRecords(TenantId tenantId, EntityId entityId, Argument argument) { long currentTime = System.currentTimeMillis(); long timeWindow = argument.getTimeWindow() == 0 ? System.currentTimeMillis() : argument.getTimeWindow(); @@ -403,7 +426,26 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas calculatedFieldExecutor); default -> throw new IllegalArgumentException("Invalid argument type '" + argument.getType() + "'."); }; - return Futures.transform(kvEntryFuture, kvEntry -> ArgumentEntry.createSingleValueArgument(kvEntry.orElse(null)), calculatedFieldExecutor); + return Futures.transform(kvEntryFuture, kvEntry -> { + if (EntityType.TENANT.equals(entityId.getEntityType()) || EntityType.CUSTOMER.equals(entityId.getEntityType())) { + updateStorage(tenantId, entityId, kvEntry); + } + return ArgumentEntry.createSingleValueArgument(kvEntry.orElse(null)); + }, calculatedFieldExecutor); + } + + private void updateStorage(TenantId tenantId, EntityId entityId, Optional kvEntry) { + if (kvEntry.isEmpty()) { + return; + } + List kvEntries = switch (entityId.getEntityType()) { + case TENANT -> tenantStorage.computeIfAbsent(tenantId, id -> new ArrayList<>()); + case CUSTOMER -> customerStorage.computeIfAbsent((CustomerId) entityId, id -> new ArrayList<>()); + default -> null; + }; + if (kvEntries != null && !kvEntries.contains(kvEntry.get())) { + kvEntries.add(kvEntry.get()); + } } private KvEntry createDefaultKvEntry(Argument argument) {