Browse Source

handle output update to avoid processing when minor strategy properties updated

pull/14535/head
IrynaMatveieva 10 months ago
parent
commit
e37332a288
  1. 10
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldEntityMessageProcessor.java
  2. 2
      application/src/main/java/org/thingsboard/server/actors/calculatedField/CalculatedFieldManagerMessageProcessor.java
  3. 3
      application/src/main/java/org/thingsboard/server/actors/calculatedField/EntityInitCalculatedFieldMsg.java
  4. 94
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java

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

@ -159,10 +159,14 @@ public class CalculatedFieldEntityMessageProcessor extends AbstractContextAwareM
} else { } else {
state.setCtx(ctx, actorCtx); state.setCtx(ctx, actorCtx);
} }
if (state.isSizeOk()) { if (msg.getStateAction() != StateAction.REFRESH_CTX) {
processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback()); if (state.isSizeOk()) {
processStateIfReady(state, Collections.emptyMap(), ctx, Collections.singletonList(ctx.getCfId()), null, null, msg.getCallback());
} else {
throw new RuntimeException(ctx.getSizeExceedsLimitMessage());
}
} else { } else {
throw new RuntimeException(ctx.getSizeExceedsLimitMessage()); msg.getCallback().onSuccess();
} }
} catch (Exception e) { } catch (Exception e) {
log.debug("[{}][{}] Failed to initialize CF state", entityId, ctx.getCfId(), e); log.debug("[{}][{}] Failed to initialize CF state", entityId, ctx.getCfId(), e);

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

@ -455,6 +455,8 @@ public class CalculatedFieldManagerMessageProcessor extends AbstractContextAware
stateAction = StateAction.REINIT; // refetch arguments, call state.init, then calculate stateAction = StateAction.REINIT; // refetch arguments, call state.init, then calculate
} else if (newCfCtx.hasContextOnlyChanges(oldCfCtx)) { } else if (newCfCtx.hasContextOnlyChanges(oldCfCtx)) {
stateAction = StateAction.REPROCESS; // call state.setCtx, then calculate stateAction = StateAction.REPROCESS; // call state.setCtx, then calculate
} else if (newCfCtx.hasRefreshContextOnlyChanges(oldCfCtx)) {
stateAction = StateAction.REFRESH_CTX;
} else { } else {
callback.onSuccess(); callback.onSuccess();
return; return;

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

@ -39,6 +39,7 @@ public class EntityInitCalculatedFieldMsg implements ToCalculatedFieldSystemMsg
INIT, INIT,
REINIT, REINIT,
RECREATE, RECREATE,
REPROCESS REPROCESS,
REFRESH_CTX
} }
} }

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

