Browse Source

scheduled reevaluation on restart

pull/14141/head
IrynaMatveieva 11 months ago
parent
commit
1499f43caf
  1. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 5
      application/src/main/java/org/thingsboard/server/service/cf/DefaultCalculatedFieldCache.java
  3. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java
  4. 2
      application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java
  5. 1
      common/proto/src/main/proto/queue.proto

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

@ -123,6 +123,9 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
if (state != null) {
state.setCtx(msg.getCtx(), actorCtx);
state.setPartition(msg.getPartition());
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState relatedEntitiesAggState) {
relatedEntitiesAggState.scheduleReevaluation();
}
states.put(cfId, state);
} else {
removeState(cfId);

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

@ -149,8 +149,9 @@ public class DefaultCalculatedFieldCache implements CalculatedFieldCache {
@Override
public List<CalculatedFieldCtx> getAggCalculatedFieldCtxsByFilter(Predicate<CalculatedFieldCtx> relatedEntityFilter) {
return calculatedFieldsCtx.values().stream()
.filter(ctx -> CalculatedFieldType.RELATED_ENTITIES_AGGREGATION.equals(ctx.getCfType()))
return calculatedFields.values().stream()
.filter(cf -> CalculatedFieldType.RELATED_ENTITIES_AGGREGATION.equals(cf.getType()))
.map(cf -> getCalculatedFieldCtx(cf.getId()))
.filter(relatedEntityFilter)
.toList();
}

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

@ -67,6 +67,10 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec());
}
public void scheduleReevaluation() {
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
}
@Override
public void reset() { // must reset everything dependent on arguments
super.reset();

2
application/src/main/java/org/thingsboard/server/utils/CalculatedFieldUtils.java

@ -118,6 +118,7 @@ public class CalculatedFieldUtils {
}
if (state instanceof RelatedEntitiesAggregationCalculatedFieldState aggState) {
builder.setLastArgsUpdateTs(aggState.getLastArgsRefreshTs());
builder.setLastMetricsEvalTs(aggState.getLastMetricsEvalTs());
}
return builder.build();
}
@ -215,6 +216,7 @@ public class CalculatedFieldUtils {
relatedEntitiesAggState.getArguments().put(argName, new RelatedEntitiesArgumentEntry(entityInputs, false));
});
relatedEntitiesAggState.setLastArgsRefreshTs(proto.getLastArgsUpdateTs());
relatedEntitiesAggState.setLastMetricsEvalTs(proto.getLastMetricsEvalTs());
return relatedEntitiesAggState;
}

1
common/proto/src/main/proto/queue.proto

@ -924,6 +924,7 @@ message CalculatedFieldStateProto {
repeated GeofencingArgumentProto geofencingArguments = 5;
AlarmStateProto alarmState = 6;
int64 lastArgsUpdateTs = 7;
int64 lastMetricsEvalTs = 8;
}
//Used to report session state to tb-Service and persist this state in the cache on the tb-Service level.

Loading…
Cancel
Save