Browse Source

removed cf name from telemetry result

pull/14441/head
IrynaMatveieva 10 months ago
parent
commit
71204e2e24
  1. 2
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 12
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  3. 2
      application/src/main/java/org/thingsboard/server/service/cf/AlarmCalculatedFieldResult.java
  4. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldProcessingService.java
  5. 2
      application/src/main/java/org/thingsboard/server/service/cf/CalculatedFieldResult.java
  6. 20
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  7. 4
      application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java
  8. 5
      application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java
  9. 2
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  10. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java
  11. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java
  12. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  13. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  14. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java
  15. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java

2
application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java

@ -492,7 +492,7 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
stateSizeChecked = true; stateSizeChecked = true;
if (state.isSizeOk()) { if (state.isSizeOk()) {
if (!calculationResult.isEmpty()) { if (!calculationResult.isEmpty()) {
cfService.processResult(tenantId, entityId, calculationResult, cfIdList, callback); cfService.processResult(tenantId, entityId, ctx.getCfName(), calculationResult, cfIdList, callback);
} else { } else {
callback.onSuccess(); callback.onSuccess();
} }

12
application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java

@ -422,13 +422,13 @@ public abstract class AbstractCalculatedFieldProcessingService {
} }
} }
protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List<CalculatedFieldId> cfIds, TbCallback callback) { protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, String cfName, TelemetryCalculatedFieldResult cfResult, List<CalculatedFieldId> cfIds, TbCallback callback) {
OutputType type = cfResult.getType(); OutputType type = cfResult.getType();
JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue())); JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue()));
log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult); log.trace("[{}][{}] Saving CF result: {}", tenantId, entityId, jsonResult);
switch (type) { switch (type) {
case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfResult.getCalculatedFieldName(), cfIds, callback); case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfName, cfIds, callback);
case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), callback); case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), callback);
} }
} }
@ -477,18 +477,20 @@ public abstract class AbstractCalculatedFieldProcessingService {
.entries(entries) .entries(entries)
.strategy(strategy) .strategy(strategy)
.previousCalculatedFieldIds(cfIds) .previousCalculatedFieldIds(cfIds)
.callback(new FutureCallback<Void>() { .callback(new FutureCallback<>() {
@Override @Override
public void onSuccess(Void result) { public void onSuccess(Void result) {
if (sendAttributesUpdatedNotification) { if (sendAttributesUpdatedNotification) {
sendAttributesUpdatedMsg(tenantId, entityId, scope, cfName, entries); sendAttributesUpdatedMsg(tenantId, entityId, scope, cfName, entries);
} }
callback.onSuccess(); callback.onSuccess();
log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, entries);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
callback.onFailure(t); callback.onFailure(t);
log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, entries, t);
} }
}) })
.build()); .build());
@ -515,15 +517,17 @@ public abstract class AbstractCalculatedFieldProcessingService {
.entityId(entityId) .entityId(entityId)
.entries(tsEntries) .entries(tsEntries)
.strategy(strategy) .strategy(strategy)
.callback(new FutureCallback<Void>() { .callback(new FutureCallback<>() {
@Override @Override
public void onSuccess(Void result) { public void onSuccess(Void result) {
callback.onSuccess(); callback.onSuccess();
log.debug("[{}][{}] Saved CF result: {}", tenantId, entityId, tsEntries);
} }
@Override @Override
public void onFailure(Throwable t) { public void onFailure(Throwable t) {
callback.onFailure(t); callback.onFailure(t);
log.error("[{}][{}] Failed to save CF result {}", tenantId, entityId, tsEntries, t);
} }
}); });
if (ttl != null) { if (ttl != null) {

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

@ -37,7 +37,7 @@ public class AlarmCalculatedFieldResult implements CalculatedFieldResult {
private final TbAlarmResult alarmResult; private final TbAlarmResult alarmResult;
@Override @Override
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { public TbMsg toTbMsg(EntityId entityId, String cfName, List<CalculatedFieldId> cfIds) {
TbMsgType msgType; TbMsgType msgType;
TbMsgMetaData metaData = new TbMsgMetaData(); TbMsgMetaData metaData = new TbMsgMetaData();
if (alarmResult.isCreated()) { if (alarmResult.isCreated()) {

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

@ -41,7 +41,7 @@ public interface CalculatedFieldProcessingService {
Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments); Map<String, ArgumentEntry> fetchArgsFromDb(TenantId tenantId, EntityId entityId, Map<String, Argument> arguments);
void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback); void processResult(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback);
ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval); ArgumentEntry fetchMetricDuringInterval(TenantId tenantId, EntityId entityId, String argKey, AggMetric metric, AggIntervalEntry interval);

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

@ -23,7 +23,7 @@ import java.util.List;
public interface CalculatedFieldResult { public interface CalculatedFieldResult {
TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds); TbMsg toTbMsg(EntityId entityId, String cfName, List<CalculatedFieldId> cfIds);
String stringValue(); String stringValue();

20
application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java

@ -133,40 +133,40 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
return super.fetchMetricDuringInterval(tenantId, entityId, argKey, metric, interval); return super.fetchMetricDuringInterval(tenantId, entityId, argKey, metric, interval);
} }
public void processResult(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) { public void processResult(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
if (result instanceof AlarmCalculatedFieldResult) { if (result instanceof AlarmCalculatedFieldResult) {
sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfIds)); sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds));
return; return;
} }
TelemetryCalculatedFieldResult telemetryResult = result instanceof TelemetryCalculatedFieldResult telemetryRes TelemetryCalculatedFieldResult telemetryResult = result instanceof TelemetryCalculatedFieldResult telemetryRes
? telemetryRes : ((PropagationCalculatedFieldResult) result).getResult(); ? telemetryRes : ((PropagationCalculatedFieldResult) result).getResult();
switch (telemetryResult.getOutputStrategy().getType()) { switch (telemetryResult.getOutputStrategy().getType()) {
case IMMEDIATE -> processImmediately(tenantId, entityId, result, cfIds, callback); case IMMEDIATE -> processImmediately(tenantId, entityId, cfName, result, cfIds, callback);
case RULE_CHAIN -> pushMsgToRuleEngine(tenantId, entityId, result, cfIds, callback); case RULE_CHAIN -> pushMsgToRuleEngine(tenantId, entityId, cfName, result, cfIds, callback);
} }
} }
private void processImmediately(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) { private void processImmediately(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
if (result instanceof TelemetryCalculatedFieldResult telemetryResult) { if (result instanceof TelemetryCalculatedFieldResult telemetryResult) {
saveTelemetryResult(tenantId, entityId, telemetryResult, cfIds, callback); saveTelemetryResult(tenantId, entityId, cfName, telemetryResult, cfIds, callback);
return; return;
} }
if (result instanceof PropagationCalculatedFieldResult propagationResult) { if (result instanceof PropagationCalculatedFieldResult propagationResult) {
handlePropagationResults(propagationResult, callback, handlePropagationResults(propagationResult, callback,
(entity, res, cb) -> saveTelemetryResult(tenantId, entity, res, cfIds, cb)); (entity, res, cb) -> saveTelemetryResult(tenantId, entity, cfName, res, cfIds, cb));
return; return;
} }
callback.onSuccess(); callback.onSuccess();
} }
private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) { private void pushMsgToRuleEngine(TenantId tenantId, EntityId entityId, String cfName, CalculatedFieldResult result, List<CalculatedFieldId> cfIds, TbCallback callback) {
if (result instanceof PropagationCalculatedFieldResult propagationResult) { if (result instanceof PropagationCalculatedFieldResult propagationResult) {
handlePropagationResults(propagationResult, callback, handlePropagationResults(propagationResult, callback,
(entity, res, cb) -> sendMsgToRuleEngine(tenantId, entity, cb, res.toTbMsg(entity, cfIds))); (entity, res, cb) -> sendMsgToRuleEngine(tenantId, entity, cb, res.toTbMsg(entity, cfName, cfIds)));
return; return;
} }
sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfIds)); sendMsgToRuleEngine(tenantId, entityId, callback, result.toTbMsg(entityId, cfName, cfIds));
} }
private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback, private void handlePropagationResults(PropagationCalculatedFieldResult propagationResult, TbCallback callback,

4
application/src/main/java/org/thingsboard/server/service/cf/PropagationCalculatedFieldResult.java

@ -32,8 +32,8 @@ public final class PropagationCalculatedFieldResult implements CalculatedFieldRe
private final TelemetryCalculatedFieldResult result; private final TelemetryCalculatedFieldResult result;
@Override @Override
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { public TbMsg toTbMsg(EntityId entityId, String cfName, List<CalculatedFieldId> cfIds) {
return result.toTbMsg(entityId, cfIds); return result.toTbMsg(entityId, cfName, cfIds);
} }
@Override @Override

5
application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java

@ -36,7 +36,6 @@ import static org.thingsboard.server.common.data.DataConstants.SCOPE;
@Builder @Builder
public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult { public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult {
private final String calculatedFieldName;
private final OutputType type; private final OutputType type;
private final AttributeScope scope; private final AttributeScope scope;
private final OutputStrategy outputStrategy; private final OutputStrategy outputStrategy;
@ -45,13 +44,13 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu
public static final TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build(); public static final TelemetryCalculatedFieldResult EMPTY = TelemetryCalculatedFieldResult.builder().result(null).build();
@Override @Override
public TbMsg toTbMsg(EntityId entityId, List<CalculatedFieldId> cfIds) { public TbMsg toTbMsg(EntityId entityId, String cfName, List<CalculatedFieldId> cfIds) {
TbMsgType msgType = switch (type) { TbMsgType msgType = switch (type) {
case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST; case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST;
case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST; case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST;
}; };
TbMsgMetaData metaData = new TbMsgMetaData(); TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue(CF_NAME_METADATA_KEY, calculatedFieldName); metaData.putValue(CF_NAME_METADATA_KEY, cfName);
if (OutputType.ATTRIBUTES == type) { if (OutputType.ATTRIBUTES == type) {
metaData.putValue(SCOPE, scope.name()); metaData.putValue(SCOPE, scope.name());
} }

2
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

@ -86,6 +86,7 @@ public class CalculatedFieldCtx implements Closeable {
private CalculatedField calculatedField; private CalculatedField calculatedField;
private CalculatedFieldId cfId; private CalculatedFieldId cfId;
private String cfName;
private TenantId tenantId; private TenantId tenantId;
private EntityId entityId; private EntityId entityId;
private CalculatedFieldType cfType; private CalculatedFieldType cfType;
@ -130,6 +131,7 @@ public class CalculatedFieldCtx implements Closeable {
this.calculatedField = calculatedField; this.calculatedField = calculatedField;
this.cfId = calculatedField.getId(); this.cfId = calculatedField.getId();
this.cfName = calculatedField.getName();
this.tenantId = calculatedField.getTenantId(); this.tenantId = calculatedField.getTenantId();
this.entityId = calculatedField.getEntityId(); this.entityId = calculatedField.getEntityId();
this.cfType = calculatedField.getType(); this.cfType = calculatedField.getType();

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java

@ -52,7 +52,6 @@ public class ScriptCalculatedFieldState extends BaseCalculatedFieldState {
Output output = ctx.getOutput(); Output output = ctx.getOutput();
return Futures.transform(resultFuture, return Futures.transform(resultFuture,
result -> TelemetryCalculatedFieldResult.builder() result -> TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(output.getStrategy()) .outputStrategy(output.getStrategy())
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java

@ -56,7 +56,6 @@ public class SimpleCalculatedFieldState extends BaseCalculatedFieldState {
JsonNode outputResult = createResultJson(ctx.isUseLatestTs(), output.getName(), result); JsonNode outputResult = createResultJson(ctx.isUseLatestTs(), output.getName(), result);
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(output.getStrategy()) .outputStrategy(output.getStrategy())
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java

@ -182,7 +182,6 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
lastMetricsEvalTs = System.currentTimeMillis(); lastMetricsEvalTs = System.currentTimeMillis();
scheduleReevaluation(); scheduleReevaluation();
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(output.getStrategy()) .outputStrategy(output.getStrategy())
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java

@ -123,7 +123,6 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
return Futures.immediateFuture(TelemetryCalculatedFieldResult.EMPTY); return Futures.immediateFuture(TelemetryCalculatedFieldResult.EMPTY);
} }
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(output.getStrategy()) .outputStrategy(output.getStrategy())
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java

@ -133,7 +133,6 @@ public class GeofencingCalculatedFieldState extends BaseCalculatedFieldState {
OutputType outputType = ctx.getOutput().getType(); OutputType outputType = ctx.getOutput().getType();
var result = TelemetryCalculatedFieldResult.builder() var result = TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(ctx.getOutput().getStrategy()) .outputStrategy(ctx.getOutput().getStrategy())
.type(outputType) .type(outputType)
.scope(ctx.getOutput().getScope()) .scope(ctx.getOutput().getScope())

1
application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java

@ -85,7 +85,6 @@ public class PropagationCalculatedFieldState extends ScriptCalculatedFieldState
Output output = ctx.getOutput(); Output output = ctx.getOutput();
TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder = TelemetryCalculatedFieldResult.TelemetryCalculatedFieldResultBuilder telemetryCfBuilder =
TelemetryCalculatedFieldResult.builder() TelemetryCalculatedFieldResult.builder()
.calculatedFieldName(ctx.getCalculatedField().getName())
.outputStrategy(output.getStrategy()) .outputStrategy(output.getStrategy())
.type(output.getType()) .type(output.getType())
.scope(output.getScope()); .scope(output.getScope());

Loading…
Cancel
Save