diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java index 7394c95f08..bf5dc8d42f 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java @@ -36,8 +36,6 @@ public interface CalculatedFieldCache { List getCalculatedFieldLinksByEntityId(EntityId entityId); - void updateCalculatedFieldLinks(CalculatedFieldId calculatedFieldId); - CalculatedFieldCtx getCalculatedFieldCtx(CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService); Set getEntitiesByProfile(TenantId tenantId, EntityId entityId); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java index 6d1d459b9b..8ba1f6dfed 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java @@ -27,8 +27,6 @@ public interface CalculatedFieldExecutionService { void onTelemetryUpdateMsg(TransportProtos.TelemetryUpdateMsgProto proto); - void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback); - void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback); void onProfileEntityMsg(TransportProtos.ProfileEntityMsgProto proto, TbCallback callback); diff --git a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java index 71fc3eff67..f762ae3530 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java @@ -36,12 +36,12 @@ import org.thingsboard.server.dao.cf.CalculatedFieldService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; -import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -72,14 +72,14 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { PageDataIterable cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf)); calculatedFields.values().forEach(cf -> - entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new ArrayList<>()).add(cf) + entityIdCalculatedFields.computeIfAbsent(cf.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(cf) ); PageDataIterable cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); - cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); + cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new CopyOnWriteArrayList<>()).add(link)); calculatedFieldLinks.values().stream() .flatMap(List::stream) .forEach(link -> - entityIdCalculatedFieldLinks.computeIfAbsent(link.getEntityId(), id -> new ArrayList<>()).add(link) + entityIdCalculatedFieldLinks.computeIfAbsent(link.getEntityId(), id -> new CopyOnWriteArrayList<>()).add(link) ); } @@ -90,41 +90,17 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { @Override public List getCalculatedFieldsByEntityId(EntityId entityId) { - return entityIdCalculatedFields.getOrDefault(entityId, new ArrayList<>()); + return entityIdCalculatedFields.getOrDefault(entityId, new CopyOnWriteArrayList<>()); } @Override public List getCalculatedFieldLinks(CalculatedFieldId calculatedFieldId) { - return calculatedFieldLinks.getOrDefault(calculatedFieldId, new ArrayList<>()); + return calculatedFieldLinks.getOrDefault(calculatedFieldId, new CopyOnWriteArrayList<>()); } @Override public List getCalculatedFieldLinksByEntityId(EntityId entityId) { - return entityIdCalculatedFieldLinks.getOrDefault(entityId, new ArrayList<>()); - } - - @Override - public void updateCalculatedFieldLinks(CalculatedFieldId calculatedFieldId) { - log.debug("Update calculated field links per entity for calculated field: [{}]", calculatedFieldId); - calculatedFieldFetchLock.lock(); - try { - List cfLinks = getCalculatedFieldLinks(calculatedFieldId); - if (cfLinks != null && !cfLinks.isEmpty()) { - cfLinks.forEach(link -> { - entityIdCalculatedFieldLinks.compute(link.getEntityId(), (id, existingList) -> { - if (existingList == null) { - existingList = new ArrayList<>(); - } else if (!(existingList instanceof ArrayList)) { - existingList = new ArrayList<>(existingList); - } - existingList.add(link); - return existingList; - }); - }); - } - } finally { - calculatedFieldFetchLock.unlock(); - } + return entityIdCalculatedFieldLinks.getOrDefault(entityId, new CopyOnWriteArrayList<>()); } @Override @@ -192,7 +168,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { calculatedFields.put(calculatedFieldId, calculatedField); - entityIdCalculatedFields.computeIfAbsent(cfEntityId, entityId -> new ArrayList<>()).add(calculatedField); + entityIdCalculatedFields.computeIfAbsent(cfEntityId, entityId -> new CopyOnWriteArrayList<>()).add(calculatedField); CalculatedFieldConfiguration configuration = calculatedField.getConfiguration(); calculatedFieldLinks.put(calculatedFieldId, configuration.buildCalculatedFieldLinks(tenantId, cfEntityId, calculatedFieldId)); @@ -200,7 +176,7 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { configuration.getReferencedEntities().stream() .filter(referencedEntityId -> !referencedEntityId.equals(cfEntityId)) .forEach(referencedEntityId -> { - entityIdCalculatedFieldLinks.computeIfAbsent(referencedEntityId, entityId -> new ArrayList<>()) + entityIdCalculatedFieldLinks.computeIfAbsent(referencedEntityId, entityId -> new CopyOnWriteArrayList<>()) .add(configuration.buildCalculatedFieldLink(tenantId, referencedEntityId, calculatedFieldId)); }); } finally { @@ -210,13 +186,8 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache { @Override public void updateCalculatedField(TenantId tenantId, CalculatedFieldId calculatedFieldId) { - calculatedFieldFetchLock.lock(); - try { - evict(calculatedFieldId); - addCalculatedField(tenantId, calculatedFieldId); - } finally { - calculatedFieldFetchLock.unlock(); - } + evict(calculatedFieldId); + addCalculatedField(tenantId, calculatedFieldId); } @Override 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 dd11c799c2..e54c1bd5c9 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 @@ -77,6 +77,7 @@ import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.gen.transport.TransportProtos; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtx; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; +import org.thingsboard.server.service.cf.ctx.CalculatedFieldStateService; import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; @@ -107,8 +108,6 @@ import java.util.function.Consumer; import java.util.stream.Collectors; import static org.thingsboard.server.common.data.DataConstants.SCOPE; -import static org.thingsboard.server.common.util.ProtoUtils.fromObjectProto; -import static org.thingsboard.server.common.util.ProtoUtils.toObjectProto; import static org.thingsboard.server.common.util.ProtoUtils.toTsKvProto; @Service @@ -122,7 +121,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private final CalculatedFieldCache calculatedFieldCache; private final AttributesService attributesService; private final TimeseriesService timeseriesService; - private final RocksDBService rocksDBService; + private final CalculatedFieldStateService stateService; private final TbClusterService clusterService; private final TbelInvokeService tbelInvokeService; @@ -148,8 +147,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field")); calculatedFieldCallbackExecutor = MoreExecutors.listeningDecorator(ThingsBoardExecutors.newWorkStealingPool( Math.max(4, Runtime.getRuntime().availableProcessors()), "calculated-field-callback")); - scheduledExecutor.submit(() -> rocksDBService.getAll() - .forEach((ctxId, ctx) -> states.put(JacksonUtil.fromString(ctxId, CalculatedFieldEntityCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldEntityCtx.class)))); + scheduledExecutor.submit(() -> states.putAll(stateService.restoreStates())); } @PreDestroy @@ -223,10 +221,9 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private void restoreState(CalculatedFieldId calculatedFieldId, EntityId entityId) { CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); - String storedState = rocksDBService.get(JacksonUtil.writeValueAsString(ctxId)); + CalculatedFieldEntityCtx restoredCtx = stateService.restoreState(ctxId); - if (storedState != null) { - CalculatedFieldEntityCtx restoredCtx = JacksonUtil.fromString(storedState, CalculatedFieldEntityCtx.class); + if (restoredCtx != null) { states.put(ctxId, restoredCtx); log.info("Restored state for CalculatedField [{}]", calculatedFieldId); } else { @@ -251,18 +248,21 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas log.info("Received CalculatedFieldMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); if (proto.getDeleted()) { log.warn("Executing onCalculatedFieldDelete, calculatedFieldId=[{}]", calculatedFieldId); + calculatedFieldCache.evict(calculatedFieldId); onCalculatedFieldDelete(calculatedFieldId, callback); callback.onSuccess(); } - CalculatedField cf = calculatedFieldCache.getCalculatedField(calculatedFieldId); + CalculatedField cf = calculatedFieldService.findById(tenantId, calculatedFieldId); if (proto.getUpdated()) { log.info("Executing onCalculatedFieldUpdate, calculatedFieldId=[{}]", calculatedFieldId); + calculatedFieldCache.updateCalculatedField(tenantId, calculatedFieldId); boolean shouldReinit = onCalculatedFieldUpdate(cf, callback); if (!shouldReinit) { return; } } if (cf != null) { + calculatedFieldCache.addCalculatedField(tenantId, calculatedFieldId); EntityId entityId = cf.getEntityId(); CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); switch (entityId.getEntityType()) { @@ -312,12 +312,13 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas private void onCalculatedFieldDelete(CalculatedFieldId calculatedFieldId, TbCallback callback) { try { cleanupEntity(calculatedFieldId); - states.keySet().removeIf(ctxId -> ctxId.cfId().equals(calculatedFieldId)); - List statesToRemove = states.keySet().stream() - .filter(ctxId -> ctxId.cfId().equals(calculatedFieldId)) - .map(JacksonUtil::writeValueAsString) - .toList(); - rocksDBService.deleteAll(statesToRemove); + states.keySet().removeIf(ctxId -> { + if (ctxId.cfId().equals(calculatedFieldId)) { + stateService.removeState(ctxId); + return true; + } + return false; + }); } catch (Exception e) { log.trace("Failed to delete calculated field: [{}]", calculatedFieldId, e); callback.onFailure(e); @@ -366,7 +367,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas TransportProtos.TelemetryUpdateMsgProto telemetryUpdateMsgProto = buildTelemetryUpdateMsgProto(calculatedFieldTelemetryUpdateRequest); clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder() .setCfTelemetryUpdateMsg(telemetryUpdateMsgProto).build(), null); - // Forward this request to a correct server based on entity id. } } } catch (Exception e) { @@ -386,20 +386,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } - private void updateTelemetryForLinkedEntity(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldLink link, Map> tpiStates) { - TenantId tenantId = request.getTenantId(); - EntityId entityId = request.getEntityId(); - CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); - - TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, targetEntity); - if (targetEntityTpi.isMyPartition()) { - mapAndProcessUpdatedTelemetry(tenantId, entityId, calculatedFieldId, request, link.getConfiguration()); - } else { - List ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); - ctxIds.add(new CalculatedFieldEntityCtxId(calculatedFieldId, targetEntity)); - } - } - private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map> tpiStates) { TenantId tenantId = request.getTenantId(); EntityId entityId = request.getEntityId(); @@ -411,14 +397,28 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas if (isProfileEntity(targetEntityId)) { calculatedFieldCache.getEntitiesByProfile(tenantId, targetEntityId).forEach(entityByProfile -> { - updateTelemetryForLinkedEntity(request, entityByProfile, link, tpiStates); + processCalculatedFieldLink(request, entityByProfile, link, tpiStates); }); } else { - updateTelemetryForLinkedEntity(request, targetEntityId, link, tpiStates); + processCalculatedFieldLink(request, targetEntityId, link, tpiStates); } }); } + private void processCalculatedFieldLink(CalculatedFieldTelemetryUpdateRequest request, EntityId targetEntity, CalculatedFieldLink link, Map> tpiStates) { + TenantId tenantId = request.getTenantId(); + EntityId entityId = request.getEntityId(); + CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); + + TopicPartitionInfo targetEntityTpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, targetEntity); + if (targetEntityTpi.isMyPartition()) { + mapAndProcessUpdatedTelemetry(tenantId, entityId, calculatedFieldId, request, link.getConfiguration()); + } else { + List ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>()); + ctxIds.add(new CalculatedFieldEntityCtxId(calculatedFieldId, targetEntity)); + } + } + private void mapAndProcessUpdatedTelemetry(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, @@ -490,31 +490,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } - @Override - public void onCalculatedFieldStateMsg(TransportProtos.CalculatedFieldStateMsgProto proto, TbCallback callback) { - try { - TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); - CalculatedFieldId calculatedFieldId = new CalculatedFieldId(new UUID(proto.getCalculatedFieldIdMSB(), proto.getCalculatedFieldIdLSB())); - EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); - log.info("Received CalculatedFieldStateMsgProto for processing: tenantId=[{}], calculatedFieldId=[{}], entityId=[{}]", tenantId, calculatedFieldId, entityId); - if (proto.getClear()) { - clearState(tenantId, calculatedFieldId, entityId); - return; - } - - List previousCalculatedFieldIds = proto.getPreviousCalculatedFieldsList().stream() - .map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB()))) - .collect(Collectors.toCollection(ArrayList::new)); - Map argumentsMap = proto.getArgumentsMap().entrySet().stream() - .collect(Collectors.toMap(Map.Entry::getKey, entry -> fromArgumentEntryProto(entry.getValue()))); - - CalculatedFieldCtx calculatedFieldCtx = calculatedFieldCache.getCalculatedFieldCtx(calculatedFieldId, tbelInvokeService); - updateOrInitializeState(calculatedFieldCtx, entityId, argumentsMap, previousCalculatedFieldIds); - } catch (Exception e) { - log.trace("Failed to process calculated field update state msg: [{}]", proto, e); - } - } - @Override public void onEntityProfileChangedMsg(TransportProtos.EntityProfileUpdateMsgProto proto, TbCallback callback) { try { @@ -522,12 +497,15 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); EntityId oldProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getOldProfileIdMSB(), proto.getOldProfileIdLSB())); EntityId newProfileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getNewProfileIdMSB(), proto.getNewProfileIdLSB())); - log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); - - calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, oldProfileId) - .forEach(cfId -> clearState(tenantId, cfId, entityId)); - initializeStateForEntityByProfile(tenantId, entityId, newProfileId, callback); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); + if (tpi.isMyPartition()) { + log.info("Received EntityProfileUpdateMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); + calculatedFieldCache.getCalculatedFieldsByEntityId(oldProfileId).forEach(cf -> clearState(cf.getId(), entityId)); + initializeStateForEntityByProfile(entityId, newProfileId, callback); + } else { + clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setEntityProfileUpdateMsg(proto).build(), null); + } } catch (Exception e) { log.trace("Failed to process entity type update msg: [{}]", proto, e); } @@ -539,37 +517,39 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas TenantId tenantId = TenantId.fromUUID(new UUID(proto.getTenantIdMSB(), proto.getTenantIdLSB())); EntityId entityId = EntityIdFactory.getByTypeAndUuid(proto.getEntityType(), new UUID(proto.getEntityIdMSB(), proto.getEntityIdLSB())); EntityId profileId = EntityIdFactory.getByTypeAndUuid(proto.getEntityProfileType(), new UUID(proto.getProfileIdMSB(), proto.getProfileIdLSB())); - log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); - if (proto.getDeleted()) { - log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - getCalculatedFieldLinks(entityId, profileId) - .forEach(link -> clearState(tenantId, link.getCalculatedFieldId(), entityId)); + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); + if (tpi.isMyPartition()) { + log.info("Received ProfileEntityMsgProto for processing: tenantId=[{}], entityId=[{}]", tenantId, entityId); + if (proto.getDeleted()) { + log.info("Executing profile entity deleted msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); + + calculatedFieldCache.getCalculatedFieldsByEntityId(entityId).forEach(cf -> clearState(cf.getId(), entityId)); + calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).forEach(cf -> clearState(cf.getId(), entityId)); + } else { + log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); + initializeStateForEntityByProfile(entityId, profileId, callback); + } } else { - log.info("Executing profile entity added msg, tenantId=[{}], entityId=[{}]", tenantId, entityId); - initializeStateForEntityByProfile(tenantId, entityId, profileId, callback); + clusterService.pushMsgToRuleEngine(tpi, UUID.randomUUID(), TransportProtos.ToRuleEngineMsg.newBuilder().setProfileEntityMsg(proto).build(), null); } + + } catch (Exception e) { log.trace("Failed to process profile entity msg: [{}]", proto, e); } } - private void clearState(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); - if (tpi.isMyPartition()) { - log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); - CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); - states.remove(ctxId); - rocksDBService.delete(JacksonUtil.writeValueAsString(ctxId)); - } else { - sendClearCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); - } + private void clearState(CalculatedFieldId calculatedFieldId, EntityId entityId) { + log.warn("Executing clearState, calculatedFieldId=[{}], entityId=[{}]", calculatedFieldId, entityId); + CalculatedFieldEntityCtxId ctxId = new CalculatedFieldEntityCtxId(calculatedFieldId, entityId); + states.remove(ctxId); + stateService.removeState(ctxId); } - private void initializeStateForEntityByProfile(TenantId tenantId, EntityId entityId, EntityId profileId, TbCallback callback) { - calculatedFieldService.findCalculatedFieldIdsByEntityId(tenantId, profileId) - .stream() - .map(cfId -> calculatedFieldCache.getCalculatedFieldCtx(cfId, tbelInvokeService)) + private void initializeStateForEntityByProfile(EntityId entityId, EntityId profileId, TbCallback callback) { + calculatedFieldCache.getCalculatedFieldsByEntityId(profileId).stream() + .map(cf -> calculatedFieldCache.getCalculatedFieldCtx(cf.getId(), tbelInvokeService)) .forEach(cfCtx -> initializeStateForEntity(cfCtx, entityId, callback)); } @@ -607,65 +587,58 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private void updateOrInitializeState(CalculatedFieldCtx calculatedFieldCtx, EntityId entityId, Map argumentValues, List previousCalculatedFieldIds) { - TenantId tenantId = calculatedFieldCtx.getTenantId(); CalculatedFieldId cfId = calculatedFieldCtx.getCfId(); Map argumentsMap = new HashMap<>(argumentValues); - TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_RULE_ENGINE, tenantId, entityId); - if (tpi.isMyPartition()) { - - CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId, entityId); + CalculatedFieldEntityCtxId entityCtxId = new CalculatedFieldEntityCtxId(cfId, entityId); - states.compute(entityCtxId, (ctxId, ctx) -> { - CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); + states.compute(entityCtxId, (ctxId, ctx) -> { + CalculatedFieldEntityCtx calculatedFieldEntityCtx = ctx != null ? ctx : fetchCalculatedFieldEntityState(ctxId, calculatedFieldCtx.getCfType()); - CompletableFuture updateFuture = new CompletableFuture<>(); + CompletableFuture updateFuture = new CompletableFuture<>(); - Consumer performUpdateState = (state) -> { - if (state.updateState(argumentsMap)) { - calculatedFieldEntityCtx.setState(state); - 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, previousCalculatedFieldIds); - } - log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); + Consumer performUpdateState = (state) -> { + if (state.updateState(argumentsMap)) { + calculatedFieldEntityCtx.setState(state); + stateService.persistState(entityCtxId, 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, previousCalculatedFieldIds); } - updateFuture.complete(null); - }; + log.info("Successfully updated state: calculatedFieldId=[{}], entityId=[{}]", calculatedFieldCtx.getCfId(), entityId); + } + updateFuture.complete(null); + }; - CalculatedFieldState state = calculatedFieldEntityCtx.getState(); + CalculatedFieldState state = calculatedFieldEntityCtx.getState(); - boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); - boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() - .anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); + boolean allKeysPresent = argumentsMap.keySet().containsAll(calculatedFieldCtx.getArguments().keySet()); + boolean requiresTsRollingUpdate = calculatedFieldCtx.getArguments().values().stream() + .anyMatch(argument -> ArgumentType.TS_ROLLING.equals(argument.getType()) && state.getArguments().get(argument.getKey()) == null); - if (!allKeysPresent || requiresTsRollingUpdate) { - Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() - .filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.TS_ROLLING.equals(entry.getValue().getType()) && state.getArguments().get(entry.getKey()) == null)) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + if (!allKeysPresent || requiresTsRollingUpdate) { + Map missingArguments = calculatedFieldCtx.getArguments().entrySet().stream() + .filter(entry -> !argumentsMap.containsKey(entry.getKey()) || (ArgumentType.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, argumentsMap::putAll) - .addListener(() -> performUpdateState.accept(state), - calculatedFieldCallbackExecutor); - } else { - performUpdateState.accept(state); - } + fetchArguments(calculatedFieldCtx.getTenantId(), entityId, missingArguments, argumentsMap::putAll) + .addListener(() -> performUpdateState.accept(state), + calculatedFieldCallbackExecutor); + } else { + performUpdateState.accept(state); + } - try { - updateFuture.join(); - } catch (Exception e) { - log.trace("Failed to update state for ctxId [{}].", ctxId, e); - throw new RuntimeException("Failed to update or initialize state.", e); - } + try { + updateFuture.join(); + } catch (Exception e) { + log.trace("Failed to update state for ctxId [{}].", ctxId, e); + throw new RuntimeException("Failed to update or initialize state.", e); + } - return calculatedFieldEntityCtx; - }); - } else { - sendUpdateCalculatedFieldStateMsg(tenantId, cfId, entityId, previousCalculatedFieldIds, argumentsMap); - } + return calculatedFieldEntityCtx; + }); } private void performCalculation(CalculatedFieldCtx calculatedFieldCtx, CalculatedFieldState state, EntityId entityId, List previousCalculatedFieldIds) { @@ -706,14 +679,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } } - private List getCalculatedFieldLinks(EntityId entityId, EntityId profileId) { - List links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId)); - if (profileId != null) { - links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(profileId)); - } - return links; - } - private ListenableFuture fetchArguments(TenantId tenantId, EntityId entityId, Map necessaryArguments, Consumer> onComplete) { Map argumentValues = new HashMap<>(); List> futures = new ArrayList<>(); @@ -777,89 +742,10 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas return Futures.transform(tsRollingFuture, tsRolling -> tsRolling == null ? TsRollingArgumentEntry.EMPTY : ArgumentEntry.createTsRollingArgument(tsRolling), calculatedFieldCallbackExecutor); } - private void sendUpdateCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId, List previousCalculatedFieldIds, Map argumentValues) { - TransportProtos.CalculatedFieldStateMsgProto.Builder msgBuilder = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId); - if (argumentValues != null) { - argumentValues.forEach((key, argumentEntry) -> msgBuilder.putArguments(key, toArgumentEntryProto(argumentEntry))); - } - if (previousCalculatedFieldIds != null) { - previousCalculatedFieldIds.forEach(cfId -> msgBuilder.addPreviousCalculatedFields( - TransportProtos.CalculatedFieldIdProto.newBuilder() - .setCalculatedFieldIdMSB(cfId.getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(cfId.getId().getLeastSignificantBits()) - .build() - )); - } - - log.info("Sending calculated field state msg from entityId [{}]", entityId); - clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msgBuilder).build(), null); - } - - private void sendClearCalculatedFieldStateMsg(TenantId tenantId, CalculatedFieldId calculatedFieldId, EntityId entityId) { - TransportProtos.CalculatedFieldStateMsgProto msg = createBaseCalculatedFieldStateMsg(tenantId, calculatedFieldId, entityId) - .setClear(true) - .build(); - - clusterService.pushMsgToCore(tenantId, calculatedFieldId, TransportProtos.ToCoreMsg.newBuilder().setCalculatedFieldStateMsg(msg).build(), null); - } - - private TransportProtos.CalculatedFieldStateMsgProto.Builder createBaseCalculatedFieldStateMsg( - TenantId tenantId, - CalculatedFieldId calculatedFieldId, - EntityId entityId - ) { - return TransportProtos.CalculatedFieldStateMsgProto.newBuilder() - .setTenantIdMSB(tenantId.getId().getMostSignificantBits()) - .setTenantIdLSB(tenantId.getId().getLeastSignificantBits()) - .setCalculatedFieldIdMSB(calculatedFieldId.getId().getMostSignificantBits()) - .setCalculatedFieldIdLSB(calculatedFieldId.getId().getLeastSignificantBits()) - .setEntityType(entityId.getEntityType().name()) - .setEntityIdMSB(entityId.getId().getMostSignificantBits()) - .setEntityIdLSB(entityId.getId().getLeastSignificantBits()); - } - - private TransportProtos.ArgumentEntryProto toArgumentEntryProto(ArgumentEntry argumentEntry) { - TransportProtos.ArgumentEntryProto.Builder argumentProtoBuilder = TransportProtos.ArgumentEntryProto.newBuilder(); - - if (argumentEntry instanceof TsRollingArgumentEntry tsRollingArgumentEntry) { - TransportProtos.TsRollingProto.Builder tsRollingProtoBuilder = TransportProtos.TsRollingProto.newBuilder(); - tsRollingArgumentEntry.getTsRecords().forEach((ts, value) -> - tsRollingProtoBuilder.putTsRecords(ts, toObjectProto(value)) - ); - argumentProtoBuilder.setTsRecords(tsRollingProtoBuilder.build()); - } else if (argumentEntry instanceof SingleValueArgumentEntry singleValueArgumentEntry) { - argumentProtoBuilder.setSingleValue( - TransportProtos.SingleValueProto.newBuilder() - .setTs(singleValueArgumentEntry.getTs()) - .setValue(toObjectProto(singleValueArgumentEntry.getValue())) - .build() - ); - } - - return argumentProtoBuilder.build(); - } - - private ArgumentEntry fromArgumentEntryProto(TransportProtos.ArgumentEntryProto entryProto) { - if (entryProto.hasTsRecords()) { - TsRollingArgumentEntry tsRollingArgumentEntry = new TsRollingArgumentEntry(); - entryProto.getTsRecords().getTsRecordsMap().forEach((ts, objectProto) -> - tsRollingArgumentEntry.getTsRecords().put(ts, fromObjectProto(objectProto)) - ); - return tsRollingArgumentEntry; - } else if (entryProto.hasSingleValue()) { - TransportProtos.SingleValueProto singleValueProto = entryProto.getSingleValue(); - return new SingleValueArgumentEntry(singleValueProto.getTs(), fromObjectProto(singleValueProto.getValue()), singleValueProto.getVersion()); - } else { - throw new IllegalArgumentException("Unsupported ArgumentEntryProto type"); - } - } - private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto(CalculatedFieldTelemetryUpdateRequest request) { return buildTelemetryUpdateMsgProto(request, Collections.emptyList()); } - ; - private TransportProtos.TelemetryUpdateMsgProto buildTelemetryUpdateMsgProto( CalculatedFieldTelemetryUpdateRequest request, List links ) { @@ -952,11 +838,12 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas } private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId, CalculatedFieldType cfType) { - String stateStr = rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)); - if (stateStr == null) { + CalculatedFieldEntityCtx state = stateService.restoreState(entityCtxId); + + if (state == null) { return new CalculatedFieldEntityCtx(entityCtxId, createStateByType(cfType)); } - return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); + return state; } private ObjectNode createJsonPayload(CalculatedFieldResult calculatedFieldResult) { diff --git a/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java index d6b2980042..3aed65eced 100644 --- a/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java +++ b/application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java @@ -19,7 +19,6 @@ import lombok.extern.slf4j.Slf4j; import org.rocksdb.RocksDB; import org.rocksdb.RocksDBException; import org.rocksdb.RocksIterator; -import org.rocksdb.WriteBatch; import org.rocksdb.WriteOptions; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.stereotype.Service; @@ -27,7 +26,6 @@ import org.thingsboard.server.utils.RocksDBConfig; import java.nio.charset.StandardCharsets; import java.util.HashMap; -import java.util.List; import java.util.Map; @Service @@ -59,17 +57,6 @@ public class RocksDBService { } } - public void deleteAll(List keys) { - try (WriteBatch batch = new WriteBatch()) { - for (String key : keys) { - batch.delete(key.getBytes(StandardCharsets.UTF_8)); - } - db.write(writeOptions, batch); - } catch (RocksDBException e) { - log.error("Failed to delete data from RocksDB", e); - } - } - public String get(String key) { try { byte[] value = db.get(key.getBytes(StandardCharsets.UTF_8)); diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java index 6baa75b3ef..c598540ff2 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java @@ -326,8 +326,6 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService future = calculatedFieldsExecutor.submit(() -> calculatedFieldExecutionService.onCalculatedFieldStateMsg(calculatedFieldStateMsgProto, callback)); - DonAsynchron.withCallback(future, - __ -> callback.onSuccess(), - t -> { - log.warn("[{}] Failed to process calculated field state message for entityId [{}]", tenantId.getId(), calculatedFieldId.getId(), t); - callback.onFailure(t); - }); - } - private void forwardToNotificationSchedulerService(TransportProtos.NotificationSchedulerServiceMsg msg, TbCallback callback) { TenantId tenantId = toTenantId(msg.getTenantIdMSB(), msg.getTenantIdLSB()); NotificationRequestId notificationRequestId = new NotificationRequestId(new UUID(msg.getRequestIdMSB(), msg.getRequestIdLSB())); diff --git a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java index d332bac64f..073f47d59b 100644 --- a/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java +++ b/common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java @@ -1316,46 +1316,6 @@ public class ProtoUtils { return builder.build(); } - public static TransportProtos.ObjectProto toObjectProto(Object value) { - if (value == null) { - throw new IllegalArgumentException("Cannot convert null to ObjectProto"); - } - - TransportProtos.ObjectProto.Builder builder = TransportProtos.ObjectProto.newBuilder(); - - if (value instanceof String) { - builder.setStringValue((String) value); - } else if (value instanceof Integer) { - builder.setIntValue((Integer) value); - } else if (value instanceof Long) { - builder.setLongValue((Long) value); - } else if (value instanceof Double) { - builder.setDoubleValue((Double) value); - } else if (value instanceof Boolean) { - builder.setBoolValue((Boolean) value); - } else { - throw new IllegalArgumentException("Unsupported value type: " + value.getClass().getName()); - } - - return builder.build(); - } - - public static Object fromObjectProto(TransportProtos.ObjectProto proto) { - try { - return switch (proto.getValueCase()) { - case STRINGVALUE -> proto.getStringValue(); - case INTVALUE -> proto.getIntValue(); - case LONGVALUE -> proto.getLongValue(); - case DOUBLEVALUE -> proto.getDoubleValue(); - case BOOLVALUE -> proto.getBoolValue(); - case VALUE_NOT_SET -> throw new IllegalArgumentException("Value not set in ObjectProto"); - }; - } catch (Exception e) { - log.error("Failed to deserialize ObjectProto: [{}]", proto, e); - return null; - } - } - private static boolean isNotNull(Object obj) { return obj != null; } diff --git a/common/proto/src/main/proto/queue.proto b/common/proto/src/main/proto/queue.proto index 8a7c2d8c03..1036d5ba67 100644 --- a/common/proto/src/main/proto/queue.proto +++ b/common/proto/src/main/proto/queue.proto @@ -841,51 +841,11 @@ message CalculatedFieldEntityCtxIdProto { int64 entityIdLSB = 5; } -message CalculatedFieldStateMsgProto { - int64 tenantIdMSB = 1; - int64 tenantIdLSB = 2; - int64 calculatedFieldIdMSB = 3; - int64 calculatedFieldIdLSB = 4; - string entityType = 5; - int64 entityIdMSB = 6; - int64 entityIdLSB = 7; - bool clear = 8; - repeated CalculatedFieldIdProto previousCalculatedFields = 9; - map arguments = 10; -} - message CalculatedFieldIdProto { int64 calculatedFieldIdMSB = 1; int64 calculatedFieldIdLSB = 2; } -message ArgumentEntryProto { - oneof entry_type { - TsRollingProto tsRecords = 1; - SingleValueProto singleValue = 2; - } -} - -message TsRollingProto { - map tsRecords = 1; -} - -message SingleValueProto { - int64 ts = 1; - ObjectProto value = 2; - int64 version = 3; -} - -message ObjectProto { - oneof value { - string stringValue = 1; - int32 intValue = 2; - int64 longValue = 3; - double doubleValue = 4; - bool boolValue = 5; - } -} - //Used to report session state to tb-Service and persist this state in the cache on the tb-Service level. message SubscriptionInfoProto { int64 lastActivityTime = 1; @@ -1632,7 +1592,6 @@ message ToCoreMsg { CalculatedFieldMsgProto calculatedFieldMsg = 53; EntityProfileUpdateMsgProto entityProfileUpdateMsg = 54; ProfileEntityMsgProto profileEntityMsg = 55; - CalculatedFieldStateMsgProto calculatedFieldStateMsg = 56; } /* High priority messages with low latency are handled by ThingsBoard Core Service separately */ @@ -1681,6 +1640,8 @@ message ToRuleEngineMsg { repeated string relationTypes = 4; string failureMessage = 5; TelemetryUpdateMsgProto cfTelemetryUpdateMsg = 6; + EntityProfileUpdateMsgProto entityProfileUpdateMsg = 7; + ProfileEntityMsgProto profileEntityMsg = 8; } message ToRuleEngineNotificationMsg {