Browse Source

fixed intervals

pull/14253/head
IrynaMatveieva 11 months ago
parent
commit
1de8e8cd35
  1. 4
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/AggIntervalEntry.java
  2. 22
      application/src/main/java/org/thingsboard/server/service/cf/ctx/state/aggregation/single/EntityAggregationCalculatedFieldState.java
  3. 8
      application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java
  4. 15
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java
  5. 114
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java
  6. 38
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java
  7. 18
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java
  8. 20
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java
  9. 18
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java
  10. 24
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java
  11. 21
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java
  12. 21
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java
  13. 22
      common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java
  14. 139
      common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalTest.java
  15. 2
      dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java

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

@ -29,4 +29,8 @@ public class AggIntervalEntry {
return ts >= startTs && ts < endTs; return ts >= startTs && ts < endTs;
} }
public long getIntervalDuration() {
return endTs - startTs;
}
} }

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

@ -37,6 +37,9 @@ import org.thingsboard.server.service.cf.ctx.state.ArgumentEntry;
import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState; import org.thingsboard.server.service.cf.ctx.state.BaseCalculatedFieldState;
import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx; import org.thingsboard.server.service.cf.ctx.state.CalculatedFieldCtx;
import java.time.Instant;
import java.time.ZoneId;
import java.time.ZonedDateTime;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Comparator; import java.util.Comparator;
import java.util.HashMap; import java.util.HashMap;
@ -49,7 +52,6 @@ import static org.thingsboard.server.utils.CalculatedFieldArgumentUtils.createDe
public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState { public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldState {
private AggInterval interval; private AggInterval interval;
private long intervalDuration;
private long watermarkDuration; private long watermarkDuration;
private long checkInterval; private long checkInterval;
private Map<String, AggMetric> metrics; private Map<String, AggMetric> metrics;
@ -65,7 +67,6 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
super.setCtx(ctx, actorCtx); super.setCtx(ctx, actorCtx);
this.cfProcessingService = ctx.getCfProcessingService(); this.cfProcessingService = ctx.getCfProcessingService();
var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration(); var configuration = (EntityAggregationCalculatedFieldConfiguration) ctx.getCalculatedField().getConfiguration();
intervalDuration = configuration.getInterval().getIntervalDurationMillis();
Watermark watermark = configuration.getWatermark(); Watermark watermark = configuration.getWatermark();
watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration()); watermarkDuration = watermark == null ? 0 : TimeUnit.SECONDS.toMillis(watermark.getDuration());
checkInterval = TimeUnit.SECONDS.toMillis(ctx.getCfCheckInterval()); checkInterval = TimeUnit.SECONDS.toMillis(ctx.getCfCheckInterval());
@ -127,18 +128,21 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
} }
public void fillMissingIntervals() { public void fillMissingIntervals() {
ZoneId zoneId = interval.getZoneId();
long currentIntervalEndTs = interval.getCurrentIntervalEndTs(); long currentIntervalEndTs = interval.getCurrentIntervalEndTs();
long intervalDuration = interval.getIntervalDurationMillis();
Map<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> intervals = getIntervals(); Map<AggIntervalEntry, Map<String, AggIntervalEntryStatus>> intervals = getIntervals();
AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null); AggIntervalEntry lastIntervalEntry = intervals.keySet().stream().max(Comparator.comparing(AggIntervalEntry::getEndTs)).orElse(null);
if (lastIntervalEntry == null) { if (lastIntervalEntry == null) {
return; return;
} }
long nextStartTs = lastIntervalEntry.getEndTs(); ZonedDateTime nextStart = Instant.ofEpochMilli(lastIntervalEntry.getEndTs()).atZone(zoneId);
long nextEndTs = nextStartTs + intervalDuration; ZonedDateTime nextEnd = interval.getNextIntervalStart(nextStart);
while (nextEndTs <= currentIntervalEndTs) { while (nextEnd.toInstant().toEpochMilli() <= currentIntervalEndTs) {
long nextStartTs = nextStart.toInstant().toEpochMilli();
long nextEndTs = nextEnd.toInstant().toEpochMilli();
AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs); AggIntervalEntry missing = new AggIntervalEntry(nextStartTs, nextEndTs);
arguments.forEach((argName, argumentEntry) -> { arguments.forEach((argName, argumentEntry) -> {
@ -147,8 +151,8 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
entityAggEntry.getAggIntervals().computeIfAbsent(missing, missingInterval -> intervalEntryStatus); entityAggEntry.getAggIntervals().computeIfAbsent(missing, missingInterval -> intervalEntryStatus);
}); });
nextStartTs = nextEndTs; nextStart = nextEnd;
nextEndTs += intervalDuration; nextEnd = interval.getNextIntervalStart(nextStart);
} }
} }
@ -174,7 +178,7 @@ public class EntityAggregationCalculatedFieldState extends BaseCalculatedFieldSt
if (now - endTs > watermarkDuration) { if (now - endTs > watermarkDuration) {
handleExpiredInterval(intervalEntry, args, results); handleExpiredInterval(intervalEntry, args, results);
expiredIntervals.add(intervalEntry); expiredIntervals.add(intervalEntry);
} else if (now - startTs >= intervalDuration) { } else if (now - startTs >= intervalEntry.getIntervalDuration()) {
handleActiveInterval(intervalEntry, args, results); handleActiveInterval(intervalEntry, args, results);
} }
} }

