From 0de5868bc5886e3c2acc2cc1c56dba8dd0fa9a7b Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 12:36:35 +0200 Subject: [PATCH] added Caffeine cache for Cassandra ts partitions saving --- .../src/main/resources/thingsboard.yml | 3 +- .../CassandraBaseTimeseriesDao.java | 46 ++++++ .../CassandraPartitionCacheKey.java | 30 ++++ .../CassandraTsPartitionsCache.java | 42 ++++++ .../nosql/CassandraPartitionsCacheTest.java | 132 ++++++++++++++++++ .../test/resources/cassandra-test.properties | 2 + 6 files changed, 254 insertions(+), 1 deletion(-) create mode 100644 dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java create mode 100644 dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java create mode 100644 dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 4cefe63e54..bf5ca8023d 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -223,8 +223,9 @@ cassandra: read_consistency_level: "${CASSANDRA_READ_CONSISTENCY_LEVEL:ONE}" write_consistency_level: "${CASSANDRA_WRITE_CONSISTENCY_LEVEL:ONE}" default_fetch_size: "${CASSANDRA_DEFAULT_FETCH_SIZE:2000}" - # Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS,INDEFINITE + # Specify partitioning size for timestamp key-value storage. Example: MINUTES, HOURS, DAYS, MONTHS, INDEFINITE ts_key_value_partitioning: "${TS_KV_PARTITIONING:MONTHS}" + ts_key_value_partitions_max_cache_size: "${TS_KV_PARTITIONS_MAX_CACHE_SIZE:100000}" ts_key_value_ttl: "${TS_KV_TTL:0}" events_ttl: "${TS_EVENTS_TTL:0}" # Specify TTL of debug log in seconds. The current value corresponds to one week 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 b96462e350..d39c6c5afb 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 @@ -29,6 +29,7 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.MoreExecutors; +import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.springframework.beans.factory.annotation.Autowired; @@ -66,6 +67,7 @@ import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Optional; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -88,12 +90,17 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem public static final String DESC_ORDER = "DESC"; private static List FIXED_PARTITION = Arrays.asList(new Long[]{0L}); + private CassandraTsPartitionsCache cassandraTsPartitionsCache; + @Autowired private Environment environment; @Value("${cassandra.query.ts_key_value_partitioning}") private String partitioning; + @Value("${cassandra.query.ts_key_value_partitions_max_cache_size}") + private long partitionsCacheSize; + @Value("${cassandra.query.ts_key_value_ttl}") private long systemTtl; @@ -126,6 +133,9 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem Optional partition = NoSqlTsPartitionDate.parse(partitioning); if (partition.isPresent()) { tsFormat = partition.get(); + if (!isFixedPartitioning() && partitionsCacheSize > 0) { + cassandraTsPartitionsCache = new CassandraTsPartitionsCache(partitionsCacheSize); + } } else { log.warn("Incorrect configuration of partitioning {}", partitioning); throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); @@ -390,6 +400,42 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } ttl = computeTtl(ttl); long partition = toPartitionTs(tsKvEntryTs); + if (cassandraTsPartitionsCache == null) { + return doSavePartition(tenantId, entityId, key, ttl, partition); + } else { + CassandraPartitionCacheKey partitionSearchKey = new CassandraPartitionCacheKey(entityId, key, partition); + CompletableFuture hasFuture = cassandraTsPartitionsCache.has(partitionSearchKey); + SettableFuture listenableFuture = SettableFuture.create(); + if (hasFuture == null) { + return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, ttl); + } else { + hasFuture.whenComplete((result, throwable) -> { + if (throwable != null) { + listenableFuture.setException(throwable); + } else { + listenableFuture.set(result); + } + }); + long finalTtl = ttl; + return Futures.transformAsync(listenableFuture, result -> { + if (result) { + return Futures.immediateFuture(null); + } else { + return processDoSavePartition(tenantId, entityId, key, partition, partitionSearchKey, finalTtl); + } + }, readResultsProcessingExecutor); + } + } + } + + private ListenableFuture processDoSavePartition(TenantId tenantId, EntityId entityId, String key, long partition, CassandraPartitionCacheKey partitionSearchKey, long ttl) { + return Futures.transformAsync(doSavePartition(tenantId, entityId, key, ttl, partition), input -> { + cassandraTsPartitionsCache.put(partitionSearchKey); + return Futures.immediateFuture(input); + }, readResultsProcessingExecutor); + } + + private ListenableFuture doSavePartition(TenantId tenantId, EntityId entityId, String key, long ttl, long partition) { log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); BoundStatement stmt = (ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind(); stmt = stmt.setString(0, entityId.getEntityType().name()) diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java new file mode 100644 index 0000000000..791ce84113 --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraPartitionCacheKey.java @@ -0,0 +1,30 @@ +/** + * Copyright © 2016-2020 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.dao.timeseries; + +import lombok.AllArgsConstructor; +import lombok.Data; +import org.thingsboard.server.common.data.id.EntityId; + +@Data +@AllArgsConstructor +public class CassandraPartitionCacheKey { + + private EntityId entityId; + private String key; + private long partition; + +} \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java new file mode 100644 index 0000000000..b467b5446f --- /dev/null +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraTsPartitionsCache.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2020 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.dao.timeseries; + +import com.github.benmanes.caffeine.cache.AsyncLoadingCache; +import com.github.benmanes.caffeine.cache.Caffeine; + +import java.util.concurrent.CompletableFuture; + +public class CassandraTsPartitionsCache { + + private AsyncLoadingCache partitionsCache; + + public CassandraTsPartitionsCache(long maxCacheSize) { + this.partitionsCache = Caffeine.newBuilder() + .maximumSize(maxCacheSize) + .buildAsync(key -> { + throw new IllegalStateException("'get' methods calls are not supported!"); + }); + } + + public CompletableFuture has(CassandraPartitionCacheKey key) { + return partitionsCache.getIfPresent(key); + } + + public void put(CassandraPartitionCacheKey key) { + partitionsCache.put(key, CompletableFuture.completedFuture(true)); + } +} \ No newline at end of file diff --git a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java new file mode 100644 index 0000000000..82c0e0d6d2 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -0,0 +1,132 @@ +/** + * Copyright © 2016-2020 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.dao.nosql; + +import com.datastax.driver.core.BoundStatement; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.CodecRegistry; +import com.datastax.driver.core.Configuration; +import com.datastax.driver.core.ConsistencyLevel; +import com.datastax.driver.core.PreparedStatement; +import com.datastax.driver.core.ResultSetFuture; +import com.datastax.driver.core.Session; +import com.datastax.driver.core.Statement; +import com.google.common.util.concurrent.Futures; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mock; +import org.mockito.runners.MockitoJUnitRunner; +import org.springframework.core.env.Environment; +import org.springframework.test.util.ReflectionTestUtils; +import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.dao.cassandra.CassandraCluster; +import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao; + +import java.util.UUID; + +import static org.mockito.Matchers.any; +import static org.mockito.Matchers.anyInt; +import static org.mockito.Matchers.anyLong; +import static org.mockito.Matchers.anyString; +import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +@RunWith(MockitoJUnitRunner.class) +public class CassandraPartitionsCacheTest { + + private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao; + + @Mock + private Environment environment; + + @Mock + private CassandraBufferedRateExecutor rateLimiter; + + @Mock + private CassandraCluster cluster; + + @Mock + private Session session; + + @Mock + private Cluster sessionCluster; + + @Mock + private Configuration configuration; + + @Mock + private PreparedStatement preparedStatement; + + @Mock + private BoundStatement boundStatement; + + @Before + public void setUp() { + when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getSession()).thenReturn(session); + when(session.getCluster()).thenReturn(sessionCluster); + when(sessionCluster.getConfiguration()).thenReturn(configuration); + when(configuration.getCodecRegistry()).thenReturn(CodecRegistry.DEFAULT_INSTANCE); + when(session.prepare(anyString())).thenReturn(preparedStatement); + when(preparedStatement.bind()).thenReturn(boundStatement); + when(boundStatement.setString(anyInt(), anyString())).thenReturn(boundStatement); + when(boundStatement.setUUID(anyInt(), any(UUID.class))).thenReturn(boundStatement); + when(boundStatement.setLong(anyInt(), anyLong())).thenReturn(boundStatement); + when(boundStatement.setInt(anyInt(), anyInt())).thenReturn(boundStatement); + + cassandraBaseTimeseriesDao = spy(new CassandraBaseTimeseriesDao()); + + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitioning", "MONTHS"); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitionsCacheSize", 100000); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "systemTtl", 0); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "setNullValuesEnabled", false); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "environment", environment); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "rateLimiter", rateLimiter); + ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "cluster", cluster); + + doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(ResultSetFuture.class), any()); + + } + + @Test + public void testPartitionSave() throws Exception { + + cassandraBaseTimeseriesDao.init(); + + + UUID id = UUID.randomUUID(); + TenantId tenantId = new TenantId(id); + long tsKvEntryTs = System.currentTimeMillis(); + + for (int i = 0; i < 50000; i++) { + cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); + } + + for (int i = 0; i < 60000; i++) { + cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test" + i, 0); + } + + verify(cassandraBaseTimeseriesDao, times(60000)).executeAsyncWrite(any(TenantId.class), any(Statement.class)); + + + } + +} \ No newline at end of file diff --git a/dao/src/test/resources/cassandra-test.properties b/dao/src/test/resources/cassandra-test.properties index 51f34a08d6..2bf57b190f 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -46,6 +46,8 @@ cassandra.query.default_fetch_size=2000 cassandra.query.ts_key_value_partitioning=HOURS +cassandra.query.ts_key_value_partitions_max_cache_size=100000 + cassandra.query.ts_key_value_ttl=0 cassandra.query.debug_events_ttl=604800