@ -37,12 +37,14 @@ import org.thingsboard.server.common.data.cf.configuration.AlarmCalculatedFieldC
import org.thingsboard.server.common.data.cf.configuration.Argument; import org.thingsboard.server.common.data.cf.configuration.Argument;
import org.thingsboard.server.common.data.cf.configuration.ArgumentType; import org.thingsboard.server.common.data.cf.configuration.ArgumentType;
import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.ArgumentsBasedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.AttributesImmediateOutputStrategy;
import org.thingsboard.server.common.data.cf.configuration.ExpressionBasedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.ExpressionBasedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.Output; import org.thingsboard.server.common.data.cf.configuration.Output;
import org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.PropagationCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey; import org.thingsboard.server.common.data.cf.configuration.ReferencedEntityKey;
import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.ScheduledUpdateSupportedCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.SimpleCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.TimeSeriesImmediateOutputStrategy;
import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunctionInput; import org.thingsboard.server.common.data.cf.configuration.aggregation.AggFunctionInput;
import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.RelatedEntitiesAggregationCalculatedFieldConfiguration;
import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration; import org.thingsboard.server.common.data.cf.configuration.aggregation.single.EntityAggregationCalculatedFieldConfiguration;
@ -623,16 +625,42 @@ public class CalculatedFieldCtx implements Closeable {
return new CalculatedFieldEntityCtxId(tenantId, cfId, entityId); return new CalculatedFieldEntityCtxId(tenantId, cfId, entityId);
} }
public boolean hasRefreshContextOnlyChanges(CalculatedFieldCtx other) { // has changes that do not require state recalculation
var thisConfig = calculatedField.getConfiguration();
var otherConfig = other.getCalculatedField().getConfiguration();
var thisOutputStrategy = thisConfig.getOutput().getStrategy();
var otherOutputStrategy = otherConfig.getOutput().getStrategy();
if (!thisOutputStrategy.getType().equals(otherOutputStrategy.getType())) {
return true;
}
if (thisOutputStrategy instanceof TimeSeriesImmediateOutputStrategy thisTimeSeriesImmediateOutputStrategy
&& otherOutputStrategy instanceof TimeSeriesImmediateOutputStrategy otherTimeSeriesImmediateOutputStrategy) {
return thisTimeSeriesImmediateOutputStrategy.getTtl() != otherTimeSeriesImmediateOutputStrategy.getTtl();
}
if (thisOutputStrategy instanceof AttributesImmediateOutputStrategy thisAttributesImmediateOutputStrategy
&& otherOutputStrategy instanceof AttributesImmediateOutputStrategy otherAttributesImmediateOutputStrategy) {
boolean updateAttributesOnlyOnValueChangeChanged = thisAttributesImmediateOutputStrategy.isUpdateAttributesOnlyOnValueChange() != otherAttributesImmediateOutputStrategy.isUpdateAttributesOnlyOnValueChange();
boolean sendAttributesUpdatedNotificationUpdated = thisAttributesImmediateOutputStrategy.isSendAttributesUpdatedNotification() != otherAttributesImmediateOutputStrategy.isSendAttributesUpdatedNotification();
return updateAttributesOnlyOnValueChangeChanged || sendAttributesUpdatedNotificationUpdated;
}
return false;
}
public boolean hasContextOnlyChanges(CalculatedFieldCtx other) { // has changes that do not require state reinit and will be picked up by the state on the fly public boolean hasContextOnlyChanges(CalculatedFieldCtx other) { // has changes that do not require state reinit and will be picked up by the state on the fly
if (calculatedField.getConfiguration() instanceof ExpressionBasedCalculatedFieldConfiguration && !Objects.equals(expression, other.expression)) { if (calculatedField.getConfiguration() instanceof ExpressionBasedCalculatedFieldConfiguration && !Objects.equals(expression, other.expression)) {
return true; return true;
} }
if (!Objects.equals(output, other.output)) { if (hasOutputChanges(other.output)) {
return true; return true;
} }
if (calculatedField.getConfiguration() instanceof SimpleCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof SimpleCalculatedFieldConfiguration thisConfig
&& other.calculatedField.getConfiguration() instanceof SimpleCalculatedFieldConfiguration otherConfig && other.calculatedField.getConfiguration() instanceof SimpleCalculatedFieldConfiguration otherConfig
&& thisConfig.isUseLatestTs() != otherConfig.isUseLatestTs()) { && thisConfig.isUseLatestTs() != otherConfig.isUseLatestTs()) {
return true; return true;
} }
if (cfType == CalculatedFieldType.ALARM) { if (cfType == CalculatedFieldType.ALARM) {
@ -654,14 +682,14 @@ public class CalculatedFieldCtx implements Closeable {
return true; return true;
} }
if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration thisConfig
&& other.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration otherConfig && other.getCalculatedField().getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration otherConfig
&& (thisConfig.getDeduplicationIntervalInSec() != otherConfig.getDeduplicationIntervalInSec() && (thisConfig.getDeduplicationIntervalInSec() != otherConfig.getDeduplicationIntervalInSec()
|| !thisConfig.getMetrics().equals(otherConfig.getMetrics()) || !thisConfig.getMetrics().equals(otherConfig.getMetrics())
|| thisConfig.isUseLatestTs() != otherConfig.isUseLatestTs())) { || thisConfig.isUseLatestTs() != otherConfig.isUseLatestTs())) {
return true; return true;
} }
if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig
&& other.getCalculatedField().getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) { && other.getCalculatedField().getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) {
boolean metricsChanged = !Objects.equals(thisConfig.getMetrics(), otherConfig.getMetrics()); boolean metricsChanged = !Objects.equals(thisConfig.getMetrics(), otherConfig.getMetrics());
boolean watermarkChanged = !Objects.equals(thisConfig.getWatermark(), otherConfig.getWatermark()); boolean watermarkChanged = !Objects.equals(thisConfig.getWatermark(), otherConfig.getWatermark());
return metricsChanged || watermarkChanged; return metricsChanged || watermarkChanged;
@ -695,9 +723,47 @@ public class CalculatedFieldCtx implements Closeable {
return false; return false;
} }
private boolean hasOutputChanges(Output otherOutput) {
if (!output.getType().equals(otherOutput.getType())) {
return true;
}
if (!output.getName().equals(otherOutput.getName())) {
return true;
}
if (output.getScope() != (otherOutput.getScope())) {
return true;
}
if (!Objects.equals(output.getDecimalsByDefault(), otherOutput.getDecimalsByDefault())) {
return true;
}
var thisOutputStrategy = output.getStrategy();
var otherOutputStrategy = otherOutput.getStrategy();
if (thisOutputStrategy instanceof TimeSeriesImmediateOutputStrategy thisTimeSeriesImmediateOutputStrategy
&& otherOutputStrategy instanceof TimeSeriesImmediateOutputStrategy otherTimeSeriesImmediateOutputStrategy) {
boolean saveTimeSeriesUpdated = thisTimeSeriesImmediateOutputStrategy.isSaveTimeSeries() != otherTimeSeriesImmediateOutputStrategy.isSaveTimeSeries();
boolean saveLatestUpdated = thisTimeSeriesImmediateOutputStrategy.isSaveLatest() != otherTimeSeriesImmediateOutputStrategy.isSaveLatest();
boolean sendWsUpdateUpdated = thisTimeSeriesImmediateOutputStrategy.isSendWsUpdate() != otherTimeSeriesImmediateOutputStrategy.isSendWsUpdate();
boolean processCfsUpdated = thisTimeSeriesImmediateOutputStrategy.isProcessCfs() != otherTimeSeriesImmediateOutputStrategy.isProcessCfs();
return saveTimeSeriesUpdated || saveLatestUpdated || sendWsUpdateUpdated || processCfsUpdated;
}
if (thisOutputStrategy instanceof AttributesImmediateOutputStrategy thisAttributesImmediateOutputStrategy
&& otherOutputStrategy instanceof AttributesImmediateOutputStrategy otherAttributesImmediateOutputStrategy) {
boolean saveTimeSeriesUpdated = thisAttributesImmediateOutputStrategy.isSaveAttribute() != otherAttributesImmediateOutputStrategy.isSaveAttribute();
boolean sendWsUpdateUpdated = thisAttributesImmediateOutputStrategy.isSendWsUpdate() != otherAttributesImmediateOutputStrategy.isSendWsUpdate();
boolean processCfsUpdated = thisAttributesImmediateOutputStrategy.isProcessCfs() != otherAttributesImmediateOutputStrategy.isProcessCfs();
return saveTimeSeriesUpdated || sendWsUpdateUpdated || processCfsUpdated;
}
return false;
}
private boolean hasGeofencingZoneGroupConfigurationChanges(CalculatedFieldCtx other) { private boolean hasGeofencingZoneGroupConfigurationChanges(CalculatedFieldCtx other) {
if (calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration thisConfig
&& other.calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration otherConfig) { && other.calculatedField.getConfiguration() instanceof GeofencingCalculatedFieldConfiguration otherConfig) {
return !thisConfig.getZoneGroups().equals(otherConfig.getZoneGroups()); return !thisConfig.getZoneGroups().equals(otherConfig.getZoneGroups());
} }
return false; return false;
@ -705,7 +771,7 @@ public class CalculatedFieldCtx implements Closeable {
private boolean hasRelatedEntitiesAggregationConfigurationChanges(CalculatedFieldCtx other) { private boolean hasRelatedEntitiesAggregationConfigurationChanges(CalculatedFieldCtx other) {
if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration thisConfig
&& other.calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration otherConfig) { && other.calculatedField.getConfiguration() instanceof RelatedEntitiesAggregationCalculatedFieldConfiguration otherConfig) {
return !thisConfig.getRelation().equals(otherConfig.getRelation()); return !thisConfig.getRelation().equals(otherConfig.getRelation());
} }
return false; return false;
@ -713,7 +779,7 @@ public class CalculatedFieldCtx implements Closeable {
private boolean hasEntityAggregationConfigurationChanges(CalculatedFieldCtx other) { private boolean hasEntityAggregationConfigurationChanges(CalculatedFieldCtx other) {
if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig if (calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration thisConfig
&& other.calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) { && other.calculatedField.getConfiguration() instanceof EntityAggregationCalculatedFieldConfiguration otherConfig) {
return !thisConfig.getInterval().equals(otherConfig.getInterval()); return !thisConfig.getInterval().equals(otherConfig.getInterval());
} }
return false; return false;
@ -738,7 +804,7 @@ public class CalculatedFieldCtx implements Closeable {
yield true; yield true;
} }
yield geofencingState.getLastDynamicArgumentsRefreshTs() < yield geofencingState.getLastDynamicArgumentsRefreshTs() <
System.currentTimeMillis() - scheduledUpdateIntervalMillis; System.currentTimeMillis() - scheduledUpdateIntervalMillis;
} }
default -> false; default -> false;
}; };
@ -782,10 +848,10 @@ public class CalculatedFieldCtx implements Closeable {
@Override @Override
public String toString() { public String toString() {
return "CalculatedFieldCtx{" + return "CalculatedFieldCtx{" +
"cfId=" + cfId + "cfId=" + cfId +
", cfType=" + cfType + ", cfType=" + cfType +
", entityId=" + entityId + ", entityId=" + entityId +
'}'; '}';
} }
} }

Loading…
Cancel
Save