8
application/src/test/java/org/thingsboard/server/cf/EntityAggregationCalculatedFieldTest.java

@ -95,7 +95,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest
public void testCreateCf_checkAggregation() throws Exception { public void testCreateCf_checkAggregation() throws Exception {
Device device = createDevice("Device", "1234567890111"); Device device = createDevice("Device", "1234567890111");
CustomInterval customInterval = new CustomInterval(30L, 0L, "Europe/Kyiv"); CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L);
long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs();
long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs();
@ -108,7 +108,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2));
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3));
long interval = customInterval.getIntervalDurationMillis(); long interval = customInterval.getCurrentIntervalDurationMillis();
Watermark watermark = new Watermark(60); Watermark watermark = new Watermark(60);
CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark);
@ -126,7 +126,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest
public void testCreateCf_checkAggregationDuringWatermark() throws Exception { public void testCreateCf_checkAggregationDuringWatermark() throws Exception {
Device device = createDevice("Device", "1234567890111"); Device device = createDevice("Device", "1234567890111");
CustomInterval customInterval = new CustomInterval(30L, 0L, "Europe/Kyiv"); CustomInterval customInterval = new CustomInterval("Europe/Kyiv", 30L, 0L);
long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs(); long currentIntervalStartTs = customInterval.getCurrentIntervalStartTs();
long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs(); long currentIntervalEndTs = customInterval.getCurrentIntervalEndTs();
@ -139,7 +139,7 @@ public class EntityAggregationCalculatedFieldTest extends AbstractControllerTest
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":180}}", tsInInterval_2));
postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3)); postTelemetry(device.getId(), String.format("{\"ts\": \"%s\", \"values\": {\"energy\":120}}", tsInInterval_3));
long interval = customInterval.getIntervalDurationMillis(); long interval = customInterval.getCurrentIntervalDurationMillis();
Watermark watermark = new Watermark(60); Watermark watermark = new Watermark(60);
CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark); CalculatedField totalConsumptionCF = createTotalConsumptionCF(device.getId(), customInterval, watermark);

15
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggInterval.java

@ -20,6 +20,9 @@ import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import com.fasterxml.jackson.annotation.JsonSubTypes; import com.fasterxml.jackson.annotation.JsonSubTypes;
import com.fasterxml.jackson.annotation.JsonTypeInfo; import com.fasterxml.jackson.annotation.JsonTypeInfo;
import java.time.ZoneId;
import java.time.ZonedDateTime;
@JsonTypeInfo( @JsonTypeInfo(
use = JsonTypeInfo.Id.NAME, use = JsonTypeInfo.Id.NAME,
include = JsonTypeInfo.As.PROPERTY, include = JsonTypeInfo.As.PROPERTY,
@ -42,15 +45,21 @@ public interface AggInterval {
AggIntervalType getType(); AggIntervalType getType();
@JsonIgnore @JsonIgnore
long getIntervalDurationMillis(); ZoneId getZoneId();
@JsonIgnore
long getCurrentIntervalDurationMillis();
@JsonIgnore @JsonIgnore
long getCurrentIntervalStartTs(); long getCurrentIntervalStartTs();
long getDateTimeIntervalStartTs(ZonedDateTime dateTime);
@JsonIgnore @JsonIgnore
long getCurrentIntervalEndTs(); long getCurrentIntervalEndTs();
@JsonIgnore long getDateTimeIntervalEndTs(ZonedDateTime dateTime);
long getDelayUntilIntervalEnd();
ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart);
} }

