Browse Source

added calculated field state service

pull/12404/head
IrynaMatveieva 2 years ago
parent
commit
5203ef7422
  1. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldExecutionService.java
  3. 51
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  4. 325
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldExecutionService.java
  5. 13
      application/src/main/java/org/thingsboard/server/service/cf/RocksDBService.java
  6. 14
      application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java
  7. 40
      common/proto/src/main/java/org/thingsboard/server/common/util/ProtoUtils.java
  8. 43
      common/proto/src/main/proto/queue.proto

2
application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldCache.java

@ -36,8 +36,6 @@ public interface CalculatedFieldCache {
List<CalculatedFieldLink> getCalculatedFieldLinksByEntityId(EntityId entityId);
void updateCalculatedFieldLinks(CalculatedFieldId calculatedFieldId);
CalculatedFieldCtx getCalculatedFieldCtx(CalculatedFieldId calculatedFieldId, TbelInvokeService tbelInvokeService);
Set<EntityId> getEntitiesByProfile(TenantId tenantId, EntityId entityId);

2
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);

51
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<CalculatedField> 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<CalculatedFieldLink> 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<CalculatedField> getCalculatedFieldsByEntityId(EntityId entityId) {
return entityIdCalculatedFields.getOrDefault(entityId, new ArrayList<>());
return entityIdCalculatedFields.getOrDefault(entityId, new CopyOnWriteArrayList<>());
}
@Override
public List<CalculatedFieldLink> getCalculatedFieldLinks(CalculatedFieldId calculatedFieldId) {
return calculatedFieldLinks.getOrDefault(calculatedFieldId, new ArrayList<>());
return calculatedFieldLinks.getOrDefault(calculatedFieldId, new CopyOnWriteArrayList<>());
}
@Override
public List<CalculatedFieldLink> 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<CalculatedFieldLink> 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

325
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<String> 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<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> 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<CalculatedFieldEntityCtxId> ctxIds = tpiStates.computeIfAbsent(targetEntityTpi, k -> new ArrayList<>());
ctxIds.add(new CalculatedFieldEntityCtxId(calculatedFieldId, targetEntity));
}
}
private void processCalculatedFieldLinks(CalculatedFieldTelemetryUpdateRequest request, Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> 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<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> 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<CalculatedFieldEntityCtxId> 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<CalculatedFieldId> previousCalculatedFieldIds = proto.getPreviousCalculatedFieldsList().stream()
.map(cfIdProto -> new CalculatedFieldId(new UUID(cfIdProto.getCalculatedFieldIdMSB(), cfIdProto.getCalculatedFieldIdLSB())))
.collect(Collectors.toCollection(ArrayList::new));
Map<String, ArgumentEntry> 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<String, ArgumentEntry> argumentValues, List<CalculatedFieldId> previousCalculatedFieldIds) {
TenantId tenantId = calculatedFieldCtx.getTenantId();
CalculatedFieldId cfId = calculatedFieldCtx.getCfId();
Map<String, ArgumentEntry> 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<Void> updateFuture = new CompletableFuture<>();
CompletableFuture<Void> updateFuture = new CompletableFuture<>();
Consumer<CalculatedFieldState> performUpdateState = (state) -> {
if (state.updateState(argumentsMap)) {
calculatedFieldEntityCtx.setState(state);
rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx));
Map<String, ArgumentEntry> 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<CalculatedFieldState> performUpdateState = (state) -> {
if (state.updateState(argumentsMap)) {
calculatedFieldEntityCtx.setState(state);
stateService.persistState(entityCtxId, calculatedFieldEntityCtx);
Map<String, ArgumentEntry> 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<String, Argument> 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<String, Argument> 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<CalculatedFieldId> previousCalculatedFieldIds) {
@ -706,14 +679,6 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas
}
}
private List<CalculatedFieldLink> getCalculatedFieldLinks(EntityId entityId, EntityId profileId) {
List<CalculatedFieldLink> links = new ArrayList<>(calculatedFieldCache.getCalculatedFieldLinksByEntityId(entityId));
if (profileId != null) {
links.addAll(calculatedFieldCache.getCalculatedFieldLinksByEntityId(profileId));
}
return links;
}
private ListenableFuture<Void> fetchArguments(TenantId tenantId, EntityId entityId, Map<String, Argument> necessaryArguments, Consumer<Map<String, ArgumentEntry>> onComplete) {
Map<String, ArgumentEntry> argumentValues = new HashMap<>();
List<ListenableFuture<ArgumentEntry>> 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<CalculatedFieldId> previousCalculatedFieldIds, Map<String, ArgumentEntry> 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<CalculatedFieldEntityCtxId> 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) {

13
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<String> 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));

14
application/src/main/java/org/thingsboard/server/service/queue/DefaultTbCoreConsumerService.java

@ -326,8 +326,6 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
forwardToCalculatedFieldService(toCoreMsg.getEntityProfileUpdateMsg(), callback);
} else if (toCoreMsg.hasProfileEntityMsg()) {
forwardToCalculatedFieldService(toCoreMsg.getProfileEntityMsg(), callback);
} else if (toCoreMsg.hasCalculatedFieldStateMsg()) {
forwardToCalculatedFieldService(toCoreMsg.getCalculatedFieldStateMsg(), callback);
}
} catch (Throwable e) {
log.warn("[{}] Failed to process message: {}", id, msg, e);
@ -740,18 +738,6 @@ public class DefaultTbCoreConsumerService extends AbstractConsumerService<ToCore
});
}
private void forwardToCalculatedFieldService(TransportProtos.CalculatedFieldStateMsgProto calculatedFieldStateMsgProto, TbCallback callback) {
var tenantId = toTenantId(calculatedFieldStateMsgProto.getTenantIdMSB(), calculatedFieldStateMsgProto.getTenantIdLSB());
var calculatedFieldId = new CalculatedFieldId(new UUID(calculatedFieldStateMsgProto.getCalculatedFieldIdMSB(), calculatedFieldStateMsgProto.getCalculatedFieldIdLSB()));
ListenableFuture<?> 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()));

40
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;
}

43
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<string, ArgumentEntryProto> arguments = 10;
}
message CalculatedFieldIdProto {
int64 calculatedFieldIdMSB = 1;
int64 calculatedFieldIdLSB = 2;
}
message ArgumentEntryProto {
oneof entry_type {
TsRollingProto tsRecords = 1;
SingleValueProto singleValue = 2;
}
}
message TsRollingProto {
map<int64, ObjectProto> 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 {

Loading…
Cancel
Save