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 808c1ea8d8..af4b3e95e5 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 @@ -92,6 +92,7 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.locks.ReentrantLock; import java.util.function.Consumer; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -119,6 +120,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private ListeningExecutorService calculatedFieldExecutor; private ListeningExecutorService calculatedFieldCallbackExecutor; + private final ConcurrentMap entityLocks = new ConcurrentHashMap<>(); + private final ConcurrentMap states = new ConcurrentHashMap<>(); private static final int MAX_LAST_RECORDS_VALUE = 1024; @@ -536,40 +539,47 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, cfId); if (tpi.isMyPartition()) { - CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); - CalculatedFieldEntityCtx calculatedFieldEntityCtx = states.computeIfAbsent(entityCtxId, ctxId -> fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType())); - - Consumer performUpdateState = (state) -> { - if (state.updateState(argumentValues)) { - calculatedFieldEntityCtx.setState(state); - states.put(entityCtxId, calculatedFieldEntityCtx); - rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); - Map arguments = state.getArguments(); - boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && - !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); - if (allArgsPresent) { - performCalculation(calculatedFieldCtx, state, entityId, calculatedFieldIds); + ReentrantLock lock = entityLocks.computeIfAbsent(entityId, id -> new ReentrantLock()); + lock.lock(); + + try { + CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId.getId(), entityId.getId()); + CalculatedFieldEntityCtx calculatedFieldEntityCtx = states.computeIfAbsent(entityCtxId, ctxId -> fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType())); + + Consumer performUpdateState = (state) -> { + if (state.updateState(argumentValues)) { + calculatedFieldEntityCtx.setState(state); + states.put(entityCtxId, calculatedFieldEntityCtx); + rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); + Map arguments = state.getArguments(); + boolean allArgsPresent = arguments.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()) && + !arguments.containsValue(SingleValueArgumentEntry.EMPTY) && !arguments.containsValue(TsRollingArgumentEntry.EMPTY); + if (allArgsPresent) { + performCalculation(calculatedFieldCtx, state, entityId, calculatedFieldIds); + } } - } - }; + }; - CalculatedFieldState state = calculatedFieldEntityCtx.getState(); + CalculatedFieldState state = calculatedFieldEntityCtx.getState(); - boolean allKeysPresent = argumentValues.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); - boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() - .anyMatch(argument -> "TS_ROLLING".equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); + boolean allKeysPresent = argumentValues.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); + boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() + .anyMatch(argument -> "TS_ROLLING".equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); - if (!allKeysPresent || requiresTsRollingUpdate) { + if (!allKeysPresent || requiresTsRollingUpdate) { - Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() - .filter(entry -> !argumentValues.containsKey(entry.getKey()) || ("TS_ROLLING".equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null)) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() + .filter(entry -> !argumentValues.containsKey(entry.getKey()) || ("TS_ROLLING".equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null)) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); - fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentValues::putAll) - .addListener(() -> performUpdateState.accept(state), - calculatedFieldCallbackExecutor); - } else { - performUpdateState.accept(state); + fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentValues::putAll) + .addListener(() -> performUpdateState.accept(state), + calculatedFieldCallbackExecutor); + } else { + performUpdateState.accept(state); + } + } finally { + lock.unlock(); } } else { sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, calculatedFieldIds, argumentValues); @@ -747,7 +757,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return tsRollingArgumentEntry; } else if (entryProto.hasSingleValue()) { TransportProtos.SingleValueProto singleValueProto = entryProto.getSingleValue(); - return new SingleValueArgumentEntry(singleValueProto.getTs(), fromObjectProto(singleValueProto.getValue())); + return new SingleValueArgumentEntry(singleValueProto.getTs(), fromObjectProto(singleValueProto.getValue()), singleValueProto.getVersion()); } else { throw new IllegalArgumentException("Unsupported ArgumentEntryProto type"); } 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 ba7e094f77..dd56405352 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 @@ -22,8 +22,6 @@ import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; import java.util.List; -import java.util.TreeMap; -import java.util.stream.Collectors; @JsonTypeInfo( use = JsonTypeInfo.Id.NAME, @@ -48,8 +46,7 @@ public interface ArgumentEntry { } static ArgumentEntry createTsRollingArgument(List kvEntries) { - return new TsRollingArgumentEntry(kvEntries.stream(). - collect(Collectors.toMap(TsKvEntry::getTs, TsKvEntry::getValue, (oldValue, newValue) -> newValue, TreeMap::new))); + return new TsRollingArgumentEntry(kvEntries); } @JsonIgnore 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 462fa19b17..be7319b930 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,7 +17,6 @@ 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 { @@ -37,34 +36,30 @@ public abstract class BaseCalculatedFieldState implements CalculatedFieldState { 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) { - if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { - arguments.put(key, argumentEntry.copy()); - stateUpdated.set(true); - } - } else if (existingArgumentEntry instanceof TsRollingArgumentEntry existingTsRollingArgumentEntry) { - if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { - existingTsRollingArgumentEntry.addAllTsRecords(tsRollingArgumentEntry.getTsRecords()); - stateUpdated.set(true); - } - } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - if (existingArgumentEntry.hasUpdatedValue(argumentEntry)) { - existingTsRollingArgumentEntry.addTsRecord(singleValueArgumentEntry.getTs(), singleValueArgumentEntry.getValue()); - stateUpdated.set(true); - } - } + + boolean stateUpdated = false; + + for (Map.Entry entry : argumentValues.entrySet()) { + String key = entry.getKey(); + ArgumentEntry newEntry = entry.getValue(); + ArgumentEntry existingEntry = arguments.get(key); + + if (existingEntry == null || existingEntry.hasUpdatedValue(newEntry)) { + if (existingEntry instanceof TsRollingArgumentEntry existingTsRollingEntry && newEntry instanceof TsRollingArgumentEntry newTsRollingEntry) { + existingTsRollingEntry.addAllTsRecords(newTsRollingEntry.getTsRecords()); + } else if (existingEntry instanceof TsRollingArgumentEntry existingTsRollingEntry && newEntry instanceof SingleValueArgumentEntry singleValueEntry) { + existingTsRollingEntry.addTsRecord(singleValueEntry.getTs(), singleValueEntry.getValue()); + } else if (existingEntry instanceof SingleValueArgumentEntry existingSingleValueEntry && newEntry instanceof SingleValueArgumentEntry singleValueEntry + && singleValueEntry.getVersion() > existingSingleValueEntry.getVersion()) { + arguments.put(key, newEntry.copy()); + } else { + arguments.put(key, newEntry.copy()); } - } else { - arguments.put(key, argumentEntry.copy()); - stateUpdated.set(true); + stateUpdated = true; } - }); - return stateUpdated.get(); + } + + return stateUpdated; } } 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 2cd5c68144..d54a3220ed 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 @@ -26,6 +26,8 @@ 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 java.util.ArrayList; +import java.util.List; import java.util.Map; @Data @@ -36,6 +38,7 @@ public class CalculatedFieldCtx { private EntityId entityId; private CalculatedFieldType cfType; private final Map arguments; + private final List argKeys; private Output output; private String expression; private TbelInvokeService tbelInvokeService; @@ -48,10 +51,11 @@ public class CalculatedFieldCtx { this.cfType = calculatedField.getType(); CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); this.arguments = configuration.getArguments(); + this.argKeys = new ArrayList<>(arguments.keySet()); this.output = configuration.getOutput(); this.expression = configuration.getExpression(); this.tbelInvokeService = tbelInvokeService; - if (!CalculatedFieldType.SIMPLE.equals(calculatedField.getType())) { + if (CalculatedFieldType.SCRIPT.equals(calculatedField.getType())) { this.calculatedFieldScriptEngine = initEngine(tenantId, expression, tbelInvokeService); } } @@ -65,7 +69,7 @@ public class CalculatedFieldCtx { tenantId, tbelInvokeService, expression, - arguments.keySet().toArray(new String[0]) + argKeys.toArray(String[]::new) ); } 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 87429050de..de7c514786 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 @@ -49,7 +49,9 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState { tsRecords.entrySet().removeIf(tsRecord -> tsRecord.getKey() < System.currentTimeMillis() - argument.getTimeWindow()); } }); - Object[] args = arguments.values().stream().map(ArgumentEntry::getValue).toArray(); + Object[] args = ctx.getArgKeys().stream() + .map(key -> arguments.get(key).getValue()) + .toArray(); ListenableFuture> resultFuture = ctx.getCalculatedFieldScriptEngine().executeToMapAsync(args); Output output = ctx.getOutput(); return Futures.transform(resultFuture, 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 3f4fd5bdce..20b531f562 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 @@ -32,11 +32,15 @@ public class SingleValueArgumentEntry implements ArgumentEntry { private long ts; private Object value; + private long version; + public SingleValueArgumentEntry(KvEntry entry) { - if (entry instanceof TsKvEntry) { - this.ts = ((TsKvEntry) entry).getTs(); - } else if (entry instanceof AttributeKvEntry) { - this.ts = ((AttributeKvEntry) entry).getLastUpdateTs(); + if (entry instanceof TsKvEntry tsKvEntry) { + this.ts = tsKvEntry.getTs(); + this.version = tsKvEntry.getVersion(); + } else if (entry instanceof AttributeKvEntry attributeKvEntry) { + this.ts = attributeKvEntry.getLastUpdateTs(); + this.version = attributeKvEntry.getVersion(); } this.value = entry.getValue(); } @@ -66,7 +70,7 @@ public class SingleValueArgumentEntry implements ArgumentEntry { @Override public ArgumentEntry copy() { - return new SingleValueArgumentEntry(this.ts, this.value); + return new SingleValueArgumentEntry(this.ts, this.value, this.version); } } 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 4ffa391550..64de2f8c8a 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 @@ -16,16 +16,20 @@ package org.thingsboard.server.service.cf.ctx.state; import com.fasterxml.jackson.annotation.JsonIgnore; +import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.math.NumberUtils; +import org.thingsboard.server.common.data.kv.TsKvEntry; +import java.util.List; import java.util.Map; import java.util.TreeMap; @Data @NoArgsConstructor +@AllArgsConstructor @Slf4j public class TsRollingArgumentEntry implements ArgumentEntry { @@ -35,8 +39,8 @@ public class TsRollingArgumentEntry implements ArgumentEntry { private TreeMap tsRecords = new TreeMap<>(); - public TsRollingArgumentEntry(TreeMap tsRecords) { - addAllTsRecords(tsRecords); + public TsRollingArgumentEntry(List kvEntries) { + kvEntries.forEach(tsKvEntry -> addTsRecord(tsKvEntry.getTs(), tsKvEntry.getValue())); } /** diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 581f9eb9f2..1f66621a5b 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -841,6 +841,7 @@ message TsRollingProto { message SingleValueProto { int64 ts = 1; ObjectProto value = 2; + int64 version = 3; } message ObjectProto {