114
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/BaseAggInterval.java

@ -15,54 +15,50 @@
*/ */
package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval;
import com.fasterxml.jackson.annotation.JsonIgnore;
import com.fasterxml.jackson.annotation.JsonInclude; import com.fasterxml.jackson.annotation.JsonInclude;
import jakarta.validation.constraints.Min;
import jakarta.validation.constraints.NotBlank; import jakarta.validation.constraints.NotBlank;
import jakarta.validation.constraints.NotNull; import lombok.AllArgsConstructor;
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor;
import java.time.DayOfWeek;
import java.time.Duration;
import java.time.LocalDate;
import java.time.LocalTime;
import java.time.ZoneId; import java.time.ZoneId;
import java.time.ZonedDateTime; import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
import java.time.temporal.TemporalAdjusters;
@Data @Data
@JsonInclude(JsonInclude.Include.NON_NULL) @JsonInclude(JsonInclude.Include.NON_NULL)
@AllArgsConstructor
@NoArgsConstructor
public abstract class BaseAggInterval implements AggInterval { public abstract class BaseAggInterval implements AggInterval {
@NotBlank @NotBlank
protected String tz; protected String tz;
protected Long offsetSec; // delay seconds since start of interval protected Long offsetSec; // delay seconds since start of interval
@JsonIgnore @Override
protected long getOffsetSec() { public ZoneId getZoneId() {
return ZoneId.of(tz);
}
protected long getOffset() {
return offsetSec != null ? offsetSec : 0L; return offsetSec != null ? offsetSec : 0L;
} }
@Override @Override
public long getIntervalDurationMillis() { public long getCurrentIntervalDurationMillis() {
return switch (getType()) { return getCurrentIntervalEndTs() - getCurrentIntervalStartTs();
case HOUR -> Duration.ofHours(1).toMillis();
case DAY -> Duration.ofDays(1).toMillis();
case WEEK, WEEK_SUN_SAT -> Duration.ofDays(7L).toMillis();
case MONTH -> Duration.ofDays(Math.round(30)).toMillis(); // average
case QUARTER -> Duration.ofDays(Math.round(91)).toMillis();
case YEAR -> Duration.ofDays(Math.round(365)).toMillis();
default -> throw new IllegalArgumentException("Unsupported type: " + getType());
};
} }
@Override @Override
public long getCurrentIntervalStartTs() { public long getCurrentIntervalStartTs() {
ZoneId zoneId = ZoneId.of(tz); ZoneId zoneId = getZoneId();
ZonedDateTime now = ZonedDateTime.now(zoneId); ZonedDateTime now = ZonedDateTime.now(zoneId);
long offset = getOffsetSec(); return getDateTimeIntervalStartTs(now);
ZonedDateTime shiftedNow = now.minusSeconds(offset); }
@Override
public long getDateTimeIntervalStartTs(ZonedDateTime dateTime) {
long offset = getOffset();
ZonedDateTime shiftedNow = dateTime.minusSeconds(offset);
ZonedDateTime alignedStart = getAlignedBoundary(shiftedNow, false); ZonedDateTime alignedStart = getAlignedBoundary(shiftedNow, false);
ZonedDateTime actualStart = alignedStart.plusSeconds(offset); ZonedDateTime actualStart = alignedStart.plusSeconds(offset);
return actualStart.toInstant().toEpochMilli(); return actualStart.toInstant().toEpochMilli();
@ -70,72 +66,20 @@ public abstract class BaseAggInterval implements AggInterval {
@Override @Override
public long getCurrentIntervalEndTs() { public long getCurrentIntervalEndTs() {
ZoneId zoneId = ZoneId.of(tz); ZoneId zoneId = getZoneId();
ZonedDateTime now = ZonedDateTime.now(zoneId); ZonedDateTime now = ZonedDateTime.now(zoneId);
long offset = getOffsetSec(); return getDateTimeIntervalEndTs(now);
ZonedDateTime shiftedNow = now.minusSeconds(offset);
ZonedDateTime alignedEnd = getAlignedBoundary(shiftedNow, true);
ZonedDateTime actualEnd = alignedEnd.plusSeconds(offset);
return actualEnd.toInstant().toEpochMilli();
} }
@Override @Override
public long getDelayUntilIntervalEnd() { public long getDateTimeIntervalEndTs(ZonedDateTime dateTime) {
long currentIntervalEndTs = getCurrentIntervalEndTs(); long offset = getOffset();
long now = System.currentTimeMillis(); ZonedDateTime shiftedNow = dateTime.minusSeconds(offset);
return currentIntervalEndTs - now; ZonedDateTime alignedEnd = getAlignedBoundary(shiftedNow, true);
} ZonedDateTime actualEnd = alignedEnd.plusSeconds(offset);
return actualEnd.toInstant().toEpochMilli();
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
return switch (getType()) {
case HOUR -> alignByHours(reference, next);
case DAY -> alignByDays(reference, next);
case WEEK -> alignByWeeks(reference, DayOfWeek.MONDAY, next);
case WEEK_SUN_SAT -> alignByWeeks(reference, DayOfWeek.SUNDAY, next);
case MONTH -> alignByMonths(reference, next);
case QUARTER -> alignByQuarters(reference, next);
case YEAR -> alignByYears(reference, next);
default -> throw new IllegalArgumentException("Unsupported interval type: " + getType());
};
}
private ZonedDateTime alignByHours(ZonedDateTime now, boolean next) {
ZonedDateTime base = now.truncatedTo(ChronoUnit.HOURS);
return next ? base.plusHours(1) : base;
}
private ZonedDateTime alignByDays(ZonedDateTime now, boolean next) {
ZonedDateTime base = now.truncatedTo(ChronoUnit.DAYS);
return next ? base.plusDays(1) : base;
}
private ZonedDateTime alignByWeeks(ZonedDateTime now, DayOfWeek startOfWeek, boolean next) {
ZonedDateTime startOfWeekDate = now.with(TemporalAdjusters.previousOrSame(startOfWeek))
.truncatedTo(ChronoUnit.DAYS);
return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate;
}
private ZonedDateTime alignByMonths(ZonedDateTime now, boolean next) {
ZonedDateTime base = now.withDayOfMonth(1).truncatedTo(ChronoUnit.DAYS);
return next ? base.plusMonths(1) : base;
}
private ZonedDateTime alignByQuarters(ZonedDateTime now, boolean next) {
int month = now.getMonthValue();
int quarterStartMonth = ((month - 1) / 3) * 3 + 1; // 1, 4, 7, 10
ZonedDateTime base = ZonedDateTime.of(
LocalDate.of(now.getYear(), quarterStartMonth, 1),
LocalTime.MIDNIGHT,
now.getZone());
return next ? base.plusMonths(3) : base;
} }
private ZonedDateTime alignByYears(ZonedDateTime now, boolean next) { protected abstract ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next);
ZonedDateTime base = ZonedDateTime.of(
LocalDate.of(now.getYear(), 1, 1),
LocalTime.MIDNIGHT,
now.getZone());
return next ? base.plusYears(1) : base;
}
} }

38
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/CustomInterval.java

@ -20,9 +20,8 @@ import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.Duration; import java.time.Duration;
import java.time.ZoneId; import java.time.Instant;
import java.time.ZonedDateTime; import java.time.ZonedDateTime;
import java.util.concurrent.TimeUnit;
@EqualsAndHashCode(callSuper = true) @EqualsAndHashCode(callSuper = true)
@Data @Data
@ -31,9 +30,8 @@ public class CustomInterval extends BaseAggInterval {
private Long durationSec; private Long durationSec;
public CustomInterval(Long durationSec, Long offsetMillis, String tz) { public CustomInterval(String tz, Long offsetSec, Long durationSec) {
this.tz = tz; super(tz, offsetSec);
this.offsetSec = offsetMillis;
this.durationSec = durationSec; this.durationSec = durationSec;
} }
@ -43,32 +41,26 @@ public class CustomInterval extends BaseAggInterval {
} }
@Override @Override
public long getIntervalDurationMillis() { public long getCurrentIntervalDurationMillis() {
return Duration.ofSeconds(durationSec).toMillis(); return getDurationMillis();
} }
@Override private long getDurationMillis() {
public long getCurrentIntervalStartTs() { return Duration.ofSeconds(durationSec).toMillis();
ZoneId zoneId = ZoneId.of(tz);
ZonedDateTime now = ZonedDateTime.now(zoneId);
ZonedDateTime shiftedNow = now.minusSeconds(getOffsetSec());
long durationMillis = getIntervalDurationMillis();
long shiftedNowMillis = shiftedNow.toInstant().toEpochMilli();
long alignedStartMillis = (shiftedNowMillis / durationMillis) * durationMillis;
long offsetMillis = TimeUnit.SECONDS.toMillis(getOffsetSec());
return alignedStartMillis + offsetMillis;
} }
@Override @Override
public long getCurrentIntervalEndTs() { protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
return getCurrentIntervalStartTs() + getIntervalDurationMillis(); long durationMillis = getDurationMillis();
long nowMillis = reference.toInstant().toEpochMilli();
long alignedStartMillis = (nowMillis / durationMillis) * durationMillis;
ZonedDateTime aligned = Instant.ofEpochMilli(alignedStartMillis).atZone(getZoneId());
return next ? aligned.plusSeconds(durationSec) : aligned;
} }
@Override @Override
public long getDelayUntilIntervalEnd() { public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return getCurrentIntervalEndTs() - System.currentTimeMillis(); return currentStart.plusSeconds(durationSec);
} }
} }

