From 4babb07fd6cfaf55dd1ba438254c4bf28d3bb501 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 11 Nov 2020 16:22:57 +0200 Subject: [PATCH 1/8] added # filter topic handling --- .../transport/mqtt/util/MqttTopicFilterFactory.java | 13 +++++++++---- .../mqtt/util/MqttTopicFilterFactoryTest.java | 11 ++++++++++- 2 files changed, 19 insertions(+), 5 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java index 4d5a9a7c2b..9893f8cbef 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java @@ -34,10 +34,15 @@ public class MqttTopicFilterFactory { } return filters.computeIfAbsent(topicFilter, filter -> { if (filter.contains("+") || filter.contains("#")) { - String regex = filter - .replace("\\", "\\\\") - .replace("+", "[^/]+") - .replace("/#", "($|/.*)"); + String regex; + if (filter.equals("#")) { + regex = filter.replace("#", "^(?!/).+"); + } else { + regex = filter + .replace("\\", "\\\\") + .replace("+", "[^/]+") + .replace("/#", "($|/.*)"); + } log.debug("Converting [{}] to [{}]", filter, regex); return new RegexTopicFilter(regex); } else { diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java index 0b854d51ef..f3a65bda14 100644 --- a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java +++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java @@ -20,7 +20,6 @@ import org.junit.runner.RunWith; import org.mockito.runners.MockitoJUnitRunner; import javax.script.ScriptException; -import java.util.regex.Pattern; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertTrue; @@ -32,6 +31,9 @@ public class MqttTopicFilterFactoryTest { private static String TEST_STR_2 = "Sensor/Temperature"; private static String TEST_STR_3 = "Sensor/Temperature2/House/48"; + private static String TEST_STR_4 = String.format("%s%n%s", "/Sensor/Temperature", "/House/48"); + private static String TEST_STR_5 = "/" + TEST_STR_1; + @Test public void metadataCanBeUpdated() throws ScriptException { MqttTopicFilter filter = MqttTopicFilterFactory.toFilter("Sensor/Temperature/House/+"); @@ -51,6 +53,13 @@ public class MqttTopicFilterFactoryTest { assertTrue(filter.filter(TEST_STR_1)); assertTrue(filter.filter(TEST_STR_2)); assertFalse(filter.filter(TEST_STR_3)); + + filter = MqttTopicFilterFactory.toFilter("#"); + assertTrue(filter.filter(TEST_STR_1)); + assertTrue(filter.filter(TEST_STR_2)); + assertTrue(filter.filter(TEST_STR_3)); + assertFalse(filter.filter(TEST_STR_4)); + assertFalse(filter.filter(TEST_STR_5)); } } From 4ee38a15d3c658dc6b276af95f6290f90f1f0dcf Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 11 Nov 2020 19:31:21 +0200 Subject: [PATCH 2/8] change regex for # filter --- .../transport/mqtt/util/MqttTopicFilterFactory.java | 2 +- .../mqtt/util/MqttTopicFilterFactoryTest.java | 11 +++++++---- 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java index 9893f8cbef..98e472ba5c 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java @@ -36,7 +36,7 @@ public class MqttTopicFilterFactory { if (filter.contains("+") || filter.contains("#")) { String regex; if (filter.equals("#")) { - regex = filter.replace("#", "^(?!/).+"); + regex = filter.replace("#", "\\S+"); } else { regex = filter .replace("\\", "\\\\") diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java index f3a65bda14..2ec05fac78 100644 --- a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java +++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java @@ -30,9 +30,10 @@ public class MqttTopicFilterFactoryTest { private static String TEST_STR_1 = "Sensor/Temperature/House/48"; private static String TEST_STR_2 = "Sensor/Temperature"; private static String TEST_STR_3 = "Sensor/Temperature2/House/48"; - - private static String TEST_STR_4 = String.format("%s%n%s", "/Sensor/Temperature", "/House/48"); - private static String TEST_STR_5 = "/" + TEST_STR_1; + private static String TEST_STR_4 = "/Sensor/Temperature2/House/48"; + private static String TEST_STR_5 = String.format("%s%n%s", "/Sensor/Temperature", "/House/48"); + private static String TEST_STR_6 = ""; + private static String TEST_STR_7 = " "; @Test public void metadataCanBeUpdated() throws ScriptException { @@ -58,8 +59,10 @@ public class MqttTopicFilterFactoryTest { assertTrue(filter.filter(TEST_STR_1)); assertTrue(filter.filter(TEST_STR_2)); assertTrue(filter.filter(TEST_STR_3)); - assertFalse(filter.filter(TEST_STR_4)); + assertTrue(filter.filter(TEST_STR_4)); assertFalse(filter.filter(TEST_STR_5)); + assertFalse(filter.filter(TEST_STR_6)); + assertFalse(filter.filter(TEST_STR_7)); } } From 3127678d17788d8e8101bda0921e4af78e2b56cb Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 16 Nov 2020 16:52:24 +0200 Subject: [PATCH 3/8] added AlwaysTrueTopicFilter --- .../mqtt/util/AlwaysTrueTopicFilter.java | 27 +++++++++++++++++++ .../mqtt/util/MqttTopicFilterFactory.java | 2 +- .../mqtt/util/MqttTopicFilterFactoryTest.java | 13 ++++----- 3 files changed, 35 insertions(+), 7 deletions(-) create mode 100644 common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/AlwaysTrueTopicFilter.java diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/AlwaysTrueTopicFilter.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/AlwaysTrueTopicFilter.java new file mode 100644 index 0000000000..9952c4b507 --- /dev/null +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/AlwaysTrueTopicFilter.java @@ -0,0 +1,27 @@ +/** + * 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.transport.mqtt.util; + +import lombok.Data; + +@Data +public class AlwaysTrueTopicFilter implements MqttTopicFilter { + + @Override + public boolean filter(String topic) { + return true; + } +} diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java index 98e472ba5c..f2b4ef27e2 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java @@ -36,7 +36,7 @@ public class MqttTopicFilterFactory { if (filter.contains("+") || filter.contains("#")) { String regex; if (filter.equals("#")) { - regex = filter.replace("#", "\\S+"); + return new AlwaysTrueTopicFilter(); } else { regex = filter .replace("\\", "\\\\") diff --git a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java index 2ec05fac78..fac2e5c01d 100644 --- a/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java +++ b/common/transport/mqtt/src/test/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactoryTest.java @@ -31,9 +31,8 @@ public class MqttTopicFilterFactoryTest { private static String TEST_STR_2 = "Sensor/Temperature"; private static String TEST_STR_3 = "Sensor/Temperature2/House/48"; private static String TEST_STR_4 = "/Sensor/Temperature2/House/48"; - private static String TEST_STR_5 = String.format("%s%n%s", "/Sensor/Temperature", "/House/48"); - private static String TEST_STR_6 = ""; - private static String TEST_STR_7 = " "; + private static String TEST_STR_5 = "Sensor/ Temperature"; + private static String TEST_STR_6 = "/"; @Test public void metadataCanBeUpdated() throws ScriptException { @@ -60,9 +59,11 @@ public class MqttTopicFilterFactoryTest { assertTrue(filter.filter(TEST_STR_2)); assertTrue(filter.filter(TEST_STR_3)); assertTrue(filter.filter(TEST_STR_4)); - assertFalse(filter.filter(TEST_STR_5)); - assertFalse(filter.filter(TEST_STR_6)); - assertFalse(filter.filter(TEST_STR_7)); + assertTrue(filter.filter(TEST_STR_5)); + assertTrue(filter.filter(TEST_STR_6)); + + filter = MqttTopicFilterFactory.toFilter("Sensor/Temperature#"); + assertFalse(filter.filter(TEST_STR_2)); } } From 22c91039553c1ef31c86a3fdfd437ff115e04320 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 16 Nov 2020 16:55:24 +0200 Subject: [PATCH 4/8] fix MqttTopicFilterFactory toFilter --- .../mqtt/util/MqttTopicFilterFactory.java | 17 +++++++---------- 1 file changed, 7 insertions(+), 10 deletions(-) diff --git a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java index f2b4ef27e2..0c3b497591 100644 --- a/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java +++ b/common/transport/mqtt/src/main/java/org/thingsboard/server/transport/mqtt/util/MqttTopicFilterFactory.java @@ -33,16 +33,13 @@ public class MqttTopicFilterFactory { throw new IllegalArgumentException("Topic filter can't be empty!"); } return filters.computeIfAbsent(topicFilter, filter -> { - if (filter.contains("+") || filter.contains("#")) { - String regex; - if (filter.equals("#")) { - return new AlwaysTrueTopicFilter(); - } else { - regex = filter - .replace("\\", "\\\\") - .replace("+", "[^/]+") - .replace("/#", "($|/.*)"); - } + if (filter.equals("#")) { + return new AlwaysTrueTopicFilter(); + } else if (filter.contains("+") || filter.contains("#")) { + String regex = filter + .replace("\\", "\\\\") + .replace("+", "[^/]+") + .replace("/#", "($|/.*)"); log.debug("Converting [{}] to [{}]", filter, regex); return new RegexTopicFilter(regex); } else { From e516cd31dc411acf84b38c22211120b84ee97aef Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 12:36:35 +0200 Subject: [PATCH 5/8] added Caffeine cache for Cassandra ts partitions saving --- .../src/main/resources/thingsboard.yml | 3 +- .../CassandraBaseTimeseriesDao.java | 104 ++++++++++++----- .../CassandraPartitionCacheKey.java | 30 +++++ .../CassandraTsPartitionsCache.java | 42 +++++++ .../nosql/CassandraPartitionsCacheTest.java | 106 ++++++++++++++++++ .../test/resources/cassandra-test.properties | 2 + 6 files changed, 259 insertions(+), 28 deletions(-) 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 791b032352..4a7c58e4db 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -192,8 +192,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 b9d8c62833..e4d22e7fef 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 @@ -79,12 +79,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD protected 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:100000}") + private long partitionsCacheSize; + @Value("${cassandra.query.ts_key_value_ttl}") private long systemTtl; @@ -111,13 +116,16 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD super.startExecutor(); if (!isInstall()) { getFetchStmt(Aggregation.NONE, DESC_ORDER); - } - Optional partition = NoSqlTsPartitionDate.parse(partitioning); - if (partition.isPresent()) { - tsFormat = partition.get(); - } else { - log.warn("Incorrect configuration of partitioning {}", partitioning); - throw new RuntimeException("Failed to parse partitioning property: " + partitioning + "!"); + 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 + "!"); + } } } @@ -168,26 +176,6 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return Futures.transform(Futures.allAsList(futures), result -> dataPointDays, MoreExecutors.directExecutor()); } - @Override - public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { - if (isFixedPartitioning()) { - return Futures.immediateFuture(null); - } - ttl = computeTtl(ttl); - long partition = toPartitionTs(tsKvEntryTs); - log.debug("Saving partition {} for the entity [{}-{}] and key {}", partition, entityId.getEntityType(), entityId.getId(), key); - BoundStatementBuilder stmtBuilder = new BoundStatementBuilder((ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind()); - stmtBuilder.setString(0, entityId.getEntityType().name()) - .setUuid(1, entityId.getId()) - .setLong(2, partition) - .setString(3, key); - if (ttl > 0) { - stmtBuilder.setInt(4, (int) ttl); - } - BoundStatement stmt = stmtBuilder.build(); - return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0); - } - @Override public ListenableFuture remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { long minPartition = toPartitionTs(query.getStartTs()); @@ -461,6 +449,68 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); } + @Override + public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { + if (isFixedPartitioning()) { + return Futures.immediateFuture(null); + } + 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); + if (!cassandraTsPartitionsCache.has(partitionSearchKey)) { + ListenableFuture result = doSavePartition(tenantId, entityId, key, ttl, partition); + Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor()); + return result; + } else { + return Futures.immediateFuture(0); + } + } + } + + 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); + PreparedStatement preparedStatement = ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt(); + BoundStatement stmt = preparedStatement.bind(); + stmt.setString(0, entityId.getEntityType().name()); + stmt.setUuid(1, entityId.getId()); + stmt.setLong(2, partition); + stmt.setString(3, key); + if (ttl > 0) { + stmt.setInt(4, (int) ttl); + } +// BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(bind); +// stmtBuilder.setString(0, entityId.getEntityType().name()) +// .setUuid(1, entityId.getId()) +// .setLong(2, partition) +// .setString(3, key); +// if (ttl > 0) { +// stmtBuilder.setInt(4, (int) ttl); +// } +// BoundStatement stmt = stmtBuilder.build(); + return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0); + } + + private class CacheCallback implements FutureCallback { + private final CassandraPartitionCacheKey key; + + private CacheCallback(CassandraPartitionCacheKey key) { + this.key = key; + } + + @Override + public void onSuccess(Void result) { + cassandraTsPartitionsCache.put(key); + } + + @Override + public void onFailure(Throwable t) { + + } + } + private long computeTtl(long ttl) { if (systemTtl > 0) { if (ttl == 0) { 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..c167fb28cb --- /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 boolean has(CassandraPartitionCacheKey key) { + return partitionsCache.getIfPresent(key) != null; + } + + public void put(CassandraPartitionCacheKey key) { + partitionsCache.put(key, CompletableFuture.completedFuture(true)); + } +} 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..4a607d0f03 --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -0,0 +1,106 @@ +/** + * 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.oss.driver.api.core.ConsistencyLevel; +import com.datastax.oss.driver.api.core.cql.BoundStatement; +import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.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.cassandra.guava.GuavaSession; +import org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao; + +import java.util.UUID; + +import static org.mockito.Matchers.any; +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 GuavaSession session; + + @Mock + private PreparedStatement preparedStatement; + + @Mock + private BoundStatement boundStatement; + + @Before + public void setUp() throws Exception { + when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getSession()).thenReturn(session); + when(session.prepare(anyString())).thenReturn(preparedStatement); + when(preparedStatement.bind()).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(TbResultSetFuture.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 43a78abac4..4b3ea0a74d 100644 --- a/dao/src/test/resources/cassandra-test.properties +++ b/dao/src/test/resources/cassandra-test.properties @@ -54,6 +54,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 From 24a31d61349e6a3e32c71fc149711822202ee281 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Mon, 21 Dec 2020 21:43:09 +0200 Subject: [PATCH 6/8] update code --- .../CassandraBaseTimeseriesDao.java | 65 ++++++++----------- .../nosql/CassandraPartitionsCacheTest.java | 19 +++--- 2 files changed, 39 insertions(+), 45 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 e4d22e7fef..acd140c0b6 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 @@ -176,6 +176,27 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return Futures.transform(Futures.allAsList(futures), result -> dataPointDays, MoreExecutors.directExecutor()); } + @Override + public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { + if (isFixedPartitioning()) { + return Futures.immediateFuture(null); + } + 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); + if (!cassandraTsPartitionsCache.has(partitionSearchKey)) { + ListenableFuture result = doSavePartition(tenantId, entityId, key, ttl, partition); + Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor()); + return result; + } else { + return Futures.immediateFuture(0); + } + } + } + @Override public ListenableFuture remove(TenantId tenantId, EntityId entityId, DeleteTsKvQuery query) { long minPartition = toPartitionTs(query.getStartTs()); @@ -449,47 +470,17 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD return getFuture(executeAsyncWrite(tenantId, stmt), rs -> null); } - @Override - public ListenableFuture savePartition(TenantId tenantId, EntityId entityId, long tsKvEntryTs, String key, long ttl) { - if (isFixedPartitioning()) { - return Futures.immediateFuture(null); - } - 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); - if (!cassandraTsPartitionsCache.has(partitionSearchKey)) { - ListenableFuture result = doSavePartition(tenantId, entityId, key, ttl, partition); - Futures.addCallback(result, new CacheCallback<>(partitionSearchKey), MoreExecutors.directExecutor()); - return result; - } else { - return Futures.immediateFuture(0); - } - } - } - 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); - PreparedStatement preparedStatement = ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt(); - BoundStatement stmt = preparedStatement.bind(); - stmt.setString(0, entityId.getEntityType().name()); - stmt.setUuid(1, entityId.getId()); - stmt.setLong(2, partition); - stmt.setString(3, key); + BoundStatementBuilder stmtBuilder = new BoundStatementBuilder((ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind()); + stmtBuilder.setString(0, entityId.getEntityType().name()) + .setUuid(1, entityId.getId()) + .setLong(2, partition) + .setString(3, key); if (ttl > 0) { - stmt.setInt(4, (int) ttl); + stmtBuilder.setInt(4, (int) ttl); } -// BoundStatementBuilder stmtBuilder = new BoundStatementBuilder(bind); -// stmtBuilder.setString(0, entityId.getEntityType().name()) -// .setUuid(1, entityId.getId()) -// .setLong(2, partition) -// .setString(3, key); -// if (ttl > 0) { -// stmtBuilder.setInt(4, (int) ttl); -// } -// BoundStatement stmt = stmtBuilder.build(); + BoundStatement stmt = stmtBuilder.build(); return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0); } 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 index 4a607d0f03..9646917ff7 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -21,6 +21,7 @@ import com.datastax.oss.driver.api.core.cql.PreparedStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.google.common.util.concurrent.Futures; import org.junit.Before; +import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; @@ -86,6 +87,7 @@ public class CassandraPartitionsCacheTest { doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(TbResultSetFuture.class), any()); } + @Ignore @Test public void testPartitionSave() throws Exception { cassandraBaseTimeseriesDao.init(); @@ -94,13 +96,14 @@ public class CassandraPartitionsCacheTest { 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)); +// 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); +// } + cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test", 0); + verify(cassandraBaseTimeseriesDao, times(1)).executeAsyncWrite(any(TenantId.class), any(Statement.class)); } } \ No newline at end of file From 637ad6cac5144ce205d0c9003b1a74813e045a96 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Tue, 22 Dec 2020 17:34:23 +0200 Subject: [PATCH 7/8] remove boundStatementBuilder from doSavePartition method --- .../CassandraBaseTimeseriesDao.java | 8 +-- .../nosql/CassandraPartitionsCacheTest.java | 53 ++++++++++--------- 2 files changed, 31 insertions(+), 30 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 acd140c0b6..db5f8f8684 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 @@ -472,15 +472,15 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD 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); - BoundStatementBuilder stmtBuilder = new BoundStatementBuilder((ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt()).bind()); - stmtBuilder.setString(0, entityId.getEntityType().name()) + PreparedStatement preparedStatement = ttl == 0 ? getPartitionInsertStmt() : getPartitionInsertTtlStmt(); + BoundStatement stmt = preparedStatement.bind(); + stmt = stmt.setString(0, entityId.getEntityType().name()) .setUuid(1, entityId.getId()) .setLong(2, partition) .setString(3, key); if (ttl > 0) { - stmtBuilder.setInt(4, (int) ttl); + stmt = stmt.setInt(4, (int) ttl); } - BoundStatement stmt = stmtBuilder.build(); return getFuture(executeAsyncWrite(tenantId, stmt), rs -> 0); } 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 index 9646917ff7..9a00692cb2 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -21,10 +21,10 @@ import com.datastax.oss.driver.api.core.cql.PreparedStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.google.common.util.concurrent.Futures; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.Mock; +import org.mockito.Spy; import org.mockito.runners.MockitoJUnitRunner; import org.springframework.core.env.Environment; import org.springframework.test.util.ReflectionTestUtils; @@ -36,9 +36,9 @@ 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.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; @@ -46,36 +46,29 @@ import static org.mockito.Mockito.when; @RunWith(MockitoJUnitRunner.class) public class CassandraPartitionsCacheTest { + @Spy private CassandraBaseTimeseriesDao cassandraBaseTimeseriesDao; @Mock - private Environment environment; + private PreparedStatement preparedStatement; @Mock - private CassandraBufferedRateExecutor rateLimiter; + private BoundStatement boundStatement; @Mock - private CassandraCluster cluster; + private Environment environment; @Mock - private GuavaSession session; + private CassandraBufferedRateExecutor rateLimiter; @Mock - private PreparedStatement preparedStatement; + private CassandraCluster cluster; @Mock - private BoundStatement boundStatement; + private GuavaSession session; @Before public void setUp() throws Exception { - when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); - when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); - when(cluster.getSession()).thenReturn(session); - when(session.prepare(anyString())).thenReturn(preparedStatement); - when(preparedStatement.bind()).thenReturn(boundStatement); - - cassandraBaseTimeseriesDao = spy(new CassandraBaseTimeseriesDao()); - ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitioning", "MONTHS"); ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "partitionsCacheSize", 100000); ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "systemTtl", 0); @@ -84,10 +77,20 @@ public class CassandraPartitionsCacheTest { ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "rateLimiter", rateLimiter); ReflectionTestUtils.setField(cassandraBaseTimeseriesDao, "cluster", cluster); + when(cluster.getDefaultReadConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getDefaultWriteConsistencyLevel()).thenReturn(ConsistencyLevel.ONE); + when(cluster.getSession()).thenReturn(session); + 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(), any(Long.class))).thenReturn(boundStatement); + doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(TbResultSetFuture.class), any()); } - @Ignore @Test public void testPartitionSave() throws Exception { cassandraBaseTimeseriesDao.init(); @@ -96,14 +99,12 @@ public class CassandraPartitionsCacheTest { 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); -// } - cassandraBaseTimeseriesDao.savePartition(tenantId, tenantId, tsKvEntryTs, "test", 0); - verify(cassandraBaseTimeseriesDao, times(1)).executeAsyncWrite(any(TenantId.class), any(Statement.class)); + 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 From fb8ddbda584393dd1abe8a59ec68120ec3949bc6 Mon Sep 17 00:00:00 2001 From: ShvaykaD Date: Wed, 23 Dec 2020 11:46:05 +0200 Subject: [PATCH 8/8] changed doReturn for cassandraBaseTimeseriesDao.getFuture in tests --- .../server/dao/nosql/CassandraPartitionsCacheTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 index 9a00692cb2..f19165a557 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/nosql/CassandraPartitionsCacheTest.java @@ -88,7 +88,7 @@ public class CassandraPartitionsCacheTest { when(boundStatement.setUuid(anyInt(), any(UUID.class))).thenReturn(boundStatement); when(boundStatement.setLong(anyInt(), any(Long.class))).thenReturn(boundStatement); - doReturn(Futures.immediateFuture(null)).when(cassandraBaseTimeseriesDao).getFuture(any(TbResultSetFuture.class), any()); + doReturn(Futures.immediateFuture(0)).when(cassandraBaseTimeseriesDao).getFuture(any(TbResultSetFuture.class), any()); } @Test