Browse Source

cancel future when stop

pull/14280/head
IrynaMatveieva 11 months ago
parent
commit
c239dfa190
  1. 21
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/RelatedEntitiesAggregationCalculatedFieldState.java

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

@ -43,6 +43,7 @@ import java.util.HashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Map.Entry; import java.util.Map.Entry;
import java.util.concurrent.ScheduledFuture;
import static java.util.concurrent.TimeUnit.SECONDS; import static java.util.concurrent.TimeUnit.SECONDS;
@ -59,6 +60,8 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
private long deduplicationIntervalMs = -1; private long deduplicationIntervalMs = -1;
private Map<String, AggMetric> metrics; private Map<String, AggMetric> metrics;
private ScheduledFuture<?> reevaluationFuture;
public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) { public RelatedEntitiesAggregationCalculatedFieldState(EntityId entityId) {
super(entityId); super(entityId);
} }
@ -71,8 +74,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec()); deduplicationIntervalMs = SECONDS.toMillis(configuration.getDeduplicationIntervalInSec());
} }
public void scheduleReevaluation() { @Override
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx); public void close() {
super.close();
if (reevaluationFuture != null) {
reevaluationFuture.cancel(true);
reevaluationFuture = null;
}
} }
@Override @Override
@ -142,6 +150,13 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
lastArgsRefreshTs = System.currentTimeMillis(); lastArgsRefreshTs = System.currentTimeMillis();
} }
public void scheduleReevaluation() {
ScheduledFuture<?> future = ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx);
if (future != null) {
reevaluationFuture = future;
}
}
@Override @Override
public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception { public ListenableFuture<CalculatedFieldResult> performCalculation(Map<String, ArgumentEntry> updatedArgs, CalculatedFieldCtx ctx) throws Exception {
boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty(); boolean cfUpdated = updatedArgs != null && updatedArgs.isEmpty();
@ -149,7 +164,7 @@ public class RelatedEntitiesAggregationCalculatedFieldState extends BaseCalculat
Output output = ctx.getOutput(); Output output = ctx.getOutput();
ObjectNode aggResult = aggregateMetrics(output); ObjectNode aggResult = aggregateMetrics(output);
lastMetricsEvalTs = System.currentTimeMillis(); lastMetricsEvalTs = System.currentTimeMillis();
ctx.scheduleReevaluation(deduplicationIntervalMs, actorCtx); scheduleReevaluation();
return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder() return Futures.immediateFuture(TelemetryCalculatedFieldResult.builder()
.type(output.getType()) .type(output.getType())
.scope(output.getScope()) .scope(output.getScope())

Loading…
Cancel
Save