18
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/DayInterval.java

@ -18,6 +18,9 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class DayInterval extends BaseAggInterval { public class DayInterval extends BaseAggInterval {
@ -27,4 +30,19 @@ public class DayInterval extends BaseAggInterval {
return AggIntervalType.DAY; return AggIntervalType.DAY;
} }
public DayInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime base = reference.truncatedTo(ChronoUnit.DAYS);
return next ? base.plusDays(1) : base;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusDays(1);
}
} }

20
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/HourInterval.java

@ -16,15 +16,35 @@
package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval; package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
@EqualsAndHashCode(callSuper = true)
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class HourInterval extends BaseAggInterval { public class HourInterval extends BaseAggInterval {
public HourInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override @Override
public AggIntervalType getType() { public AggIntervalType getType() {
return AggIntervalType.HOUR; return AggIntervalType.HOUR;
} }
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime base = reference.truncatedTo(ChronoUnit.HOURS);
return next ? base.plusHours(1) : base;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusHours(1);
}
} }

18
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/MonthInterval.java

@ -18,6 +18,9 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class MonthInterval extends BaseAggInterval { public class MonthInterval extends BaseAggInterval {
@ -27,4 +30,19 @@ public class MonthInterval extends BaseAggInterval {
return AggIntervalType.MONTH; return AggIntervalType.MONTH;
} }
public MonthInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime base = reference.withDayOfMonth(1).truncatedTo(ChronoUnit.DAYS);
return next ? base.plusMonths(1) : base;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusMonths(1);
}
} }

