Browse Source

Fix reevaluation interval for entity aggregation CFs

pull/14391/head
Viacheslav Klimov 10 months ago
parent
commit
2e8fe17871
  1. 11
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/CalculatedFieldCtx.java
  2. 7
      dao/src/main/java/org/thingsboard/server/dao/util/TimeUtils.java

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

@ -58,6 +58,7 @@ import org.thingsboard.server.common.data.tenant.profile.DefaultTenantProfileCon
import org.thingsboard.server.common.data.util.CollectionsUtil; import org.thingsboard.server.common.data.util.CollectionsUtil;
import org.thingsboard.server.common.util.ProtoUtils; import org.thingsboard.server.common.util.ProtoUtils;
import org.thingsboard.server.dao.relation.RelationService; import org.thingsboard.server.dao.relation.RelationService;
import org.thingsboard.server.dao.util.TimeUtils;
import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto; import org.thingsboard.server.gen.transport.TransportProtos.CalculatedFieldTelemetryMsgProto;
import org.thingsboard.server.service.cf.CalculatedFieldProcessingService; import org.thingsboard.server.service.cf.CalculatedFieldProcessingService;
import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId; import org.thingsboard.server.service.cf.ctx.CalculatedFieldEntityCtxId;
@ -66,6 +67,7 @@ import org.thingsboard.server.service.cf.ctx.state.geofencing.GeofencingCalculat
import org.thingsboard.server.service.telemetry.AlarmSubscriptionService; import org.thingsboard.server.service.telemetry.AlarmSubscriptionService;
import java.io.Closeable; import java.io.Closeable;
import java.time.ZonedDateTime;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.HashMap; import java.util.HashMap;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
@ -215,8 +217,13 @@ public class CalculatedFieldCtx implements Closeable {
if (watermark != null && watermark.getDuration() > 0) { if (watermark != null && watermark.getDuration() > 0) {
return true; return true;
} }
long intervalDurationMillis = entityAggregationConfig.getInterval().getCurrentIntervalDurationMillis(); if (lastReevaluationTs == 0) {
if (now - lastReevaluationTs >= intervalDurationMillis) { lastReevaluationTs = now;
return true;
}
ZonedDateTime lastReevaluationTime = TimeUtils.toZonedDateTime(lastReevaluationTs, entityAggregationConfig.getInterval().getZoneId());
long previousIntervalEndTs = entityAggregationConfig.getInterval().getDateTimeIntervalEndTs(lastReevaluationTime);
if (now >= previousIntervalEndTs) {
lastReevaluationTs = now; lastReevaluationTs = now;
return true; return true;
} }

7
dao/src/main/java/org/thingsboard/server/dao/util/TimeUtils.java

@ -15,6 +15,8 @@
*/ */
package org.thingsboard.server.dao.util; package org.thingsboard.server.dao.util;
import lombok.AccessLevel;
import lombok.NoArgsConstructor;
import org.thingsboard.server.common.data.kv.IntervalType; import org.thingsboard.server.common.data.kv.IntervalType;
import java.time.Instant; import java.time.Instant;
@ -24,6 +26,7 @@ import java.time.temporal.ChronoUnit;
import java.time.temporal.IsoFields; import java.time.temporal.IsoFields;
import java.time.temporal.WeekFields; import java.time.temporal.WeekFields;
@NoArgsConstructor(access = AccessLevel.PRIVATE)
public class TimeUtils { public class TimeUtils {
public static long calculateIntervalEnd(long startTs, IntervalType intervalType, ZoneId tzId) { public static long calculateIntervalEnd(long startTs, IntervalType intervalType, ZoneId tzId) {
@ -42,4 +45,8 @@ public class TimeUtils {
} }
} }
public static ZonedDateTime toZonedDateTime(long ts, ZoneId zoneId) {
return ZonedDateTime.ofInstant(Instant.ofEpochMilli(ts), zoneId);
}
} }

Loading…
Cancel
Save