|
|
|
@ -36,16 +36,20 @@ import org.thingsboard.script.api.tbel.TbelInvokeService; |
|
|
|
import org.thingsboard.server.cluster.TbClusterService; |
|
|
|
import org.thingsboard.server.common.data.EntityType; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedField; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedFieldLink; |
|
|
|
import org.thingsboard.server.common.data.cf.CalculatedFieldType; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.Argument; |
|
|
|
import org.thingsboard.server.common.data.cf.configuration.CalculatedFieldConfiguration; |
|
|
|
import org.thingsboard.server.common.data.id.AssetId; |
|
|
|
import org.thingsboard.server.common.data.id.AssetProfileId; |
|
|
|
import org.thingsboard.server.common.data.id.CalculatedFieldId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceId; |
|
|
|
import org.thingsboard.server.common.data.id.DeviceProfileId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityId; |
|
|
|
import org.thingsboard.server.common.data.id.EntityIdFactory; |
|
|
|
import org.thingsboard.server.common.data.id.TenantId; |
|
|
|
import org.thingsboard.server.common.data.kv.Aggregation; |
|
|
|
import org.thingsboard.server.common.data.kv.AttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseAttributeKvEntry; |
|
|
|
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; |
|
|
|
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; |
|
|
|
@ -76,6 +80,8 @@ import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldState; |
|
|
|
import org.thingsboard.server.service.cf.ctx.state.ScriptCalculatedFieldState; |
|
|
|
import org.thingsboard.server.service.cf.ctx.state.SimpleCalculatedFieldState; |
|
|
|
import org.thingsboard.server.service.partition.AbstractPartitionBasedService; |
|
|
|
import org.thingsboard.server.service.profile.TbAssetProfileCache; |
|
|
|
import org.thingsboard.server.service.profile.TbDeviceProfileCache; |
|
|
|
|
|
|
|
import java.util.ArrayList; |
|
|
|
import java.util.Collections; |
|
|
|
@ -102,6 +108,8 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private final CalculatedFieldService calculatedFieldService; |
|
|
|
private final AssetService assetService; |
|
|
|
private final DeviceService deviceService; |
|
|
|
private final TbAssetProfileCache assetProfileCache; |
|
|
|
private final TbDeviceProfileCache deviceProfileCache; |
|
|
|
private final AttributesService attributesService; |
|
|
|
private final TimeseriesService timeseriesService; |
|
|
|
private final RocksDBService rocksDBService; |
|
|
|
@ -112,6 +120,7 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
private ListeningExecutorService calculatedFieldCallbackExecutor; |
|
|
|
|
|
|
|
private final ConcurrentMap<CalculatedFieldId, CalculatedField> calculatedFields = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<CalculatedFieldId, List<CalculatedFieldLink>> calculatedFieldLinks = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<CalculatedFieldId, CalculatedFieldCtx> calculatedFieldsCtx = new ConcurrentHashMap<>(); |
|
|
|
private final ConcurrentMap<CalculatedFieldEntityCtxId, CalculatedFieldEntityCtx> states = new ConcurrentHashMap<>(); |
|
|
|
|
|
|
|
@ -130,6 +139,16 @@ 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(this::fetchCalculatedFields); |
|
|
|
} |
|
|
|
|
|
|
|
private void fetchCalculatedFields() { |
|
|
|
PageDataIterable<CalculatedField> cfs = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFields, initFetchPackSize); |
|
|
|
cfs.forEach(cf -> calculatedFields.putIfAbsent(cf.getId(), cf)); |
|
|
|
PageDataIterable<CalculatedFieldLink> cfls = new PageDataIterable<>(calculatedFieldService::findAllCalculatedFieldLinks, initFetchPackSize); |
|
|
|
cfls.forEach(link -> calculatedFieldLinks.computeIfAbsent(link.getCalculatedFieldId(), id -> new ArrayList<>()).add(link)); |
|
|
|
rocksDBService.getAll().forEach((ctxId, ctx) -> states.put(JacksonUtil.fromString(ctxId, CalculatedFieldEntityCtxId.class), JacksonUtil.fromString(ctx, CalculatedFieldEntityCtx.class))); |
|
|
|
states.keySet().removeIf(ctxId -> calculatedFields.keySet().stream().noneMatch(id -> ctxId.cfId().equals(id.getId()))); |
|
|
|
} |
|
|
|
|
|
|
|
@PreDestroy |
|
|
|
@ -216,30 +235,77 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map<String, KvEntry> updatedTelemetry) { |
|
|
|
public void onTelemetryUpdate(TenantId tenantId, EntityId entityId, List<? extends KvEntry> telemetry) { |
|
|
|
try { |
|
|
|
log.info("Received telemetry update msg: tenantId=[{}], calculatedFieldId=[{}]", tenantId, calculatedFieldId); |
|
|
|
CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); |
|
|
|
Map<String, ArgumentEntry> argumentValues = updatedTelemetry.entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); |
|
|
|
|
|
|
|
EntityId cfEntityId = calculatedField.getEntityId(); |
|
|
|
switch (cfEntityId.getEntityType()) { |
|
|
|
case ASSET_PROFILE, DEVICE_PROFILE -> { |
|
|
|
boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); |
|
|
|
if (isCommonEntity) { |
|
|
|
getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); |
|
|
|
} else { |
|
|
|
updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); |
|
|
|
} |
|
|
|
EntityType entityType = entityId.getEntityType(); |
|
|
|
if (EntityType.DEVICE.equals(entityType) || EntityType.ASSET.equals(entityType) || EntityType.CUSTOMER.equals(entityType) || EntityType.TENANT.equals(entityType)) { |
|
|
|
EntityId profileId = null; |
|
|
|
if (EntityType.ASSET.equals(entityType)) { |
|
|
|
profileId = assetProfileCache.get(tenantId, (AssetId) entityId).getId(); |
|
|
|
} else if (EntityType.DEVICE.equals(entityType)) { |
|
|
|
profileId = deviceProfileCache.get(tenantId, (DeviceId) entityId).getId(); |
|
|
|
} |
|
|
|
default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); |
|
|
|
List<CalculatedFieldLink> cfLinks = new ArrayList<>(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, entityId)); |
|
|
|
Optional.ofNullable(profileId).ifPresent(id -> cfLinks.addAll(calculatedFieldService.findAllCalculatedFieldLinksByEntityId(tenantId, id))); |
|
|
|
cfLinks.forEach(link -> { |
|
|
|
CalculatedFieldId calculatedFieldId = link.getCalculatedFieldId(); |
|
|
|
Map<String, String> attributes = link.getConfiguration().getAttributes(); |
|
|
|
Map<String, String> timeSeries = link.getConfiguration().getTimeSeries(); |
|
|
|
Map<String, KvEntry> updatedTelemetry = telemetry.stream() |
|
|
|
.filter(entry -> attributes.containsValue(entry.getKey()) || timeSeries.containsValue(entry.getKey())) |
|
|
|
.collect(Collectors.toMap( |
|
|
|
entry -> getMappedKey(entry, attributes, timeSeries), |
|
|
|
entry -> entry, |
|
|
|
(v1, v2) -> v1 |
|
|
|
)); |
|
|
|
|
|
|
|
if (!updatedTelemetry.isEmpty()) { |
|
|
|
executeTelemetryUpdate(tenantId, entityId, calculatedFieldId, updatedTelemetry); |
|
|
|
} |
|
|
|
}); |
|
|
|
} |
|
|
|
log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); |
|
|
|
} catch (Exception e) { |
|
|
|
log.trace("Failed to update telemetry for calculatedFieldId: [{}]", calculatedFieldId, e); |
|
|
|
log.trace("Failed to update telemetry entityId: [{}]", entityId, e); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private void executeTelemetryUpdate(TenantId tenantId, EntityId entityId, CalculatedFieldId calculatedFieldId, Map<String, KvEntry> updatedTelemetry) { |
|
|
|
log.info("Received telemetry update msg: tenantId=[{}], entityId=[{}], calculatedFieldId=[{}]", tenantId, entityId, calculatedFieldId); |
|
|
|
CalculatedField calculatedField = getOrFetchFromDb(tenantId, calculatedFieldId); |
|
|
|
CalculatedFieldCtx calculatedFieldCtx = calculatedFieldsCtx.computeIfAbsent(calculatedFieldId, id -> new CalculatedFieldCtx(calculatedField, tbelInvokeService)); |
|
|
|
Map<String, ArgumentEntry> argumentValues = updatedTelemetry.entrySet().stream() |
|
|
|
.collect(Collectors.toMap(Map.Entry::getKey, entry -> ArgumentEntry.createSingleValueArgument(entry.getValue()))); |
|
|
|
|
|
|
|
EntityId cfEntityId = calculatedField.getEntityId(); |
|
|
|
switch (cfEntityId.getEntityType()) { |
|
|
|
case ASSET_PROFILE, DEVICE_PROFILE -> { |
|
|
|
boolean isCommonEntity = calculatedField.getConfiguration().getReferencedEntities().contains(entityId); |
|
|
|
if (isCommonEntity) { |
|
|
|
getOrFetchFromDBProfileEntities(tenantId, cfEntityId).forEach(id -> updateOrInitializeState(calculatedFieldCtx, id, argumentValues)); |
|
|
|
} else { |
|
|
|
updateOrInitializeState(calculatedFieldCtx, entityId, argumentValues); |
|
|
|
} |
|
|
|
} |
|
|
|
default -> updateOrInitializeState(calculatedFieldCtx, cfEntityId, argumentValues); |
|
|
|
} |
|
|
|
log.info("Successfully updated telemetry for calculatedFieldId: [{}]", calculatedFieldId); |
|
|
|
} |
|
|
|
|
|
|
|
private String getMappedKey(KvEntry entry, Map<String, String> attributes, Map<String, String> timeSeries) { |
|
|
|
if (entry instanceof AttributeKvEntry) { |
|
|
|
return attributes.entrySet().stream() |
|
|
|
.filter(attr -> attr.getValue().equals(entry.getKey())) |
|
|
|
.map(Map.Entry::getKey) |
|
|
|
.findFirst() |
|
|
|
.orElse(entry.getKey()); |
|
|
|
} else if (entry instanceof TsKvEntry) { |
|
|
|
return timeSeries.entrySet().stream() |
|
|
|
.filter(ts -> ts.getValue().equals(entry.getKey())) |
|
|
|
.map(Map.Entry::getKey) |
|
|
|
.findFirst() |
|
|
|
.orElse(entry.getKey()); |
|
|
|
} |
|
|
|
return entry.getKey(); |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
@ -493,26 +559,29 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
if (state == null) { |
|
|
|
state = createStateByType(calculatedFieldCtx.getCfType()); |
|
|
|
} |
|
|
|
state.initState(argumentValues); |
|
|
|
calculatedFieldEntityCtx.setState(state); |
|
|
|
states.put(entityCtxId, calculatedFieldEntityCtx); |
|
|
|
rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); |
|
|
|
|
|
|
|
ListenableFuture<CalculatedFieldResult> resultFuture = state.performCalculation(calculatedFieldCtx); |
|
|
|
Futures.addCallback(resultFuture, new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(CalculatedFieldResult result) { |
|
|
|
if (result != null) { |
|
|
|
pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); |
|
|
|
} |
|
|
|
} |
|
|
|
if (state.updateState(argumentValues)) { |
|
|
|
calculatedFieldEntityCtx.setState(state); |
|
|
|
states.put(entityCtxId, calculatedFieldEntityCtx); |
|
|
|
rocksDBService.put(JacksonUtil.writeValueAsString(entityCtxId), JacksonUtil.writeValueAsString(calculatedFieldEntityCtx)); |
|
|
|
|
|
|
|
boolean allArgsPresent = calculatedFieldCtx.getArguments().keySet().containsAll(state.getArguments().keySet()); |
|
|
|
if (allArgsPresent) { |
|
|
|
ListenableFuture<CalculatedFieldResult> resultFuture = state.performCalculation(calculatedFieldCtx); |
|
|
|
Futures.addCallback(resultFuture, new FutureCallback<>() { |
|
|
|
@Override |
|
|
|
public void onSuccess(CalculatedFieldResult result) { |
|
|
|
if (result != null) { |
|
|
|
pushMsgToRuleEngine(calculatedFieldCtx.getTenantId(), entityId, result); |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); |
|
|
|
@Override |
|
|
|
public void onFailure(Throwable t) { |
|
|
|
log.warn("[{}] Failed to perform calculation. entityId: [{}]", calculatedFieldCtx.getCfId(), entityId, t); |
|
|
|
} |
|
|
|
}, MoreExecutors.directExecutor()); |
|
|
|
} |
|
|
|
}, MoreExecutors.directExecutor()); |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
private CalculatedFieldEntityCtx fetchCalculatedFieldEntityState(CalculatedFieldEntityCtxId entityCtxId) { |
|
|
|
@ -520,14 +589,14 @@ public class DefaultCalculatedFieldExecutionService extends AbstractPartitionBas |
|
|
|
if (stateStr == null) { |
|
|
|
return new CalculatedFieldEntityCtx(entityCtxId, null); |
|
|
|
} |
|
|
|
return JacksonUtil.fromString(rocksDBService.get(JacksonUtil.writeValueAsString(entityCtxId)), CalculatedFieldEntityCtx.class); |
|
|
|
return JacksonUtil.fromString(stateStr, CalculatedFieldEntityCtx.class); |
|
|
|
} |
|
|
|
|
|
|
|
private void pushMsgToRuleEngine(TenantId tenantId, EntityId originatorId, CalculatedFieldResult calculatedFieldResult) { |
|
|
|
try { |
|
|
|
String type = calculatedFieldResult.getType(); |
|
|
|
TbMsgType msgType = "ATTRIBUTES".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; |
|
|
|
TbMsgMetaData md = "ATTRIBUTES".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; |
|
|
|
TbMsgType msgType = "ATTRIBUTE".equals(type) ? TbMsgType.POST_ATTRIBUTES_REQUEST : TbMsgType.POST_TELEMETRY_REQUEST; |
|
|
|
TbMsgMetaData md = "ATTRIBUTE".equals(type) ? new TbMsgMetaData(Map.of(SCOPE, calculatedFieldResult.getScope().name())) : TbMsgMetaData.EMPTY; |
|
|
|
ObjectNode payload = createJsonPayload(calculatedFieldResult); |
|
|
|
TbMsg msg = TbMsg.newMsg(msgType, originatorId, md, JacksonUtil.writeValueAsString(payload)); |
|
|
|
clusterService.pushMsgToRuleEngine(tenantId, originatorId, msg, null); |
|
|
|
|