24
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/QuarterInterval.java

@ -18,6 +18,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.LocalDate;
import java.time.LocalTime;
import java.time.ZonedDateTime;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class QuarterInterval extends BaseAggInterval { public class QuarterInterval extends BaseAggInterval {
@ -27,4 +31,24 @@ public class QuarterInterval extends BaseAggInterval {
return AggIntervalType.QUARTER; return AggIntervalType.QUARTER;
} }
public QuarterInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
int month = reference.getMonthValue();
int quarterStartMonth = ((month - 1) / 3) * 3 + 1; // 1, 4, 7, 10
ZonedDateTime base = ZonedDateTime.of(
LocalDate.of(reference.getYear(), quarterStartMonth, 1),
LocalTime.MIDNIGHT,
reference.getZone());
return next ? base.plusMonths(3) : base;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusMonths(3);
}
} }

21
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekInterval.java

@ -18,6 +18,11 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.DayOfWeek;
import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
import java.time.temporal.TemporalAdjusters;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class WeekInterval extends BaseAggInterval { public class WeekInterval extends BaseAggInterval {
@ -27,4 +32,20 @@ public class WeekInterval extends BaseAggInterval {
return AggIntervalType.WEEK; return AggIntervalType.WEEK;
} }
public WeekInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime startOfWeekDate = reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.MONDAY))
.truncatedTo(ChronoUnit.DAYS);
return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusWeeks(1);
}
} }

