From e19ce5a823d65848509193568224346df1e73929 Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Wed, 4 May 2022 14:47:02 +0300 Subject: [PATCH] testCalculatePartitions for cassandra ts dao --- .../CassandraBaseTimeseriesDao.java | 11 +++-- .../CassandraBaseTimeseriesDaoTest.java | 41 +++++++++++++++++-- 2 files changed, 44 insertions(+), 8 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 639693acac..6521c88a6a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -82,8 +82,9 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD protected static final int MIN_AGGREGATION_STEP_MS = 1000; public static final String ASC_ORDER = "ASC"; public static final long SECONDS_IN_DAY = TimeUnit.DAYS.toSeconds(1); + static final long DAYS_32_MS = TimeUnit.DAYS.toMillis(32); - protected static List FIXED_PARTITION = Arrays.asList(new Long[]{0L}); + protected static final List FIXED_PARTITION = List.of(0L); private CassandraTsPartitionsCache cassandraTsPartitionsCache; @@ -341,7 +342,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD }, MoreExecutors.directExecutor()); } - private long toPartitionTs(long ts) { + long toPartitionTs(long ts) { LocalDateTime time = LocalDateTime.ofInstant(Instant.ofEpochMilli(ts), ZoneOffset.UTC); return tsFormat.truncatedTo(time).toInstant(ZoneOffset.UTC).toEpochMilli(); } @@ -432,13 +433,15 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD if (minPartition == maxPartition) { return Collections.singletonList(minPartition); } + List partitions = new ArrayList<>(); + partitions.add(minPartition); - List partitions = Arrays.asList(minPartition, maxPartition); long currentPartition = minPartition; - while (maxPartition > (currentPartition = toPartitionTs(currentPartition + TimeUnit.DAYS.toMillis(32)))){ + while (maxPartition > (currentPartition = toPartitionTs(currentPartition + DAYS_32_MS))){ partitions.add(currentPartition); } + partitions.add(maxPartition); return partitions; } diff --git a/dao/src/test/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDaoTest.java index c2ca768649..fd9736c96c 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDaoTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDaoTest.java @@ -34,22 +34,25 @@ import lombok.extern.slf4j.Slf4j; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Answers; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.boot.test.context.SpringBootTest; import org.springframework.boot.test.mock.mockito.MockBean; -import org.springframework.boot.test.mock.mockito.SpyBean; import org.springframework.test.context.TestPropertySource; import org.springframework.test.context.junit4.SpringRunner; import org.thingsboard.server.dao.cassandra.CassandraCluster; import org.thingsboard.server.dao.nosql.CassandraBufferedRateExecutor; +import java.text.ParseException; import java.util.List; +import static org.apache.commons.lang3.time.DateFormatUtils.ISO_DATETIME_TIME_ZONE_FORMAT; import static org.assertj.core.api.Assertions.assertThat; @RunWith(SpringRunner.class) @SpringBootTest(classes = CassandraBaseTimeseriesDao.class) @TestPropertySource(properties = { + "database.ts.type=cassandra", "cassandra.query.ts_key_value_partitioning=MONTHS", "cassandra.query.ts_key_value_partitioning_always_exist_in_reading=true", "cassandra.query.ts_key_value_partitions_max_cache_size=100000", @@ -61,7 +64,7 @@ import static org.assertj.core.api.Assertions.assertThat; @Slf4j public class CassandraBaseTimeseriesDaoTest { - @SpyBean + @Autowired CassandraBaseTimeseriesDao tsDao; @MockBean(answer = Answers.RETURNS_MOCKS) @@ -71,8 +74,38 @@ public class CassandraBaseTimeseriesDaoTest { CassandraBufferedRateExecutor cassandraBufferedRateExecutor; @Test - public void testCalculatePartitions() { + public void testToPartitionsMonths() throws ParseException { + assertThat(tsDao.getPartitioning()).isEqualTo("MONTHS"); + assertThat(tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2022-01-01T00:00:00Z").getTime())).isEqualTo(1640995200000L); + assertThat(tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2022-05-01T00:00:00Z").getTime())).isEqualTo(1651363200000L); + assertThat(tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2022-05-01T00:00:01Z").getTime())).isEqualTo(1651363200000L); + assertThat(tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2022-05-31T23:59:59Z").getTime())).isEqualTo(1651363200000L); + assertThat(tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2023-12-31T23:59:59Z").getTime())).isEqualTo(1701388800000L); + } + + @Test + public void testCalculatePartitions() throws ParseException { + long startTs = tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2019-12-12T00:00:00Z").getTime()); + long nextTs = tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2020-01-31T23:59:59Z").getTime()); + long leapTs = tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2020-02-29T23:59:59Z").getTime()); + long endTs = tsDao.toPartitionTs(ISO_DATETIME_TIME_ZONE_FORMAT.parse("2021-01-31T23:59:59Z").getTime()); + + log.warn("startTs {}, nextTs {}, leapTs {}, endTs {}", startTs, nextTs, leapTs, endTs); + assertThat(tsDao.calculatePartitions(0, 0)).isEqualTo(List.of(0L)); assertThat(tsDao.calculatePartitions(0, 1)).isEqualTo(List.of(0L, 1L)); + assertThat(tsDao.calculatePartitions(startTs, startTs)).isEqualTo(List.of(1575158400000L)); + assertThat(tsDao.calculatePartitions(startTs, nextTs)).isEqualTo(List.of(1575158400000L, 1577836800000L)); + assertThat(tsDao.calculatePartitions(startTs, leapTs)).isEqualTo(List.of(1575158400000L, 1577836800000L, 1580515200000L)); + + assertThat(tsDao.calculatePartitions(startTs, endTs)).hasSize(14); + assertThat(tsDao.calculatePartitions(startTs, endTs)).isEqualTo(List.of( + 1575158400000L, + 1577836800000L, 1580515200000L, 1583020800000L, + 1585699200000L, 1588291200000L, 1590969600000L, + 1593561600000L, 1596240000000L, 1598918400000L, + 1601510400000L, 1604188800000L, 1606780800000L, + 1609459200000L)); } -} \ No newline at end of file + +}