Browse Source

testCalculatePartitions for cassandra ts dao

pull/7629/head
Sergey Matvienko 4 years ago
parent
commit
e19ce5a823
  1. 11
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  2. 41
      dao/src/test/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDaoTest.java

11
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<Long> FIXED_PARTITION = Arrays.asList(new Long[]{0L});
protected static final List<Long> 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<Long> partitions = new ArrayList<>();
partitions.add(minPartition);
List<Long> 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;
}

41
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));
}
}
}

Loading…
Cancel
Save