21
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/WeekSunSatInterval.java

@ -18,6 +18,11 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.DayOfWeek;
import java.time.ZonedDateTime;
import java.time.temporal.ChronoUnit;
import java.time.temporal.TemporalAdjusters;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class WeekSunSatInterval extends BaseAggInterval { public class WeekSunSatInterval extends BaseAggInterval {
@ -27,4 +32,20 @@ public class WeekSunSatInterval extends BaseAggInterval {
return AggIntervalType.WEEK_SUN_SAT; return AggIntervalType.WEEK_SUN_SAT;
} }
public WeekSunSatInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime startOfWeekDate = reference.with(TemporalAdjusters.previousOrSame(DayOfWeek.SUNDAY))
.truncatedTo(ChronoUnit.DAYS);
return next ? startOfWeekDate.plusWeeks(1) : startOfWeekDate;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusWeeks(1);
}
} }

22
common/data/src/main/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/YearInterval.java

@ -18,6 +18,10 @@ package org.thingsboard.server.common.data.cf.configuration.aggregation.single.i
import lombok.Data; import lombok.Data;
import lombok.NoArgsConstructor; import lombok.NoArgsConstructor;
import java.time.LocalDate;
import java.time.LocalTime;
import java.time.ZonedDateTime;
@Data @Data
@NoArgsConstructor @NoArgsConstructor
public class YearInterval extends BaseAggInterval { public class YearInterval extends BaseAggInterval {
@ -27,4 +31,22 @@ public class YearInterval extends BaseAggInterval {
return AggIntervalType.YEAR; return AggIntervalType.YEAR;
} }
public YearInterval(String tz, Long offsetSec) {
super(tz, offsetSec);
}
@Override
protected ZonedDateTime getAlignedBoundary(ZonedDateTime reference, boolean next) {
ZonedDateTime base = ZonedDateTime.of(
LocalDate.of(reference.getYear(), 1, 1),
LocalTime.MIDNIGHT,
reference.getZone());
return next ? base.plusYears(1) : base;
}
@Override
public ZonedDateTime getNextIntervalStart(ZonedDateTime currentStart) {
return currentStart.plusYears(1);
}
} }

139
common/data/src/test/java/org/thingsboard/server/common/data/cf/configuration/aggregation/single/interval/AggIntervalTest.java

