Browse Source

added sendAttributesUpdated flag to cf output strategy and added cf name to metadata

pull/14441/head
IrynaMatveieva 8 months ago
parent
commit
c038bd8fb9
  1. 85
      application/src/main/java/org/thingsboard/server/service/cf/AbstractCalculatedFieldProcessingService.java
  2. 25
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldProcessingService.java
  3. 12
      application/src/main/java/org/thingsboard/server/service/cf/TelemetryCalculatedFieldResult.java
  4. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/ScriptCalculatedFieldState.java
  5. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/SimpleCalculatedFieldState.java
  6. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  7. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  8. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/geofencing/GeofencingCalculatedFieldState.java
  9. 1
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/propagation/PropagationCalculatedFieldState.java
  10. 1
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java
  11. 1
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java

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

@ -15,6 +15,7 @@
*/
package org.thingsboard.server.service.cf;
import com.fasterxml.jackson.databind.node.ObjectNode;
import com.google.common.util.concurrent.FutureCallback;
import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture;
@ -28,10 +29,12 @@ import jakarta.annotation.PreDestroy;
import lombok.Data;
import lombok.extern.slf4j.Slf4j;
import org.thingsboard.common.util.DonAsynchron;
import org.thingsboard.common.util.JacksonUtil;
import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.rule.engine.api.AttributesSaveRequest;
import org.thingsboard.rule.engine.api.AttributesSaveRequest.Strategy;
import org.thingsboard.rule.engine.api.TimeseriesSaveRequest;
import org.thingsboard.server.cluster.TbClusterService;
import org.thingsboard.server.common.adaptor.JsonConverter;
import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.cf.CalculatedField;
@ -62,11 +65,15 @@ import org.thingsboard.server.common.data.relation.EntityRelation;
import org.thingsboard.server.common.data.relation.EntityRelationPathQuery;
import org.thingsboard.server.common.data.relation.RelationPathLevel;
import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileConfiguration;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.dao.attributes.AttributesService;
import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.timeseries.TimeseriesService;
import org.thingsboard.server.dao.usagerecord.ApiLimitService;
import org.thingsboard.server.queue.TbQueueCallback;
import org.thingsboard.server.queue.TbQueueMsgMetadata;
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.SingleValueArgumentEntry;
@ -85,10 +92,13 @@ import java.util.function.Function;
import java.util.function.Predicate;
import java.util.stream.Collectors;
import static org.thingsboard.server.common.data.DataConstants.CF_NAME_METADATA_KEY;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
import static org.thingsboard.server.common.data.cf.CalculatedFieldType.PROPAGATION;
import static org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration.PROPAGATION_CONFIG_ARGUMENT;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LATITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.common.data.cf.configuration.geofencing.EntityCoordinates.ENTITY_ID_LONGITUDE_ARGUMENT_KEY;
import static org.thingsboard.server.common.data.msg.TbMsgType.ATTRIBUTES_UPDATED;
import static org.thingsboard.server.dao.util.KvUtils.filterChangedAttr;
import static org.thingsboard.server.dao.util.KvUtils.toTsKvEntryList;
import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDefaultAttributeEntry;
@ -108,6 +118,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
protected final ApiLimitService apiLimitService;
protected final RelationService relationService;
protected final OwnerService ownerService;
protected final TbClusterService clusterService;
protected ListeningExecutorService calculatedFieldCallbackExecutor;
@ -392,6 +403,26 @@ public abstract class AbstractCalculatedFieldProcessingService {
return new BaseReadTsKvQuery(argument.getRefEntityKey().getKey(), startTs, endTs, 0, limit, Aggregation.NONE);
}
protected void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) {
try {
clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
log.trace("[{}][{}] Pushed message to rule engine: {} ", tenantId, entityId, msg);
callback.onSuccess();
}
@Override
public void onFailure(Throwable t) {
callback.onFailure(t);
}
});
} catch (Exception e) {
log.warn("[{}][{}] Failed to push message to rule engine: {}", tenantId, entityId, msg, e);
callback.onFailure(e);
}
}
protected void saveTelemetryResult(TenantId tenantId, EntityId entityId, TelemetryCalculatedFieldResult cfResult, List<CalculatedFieldId> cfIds, TbCallback callback) {
OutputType type = cfResult.getType();
JsonElement jsonResult = JsonParser.parseString(Objects.requireNonNull(cfResult.stringValue()));
@ -400,7 +431,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
SettableFuture<Void> future = SettableFuture.create();
switch (type) {
case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfIds, future);
case ATTRIBUTES -> saveAttributes(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfResult.getScope(), cfResult.getCalculatedFieldName(), cfIds, future);
case TIME_SERIES -> saveTimeSeries(tenantId, entityId, jsonResult, cfResult.getOutputStrategy(), cfIds, System.currentTimeMillis(), future);
}
@ -419,7 +450,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
}, MoreExecutors.directExecutor());
}
private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, List<CalculatedFieldId> cfIds, SettableFuture<Void> future) {
private void saveAttributes(TenantId tenantId, EntityId entityId, JsonElement jsonResult, OutputStrategy outputStrategy, AttributeScope scope, String cfName, List<CalculatedFieldId> cfIds, SettableFuture<Void> future) {
if (!(outputStrategy instanceof AttributesImmediateOutputStrategy attOutputStrategy)) {
future.setException(new IllegalArgumentException("Only AttributeImmediateOutputStrategy is supported."));
} else {
@ -427,7 +458,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
List<AttributeKvEntry> newAttributes = JsonConverter.convertToAttributes(jsonResult);
if (!attOutputStrategy.isUpdateAttributesOnlyOnValueChange()) {
saveAttributesInternal(tenantId, entityId, scope, cfIds, newAttributes, strategy, future);
saveAttributesInternal(tenantId, entityId, scope, cfName, cfIds, newAttributes, strategy, attOutputStrategy.isSendAttributesUpdatedNotification(), future);
return;
}
@ -441,7 +472,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
future.set(null);
return;
}
saveAttributesInternal(tenantId, entityId, scope, cfIds, changed, strategy, future);
saveAttributesInternal(tenantId, entityId, scope, cfName, cfIds, changed, strategy, attOutputStrategy.isSendAttributesUpdatedNotification(), future);
},
future::setException,
MoreExecutors.directExecutor());
@ -450,10 +481,15 @@ public abstract class AbstractCalculatedFieldProcessingService {
private void saveAttributesInternal(TenantId tenantId, EntityId entityId,
AttributeScope scope,
String cfName,
List<CalculatedFieldId> cfIds,
List<AttributeKvEntry> entries,
AttributesSaveRequest.Strategy strategy,
boolean sendAttributesUpdatedNotification,
SettableFuture<Void> future) {
Runnable onSuccess = sendAttributesUpdatedNotification
? () -> sendAttributesUpdatedMsg(tenantId, entityId, scope, cfName, entries)
: null;
tsSubService.saveAttributes(AttributesSaveRequest.builder()
.tenantId(tenantId)
.entityId(entityId)
@ -461,7 +497,7 @@ public abstract class AbstractCalculatedFieldProcessingService {
.entries(entries)
.strategy(strategy)
.previousCalculatedFieldIds(cfIds)
.future(future)
.callback(wrapWithSuccessHandler(future, onSuccess))
.build());
}
@ -496,4 +532,43 @@ public abstract class AbstractCalculatedFieldProcessingService {
tsSubService.saveTimeseries(builder.build());
}
private void sendAttributesUpdatedMsg(TenantId tenantId, EntityId entityId,
AttributeScope scope,
String cfName,
List<AttributeKvEntry> entries) {
ObjectNode entityNode = JacksonUtil.newObjectNode();
if (entries != null) {
entries.forEach(attributeKvEntry -> JacksonUtil.addKvEntry(entityNode, attributeKvEntry));
}
TbMsg attributesUpdatedMsg = TbMsg.newMsg()
.type(ATTRIBUTES_UPDATED)
.originator(entityId)
.data(JacksonUtil.toString(entityNode))
.metaData(new TbMsgMetaData(Map.of(
CF_NAME_METADATA_KEY, cfName,
SCOPE, scope.name()
)))
.build();
sendMsgToRuleEngine(tenantId, entityId, TbCallback.EMPTY, attributesUpdatedMsg);
}
private FutureCallback<Void> wrapWithSuccessHandler(SettableFuture<Void> future, Runnable onSuccess) {
return new FutureCallback<>() {
@Override
public void onSuccess(Void result) {
future.set(result);
if (onSuccess != null) {
onSuccess.run();
}
}
@Override
public void onFailure(Throwable t) {
future.setException(t);
}
};
}
}

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

@ -30,7 +30,6 @@ import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEn
import org.thingsboard.server.common.data.id.CalculatedFieldId;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.queue.ServiceType;
import org.thingsboard.server.common.msg.queue.TbCallback;
import org.thingsboard.server.common.msg.queue.TopicPartitionInfo;
@ -69,7 +68,6 @@ import static org.thingsboard.server.utils.CalculatedFieldUtils.toProto;
@Slf4j
public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedFieldProcessingService implements CalculatedFieldProcessingService {
private final TbClusterService clusterService;
private final PartitionService partitionService;
public DefaultCalculatedFieldProcessingService(AttributesService attributesService,
@ -80,8 +78,7 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
TbClusterService clusterService,
TelemetrySubscriptionService tsSubService,
PartitionService partitionService) {
super(attributesService, timeseriesService, tsSubService, apiLimitService, relationService, ownerService);
this.clusterService = clusterService;
super(attributesService, timeseriesService, tsSubService, apiLimitService, relationService, ownerService, clusterService);
this.partitionService = partitionService;
}
@ -190,26 +187,6 @@ public class DefaultCalculatedFieldProcessingService extends AbstractCalculatedF
}
}
private void sendMsgToRuleEngine(TenantId tenantId, EntityId entityId, TbCallback callback, TbMsg msg) {
try {
clusterService.pushMsgToRuleEngine(tenantId, entityId, msg, new TbQueueCallback() {
@Override
public void onSuccess(TbQueueMsgMetadata metadata) {
log.trace("[{}][{}] Pushed message to rule engine: {} ", tenantId, entityId, msg);
callback.onSuccess();
}
@Override
public void onFailure(Throwable t) {
callback.onFailure(t);
}
});
} catch (Exception e) {
log.warn("[{}][{}] Failed to push message to rule engine: {}", tenantId, entityId, msg, e);
callback.onFailure(e);
}
}
@Override
public void pushMsgToLinks(CalculatedFieldTelemetryMsg msg, List<CalculatedFieldEntityCtxId> linkedCalculatedFields, TbCallback callback) {
Map<TopicPartitionInfo, List<CalculatedFieldEntityCtxId>> unicasts = new HashMap<>();

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

@ -28,14 +28,15 @@ import org.thingsboard.server.common.msg.TbMsg;
import org.thingsboard.server.common.msg.TbMsgMetaData;
import java.util.List;
import java.util.Map;
import static org.thingsboard.server.common.data.DataConstants.CF_NAME_METADATA_KEY;
import static org.thingsboard.server.common.data.DataConstants.SCOPE;
@Data
@Builder
public final class TelemetryCalculatedFieldResult implements CalculatedFieldResult {
private final String calculatedFieldName;
private final OutputType type;
private final AttributeScope scope;
private final OutputStrategy outputStrategy;
@ -49,10 +50,11 @@ public final class TelemetryCalculatedFieldResult implements CalculatedFieldResu
case ATTRIBUTES -> TbMsgType.POST_ATTRIBUTES_REQUEST;
case TIME_SERIES -> TbMsgType.POST_TELEMETRY_REQUEST;
};
TbMsgMetaData metaData = switch (type) {
case ATTRIBUTES -> new TbMsgMetaData(Map.of(SCOPE, scope.name()));
case TIME_SERIES -> TbMsgMetaData.EMPTY;
};
TbMsgMetaData metaData = new TbMsgMetaData();
metaData.putValue(CF_NAME_METADATA_KEY, calculatedFieldName);
if (OutputType.ATTRIBUTES == type) {
metaData.putValue(SCOPE, scope.name());
}
return TbMsg.newMsg()
.type(msgType)
.originator(entityId)

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

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

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

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

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

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

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

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

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

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

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

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

1
common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

@ -41,6 +41,7 @@ public class DataConstants {
public static final String EDGE_ID = "edgeId";
public static final String DEVICE_ID = "deviceId";
public static final String GATEWAY_PARAMETER = "gateway";
public static final String CF_NAME_METADATA_KEY = "calculatedFieldName";
public static final String OVERWRITE_ACTIVITY_TIME_PARAMETER = "overwriteActivityTime";
public static final String COAP_TRANSPORT_NAME = "COAP";

1
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/AttributesImmediateOutputStrategy.java

@ -24,6 +24,7 @@ import lombok.NoArgsConstructor;
@NoArgsConstructor
public class AttributesImmediateOutputStrategy implements AttributesOutputStrategy {
private boolean sendAttributesUpdatedNotification;
private boolean updateAttributesOnlyOnValueChange;
private boolean saveAttribute;

Loading…
Cancel
Save