@ -0,0 +1,139 @@
/**
* Copyright © 2016-2025 The Thingsboard Authors
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.common.data.cf.configuration.aggregation.single.interval;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;
import java.time.Duration;
import java.time.Instant;
import java.time.ZoneId;
import java.time.ZonedDateTime;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import java.util.function.LongFunction;
import java.util.stream.Stream;
import static org.assertj.core.api.Assertions.assertThat;
public class AggIntervalTest {
private static final String TZ = "Europe/Kiev";
@ParameterizedTest
@MethodSource("intervals")
void testGetStartAndEndWithoutOffset(LongFunction<AggInterval> intervalCreator) {
AggInterval interval = intervalCreator.apply(0L);
ZonedDateTime dateTime = ZonedDateTime.of(
// 2025.11.11 00:00:00
2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ)
);
long startTs = interval.getDateTimeIntervalStartTs(dateTime);
long endTs = interval.getDateTimeIntervalEndTs(dateTime);
assertThat(endTs).isGreaterThan(startTs);
assertThat(endTs - startTs).isEqualTo(interval.getCurrentIntervalDurationMillis());
}
@ParameterizedTest
@MethodSource("intervals")
void testApplyOffset(LongFunction<AggInterval> intervalCreator) {
long offsetSec = TimeUnit.MINUTES.toSeconds(15);
AggInterval intervalWithOffset = intervalCreator.apply(offsetSec);
AggInterval intervalNoOffset = intervalCreator.apply(0L);
ZonedDateTime dateTime = ZonedDateTime.of(
// 2025.11.11 11:20:00 - chosen so 15m offset shifts into a new interval
2025, 11, 11, 11, 20, 0, 0, ZoneId.of(TZ)
);
long startWithOffsetTs = intervalWithOffset.getDateTimeIntervalStartTs(dateTime);
long startNoOffsetTs = intervalNoOffset.getDateTimeIntervalStartTs(dateTime);
ZonedDateTime startWithOffset = Instant.ofEpochMilli(startWithOffsetTs).atZone(intervalWithOffset.getZoneId());
ZonedDateTime startNoOffset = Instant.ofEpochMilli(startNoOffsetTs).atZone(intervalNoOffset.getZoneId());
long actualOffset = Duration.between(startNoOffset, startWithOffset).toSeconds();
assertThat(actualOffset).isEqualTo(offsetSec);
}
private static Stream<Arguments> intervals() {
return Stream.of(
Arguments.of((LongFunction<AggInterval>) offset -> new HourInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new DayInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new WeekInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new WeekSunSatInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new MonthInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new QuarterInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new YearInterval(TZ, offset)),
Arguments.of((LongFunction<AggInterval>) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4)))
);
}
@ParameterizedTest
@MethodSource("nextIntervalFromExactDate")
void testNextIntervalFromExactDate(LongFunction<AggInterval> intervalCreator, Function<ZonedDateTime, ZonedDateTime> expectedDateTimeFunction) {
AggInterval interval = intervalCreator.apply(0L);
ZonedDateTime currentStart = ZonedDateTime.of(
2025, 11, 11, 0, 0, 0, 0, ZoneId.of(TZ)
);
ZonedDateTime nextStart = interval.getNextIntervalStart(currentStart);
assertThat(nextStart).isEqualTo(expectedDateTimeFunction.apply(currentStart));
}
private static Stream<Arguments> nextIntervalFromExactDate() {
return Stream.of(
Arguments.of(
(LongFunction<AggInterval>) offset -> new HourInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusHours(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new DayInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusDays(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new WeekInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusWeeks(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new WeekSunSatInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusWeeks(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new MonthInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusMonths(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new QuarterInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusMonths(3)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new YearInterval(TZ, offset),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusYears(1)
),
Arguments.of(
(LongFunction<AggInterval>) offset -> new CustomInterval(TZ, offset, TimeUnit.HOURS.toSeconds(4)),
(Function<ZonedDateTime, ZonedDateTime>) currentInterval -> currentInterval.plusHours(4)
)
);
}
}

2
dao/src/main/java/org/thingsboard/server/dao/service/validator/CalculatedFieldDataValidator.java

@ -138,7 +138,7 @@ public class CalculatedFieldDataValidator extends DataValidator<CalculatedField>
if (minAggregationIntervalInSec <= 0) { if (minAggregationIntervalInSec <= 0) {
return; return;
} }
if (aggConfiguration.getInterval().getIntervalDurationMillis() < TimeUnit.SECONDS.toMillis(minAggregationIntervalInSec)) { if (aggConfiguration.getInterval().getCurrentIntervalDurationMillis() < TimeUnit.SECONDS.toMillis(minAggregationIntervalInSec)) {
throw new IllegalArgumentException("Aggregation interval duration is less than configured " + throw new IllegalArgumentException("Aggregation interval duration is less than configured " +
"minimum allowed aggregation interval in tenant profile: " + minAggregationIntervalInSec + " sec."); "minimum allowed aggregation interval in tenant profile: " + minAggregationIntervalInSec + " sec.");
} }

Loading…
Cancel
Save