From 349bf398e2d7d1d58502e6bb03330066f825daf1 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Thu, 3 May 2018 15:57:03 +0300 Subject: [PATCH 1/7] first stage delete timeseries records, implemented orderBy query --- .../server/common/data/kv/BaseTsKvQuery.java | 6 +- .../server/common/data/kv/TsKvQuery.java | 1 + .../dao/sql/timeseries/JpaTimeseriesDao.java | 12 ++ .../dao/timeseries/BaseTimeseriesService.java | 17 ++ .../CassandraBaseTimeseriesDao.java | 159 +++++++++++++++--- .../server/dao/timeseries/TimeseriesDao.java | 4 + .../dao/timeseries/TimeseriesService.java | 2 + .../dao/timeseries/TsKvQueryCursor.java | 19 ++- .../timeseries/BaseTimeseriesServiceTest.java | 31 +++- .../handlers/TelemetryRestMsgHandler.java | 2 +- .../TelemetryWebsocketMsgHandler.java | 5 +- 11 files changed, 220 insertions(+), 38 deletions(-) diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java index 0afe00b4c3..51d4ad2004 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java @@ -26,18 +26,20 @@ public class BaseTsKvQuery implements TsKvQuery { private final long interval; private final int limit; private final Aggregation aggregation; + private final String orderBy; - public BaseTsKvQuery(String key, long startTs, long endTs, long interval, int limit, Aggregation aggregation) { + public BaseTsKvQuery(String key, long startTs, long endTs, long interval, int limit, Aggregation aggregation, String orderBy) { this.key = key; this.startTs = startTs; this.endTs = endTs; this.interval = interval; this.limit = limit; this.aggregation = aggregation; + this.orderBy = orderBy; } public BaseTsKvQuery(String key, long startTs, long endTs) { - this(key, startTs, endTs, endTs-startTs, 1, Aggregation.AVG); + this(key, startTs, endTs, endTs - startTs, 1, Aggregation.AVG, "DESC"); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java index ca9f90c4e8..9b907c3440 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java @@ -29,4 +29,5 @@ public interface TsKvQuery { Aggregation getAggregation(); + String getOrderBy(); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java index 6350352961..5503f49262 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java @@ -299,6 +299,18 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp }); } + @Override + public ListenableFuture remove(EntityId entityId, TsKvQuery query) { + //TODO: implement + return null; + } + + @Override + public ListenableFuture removeLatest(EntityId entityId, TsKvQuery query) { + //TODO: implement + return null; + } + @PreDestroy void onDestroy() { if (insertService != null) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index c981378939..a075885308 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -40,6 +40,7 @@ import static org.apache.commons.lang3.StringUtils.isBlank; public class BaseTimeseriesService implements TimeseriesService { public static final int INSERTS_PER_ENTRY = 3; + public static final int DELETES_PER_ENTRY = 2; @Autowired private TimeseriesDao timeseriesDao; @@ -95,6 +96,22 @@ public class BaseTimeseriesService implements TimeseriesService { futures.add(timeseriesDao.save(entityId, tsKvEntry, ttl)); } + @Override + public ListenableFuture> remove(EntityId entityId, List tsKvQueries) { + validate(entityId); + tsKvQueries.forEach(BaseTimeseriesService::validate); + List> futures = Lists.newArrayListWithExpectedSize(tsKvQueries.size() * DELETES_PER_ENTRY); + for (TsKvQuery tsKvQuery : tsKvQueries) { + deleteAndRegisterFutures(futures, entityId, tsKvQuery); + } + return Futures.allAsList(futures); + } + + private void deleteAndRegisterFutures(List> futures, EntityId entityId, TsKvQuery query) { + futures.add(timeseriesDao.remove(entityId, query)); + futures.add(timeseriesDao.removeLatest(entityId, query)); + } + private static void validate(EntityId entityId) { Validator.validateEntityId(entityId, "Incorrect entityId " + entityId); } 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 cda4b1669b..7895b59f20 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 @@ -62,6 +62,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem public static final String GENERATED_QUERY_FOR_ENTITY_TYPE_AND_ENTITY_ID = "Generated query [{}] for entityType {} and entityId {}"; public static final String SELECT_PREFIX = "SELECT "; public static final String EQUALS_PARAM = " = ? "; + public static final String ASC_ORDER = "ASC"; + public static final String DESC_ORDER = "DESC"; @Autowired private Environment environment; @@ -76,7 +78,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem private PreparedStatement latestInsertStmt; private PreparedStatement[] saveStmts; private PreparedStatement[] saveTtlStmts; - private PreparedStatement[] fetchStmts; + private PreparedStatement[] fetchStmtsAsc; + private PreparedStatement[] fetchStmtsDesc; private PreparedStatement findLatestStmt; private PreparedStatement findAllLatestStmt; @@ -88,7 +91,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem public void init() { super.startExecutor(); if (!isInstall()) { - getFetchStmt(Aggregation.NONE); + getFetchStmt(Aggregation.NONE, DESC_ORDER); Optional partition = TsPartitionDate.parse(partitioning); if (partition.isPresent()) { tsFormat = partition.get(); @@ -132,7 +135,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem while (stepTs < query.getEndTs()) { long startTs = stepTs; long endTs = stepTs + step; - TsKvQuery subQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, step, 1, query.getAggregation()); + TsKvQuery subQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, step, 1, query.getAggregation(), query.getOrderBy()); futures.add(findAndAggregateAsync(entityId, subQuery, toPartitionTs(startTs), toPartitionTs(endTs))); stepTs = endTs; } @@ -181,7 +184,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem if (cursor.isFull() || !cursor.hasNextPartition()) { resultFuture.set(cursor.getData()); } else { - PreparedStatement proto = getFetchStmt(Aggregation.NONE); + PreparedStatement proto = getFetchStmt(Aggregation.NONE, cursor.getOrderBy()); BoundStatement stmt = proto.bind(); stmt.setString(0, cursor.getEntityType()); stmt.setUUID(1, cursor.getEntityId()); @@ -231,7 +234,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem private AsyncFunction, List> getFetchChunksAsyncFunction(EntityId entityId, String key, Aggregation aggregation, long startTs, long endTs) { return partitions -> { try { - PreparedStatement proto = getFetchStmt(aggregation); + PreparedStatement proto = getFetchStmt(aggregation, DESC_ORDER); List futures = new ArrayList<>(partitions.size()); for (Long partition : partitions) { log.trace("Fetching data for partition [{}] for entityType {} and entityId {}", partition, entityId.getEntityType(), entityId.getId()); @@ -318,6 +321,99 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return getFuture(executeAsyncWrite(stmt), rs -> null); } + @Override + public ListenableFuture remove(EntityId entityId, TsKvQuery query) { + long minPartition = toPartitionTs(query.getStartTs()); + long maxPartition = toPartitionTs(query.getEndTs()); + + ResultSetFuture partitionsFuture = fetchPartitions(entityId, query.getKey(), minPartition, maxPartition); + + final SimpleListenableFuture resultFuture = new SimpleListenableFuture<>(); + final ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor); + + Futures.addCallback(partitionsListFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List partitions) { + TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions); + deleteAsync(cursor, resultFuture); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); + } + }, readResultsProcessingExecutor); + return resultFuture; + } + + private void deleteAsync(final TsKvQueryCursor cursor, final SimpleListenableFuture resultFuture) { + if (!cursor.hasNextPartition()) { + resultFuture.set(null); + } else { + PreparedStatement proto = getDeleteStmt(); + BoundStatement stmt = proto.bind(); + stmt.setString(0, cursor.getEntityType()); + stmt.setUUID(1, cursor.getEntityId()); + stmt.setString(2, cursor.getKey()); + stmt.setLong(3, cursor.getNextPartition()); + stmt.setLong(4, cursor.getStartTs()); + stmt.setLong(5, cursor.getEndTs()); + + Futures.addCallback(executeAsyncWrite(stmt), new FutureCallback() { + @Override + public void onSuccess(@Nullable ResultSet result) { + deleteAsync(cursor, resultFuture); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to delete data for query {}-{}", stmt, t); + } + }, readResultsProcessingExecutor); + } + } + + private PreparedStatement getDeleteStmt() { + return prepare("DELETE FROM " + ModelConstants.TS_KV_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " > ? " + + "AND " + ModelConstants.TS_COLUMN + " <= ?"); + } + + @Override + public ListenableFuture removeLatest(EntityId entityId, TsKvQuery query) { + ListenableFuture future = findLatest(entityId, query.getKey()); + return Futures.transform(future, new Function() { + @Nullable + @Override + public Void apply(@Nullable TsKvEntry latestEntry) { + if (latestEntry != null) { + long ts = latestEntry.getTs(); + if (ts >= query.getStartTs() && ts <= query.getEndTs()) { + deleteLatest(entityId, latestEntry.getKey()); + + //TODO: save new latest entry(< query.getStartTs() - if present) to TS_KV_LATEST_CF + } else { + log.trace("Won't be deleted latest value for [{}], key - {}", entityId, query.getKey()); + } + } + return null; + } + }); + } + + private ListenableFuture deleteLatest(EntityId entityId, String key) { + Statement delete = QueryBuilder.delete().from(ModelConstants.TS_KV_LATEST_CF) + .where(eq(ModelConstants.ENTITY_TYPE_COLUMN, entityId.getEntityType())) + .and(eq(ModelConstants.ENTITY_ID_COLUMN, entityId.getId())) + .and(eq(ModelConstants.KEY_COLUMN, key)); + log.debug("Remove request: {}", delete.toString()); + return getFuture(executeAsyncWrite(delete), rs -> null); + } + private List convertResultToTsKvEntryList(List rows) { List entries = new ArrayList<>(rows.size()); if (!rows.isEmpty()) { @@ -413,28 +509,43 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return saveTtlStmts[dataType.ordinal()]; } - private PreparedStatement getFetchStmt(Aggregation aggType) { - if (fetchStmts == null) { - fetchStmts = new PreparedStatement[Aggregation.values().length]; - for (Aggregation type : Aggregation.values()) { - if (type == Aggregation.SUM && fetchStmts[Aggregation.AVG.ordinal()] != null) { - fetchStmts[type.ordinal()] = fetchStmts[Aggregation.AVG.ordinal()]; - } else if (type == Aggregation.AVG && fetchStmts[Aggregation.SUM.ordinal()] != null) { - fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()]; - } else { - fetchStmts[type.ordinal()] = prepare(SELECT_PREFIX + - String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF - + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.TS_COLUMN + " > ? " - + "AND " + ModelConstants.TS_COLUMN + " <= ?" - + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " DESC LIMIT ?" : "")); + private PreparedStatement getFetchStmt(Aggregation aggType, String orderBy) { + switch (orderBy) { + case ASC_ORDER: + if (fetchStmtsAsc == null) { + fetchStmtsAsc = initFetchStmt(orderBy); + } + return fetchStmtsAsc[aggType.ordinal()]; + case DESC_ORDER: + if (fetchStmtsDesc == null) { + fetchStmtsDesc = initFetchStmt(orderBy); } + return fetchStmtsDesc[aggType.ordinal()]; + default: + throw new RuntimeException("Not supported" + orderBy + "order!"); + } + } + + private PreparedStatement[] initFetchStmt(String orderBy) { + PreparedStatement[] fetchStmts = new PreparedStatement[Aggregation.values().length]; + for (Aggregation type : Aggregation.values()) { + if (type == Aggregation.SUM && fetchStmts[Aggregation.AVG.ordinal()] != null) { + fetchStmts[type.ordinal()] = fetchStmts[Aggregation.AVG.ordinal()]; + } else if (type == Aggregation.AVG && fetchStmts[Aggregation.SUM.ordinal()] != null) { + fetchStmts[type.ordinal()] = fetchStmts[Aggregation.SUM.ordinal()]; + } else { + fetchStmts[type.ordinal()] = prepare(SELECT_PREFIX + + String.join(", ", ModelConstants.getFetchColumnNames(type)) + " FROM " + ModelConstants.TS_KV_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " > ? " + + "AND " + ModelConstants.TS_COLUMN + " <= ?" + + (type == Aggregation.NONE ? " ORDER BY " + ModelConstants.TS_COLUMN + " " + orderBy + " LIMIT ?" : "")); } } - return fetchStmts[aggType.ordinal()]; + return fetchStmts; } private PreparedStatement getLatestStmt() { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java index 1e3f4cecb7..22bb166585 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java @@ -38,4 +38,8 @@ public interface TimeseriesDao { ListenableFuture savePartition(EntityId entityId, long tsKvEntryTs, String key, long ttl); ListenableFuture saveLatest(EntityId entityId, TsKvEntry tsKvEntry); + + ListenableFuture remove(EntityId entityId, TsKvQuery query); + + ListenableFuture removeLatest(EntityId entityId, TsKvQuery query); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index 2cd2d8dab9..a14919185d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java @@ -37,4 +37,6 @@ public interface TimeseriesService { ListenableFuture> save(EntityId entityId, TsKvEntry tsKvEntry); ListenableFuture> save(EntityId entityId, List tsKvEntry, long ttl); + + ListenableFuture> remove(EntityId entityId, List queries); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TsKvQueryCursor.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TsKvQueryCursor.java index d6b6bbd5c0..c4925ee9be 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TsKvQueryCursor.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TsKvQueryCursor.java @@ -23,6 +23,8 @@ import java.util.ArrayList; import java.util.List; import java.util.UUID; +import static org.thingsboard.server.dao.timeseries.CassandraBaseTimeseriesDao.DESC_ORDER; + /** * Created by ashvayka on 21.02.17. */ @@ -40,6 +42,8 @@ public class TsKvQueryCursor { private final List partitions; @Getter private final List data; + @Getter + private String orderBy; private int partitionIndex; private int currentLimit; @@ -51,13 +55,14 @@ public class TsKvQueryCursor { this.startTs = baseQuery.getStartTs(); this.endTs = baseQuery.getEndTs(); this.partitions = partitions; - this.partitionIndex = partitions.size() - 1; + this.orderBy = baseQuery.getOrderBy(); + this.partitionIndex = isDesc() ? partitions.size() - 1 : 0; this.data = new ArrayList<>(); this.currentLimit = baseQuery.getLimit(); } public boolean hasNextPartition() { - return partitionIndex >= 0; + return isDesc() ? partitionIndex >= 0 : partitionIndex <= partitions.size() - 1; } public boolean isFull() { @@ -66,7 +71,11 @@ public class TsKvQueryCursor { public long getNextPartition() { long partition = partitions.get(partitionIndex); - partitionIndex--; + if (isDesc()) { + partitionIndex--; + } else { + partitionIndex++; + } return partition; } @@ -78,4 +87,8 @@ public class TsKvQueryCursor { currentLimit -= newData.size(); data.addAll(newData); } + + private boolean isDesc() { + return orderBy.equals(DESC_ORDER); + } } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index 0cb3f7fbdd..2130ab97d0 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -45,6 +45,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { private static final String BOOLEAN_KEY = "booleanKey"; private static final long TS = 42L; + private static final String DESC_ORDER = "DESC"; KvEntry stringKvEntry = new StringDataEntry(STRING_KEY, "value"); KvEntry longKvEntry = new LongDataEntry(LONG_KEY, Long.MAX_VALUE); @@ -92,6 +93,24 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { Assert.assertEquals(toTsEntry(TS, stringKvEntry), entries.get(0)); } + @Test + public void testDeleteDeviceTsData() throws Exception { + DeviceId deviceId = new DeviceId(UUIDs.timeBased()); + + saveEntries(deviceId, TS - 3); + saveEntries(deviceId, TS - 2); + saveEntries(deviceId, TS - 1); + saveEntries(deviceId, TS); + + tsService.remove(deviceId, Collections.singletonList( + new BaseTsKvQuery(STRING_KEY, TS - 4, TS - 2))).get(); + + List list = tsService.findAll(deviceId, Collections.singletonList( + new BaseTsKvQuery(STRING_KEY, 0, 60000, 60000, 5, Aggregation.NONE, DESC_ORDER))).get(); + + Assert.assertEquals(2, list.size()); + } + @Test public void testFindDeviceTsData() throws Exception { DeviceId deviceId = new DeviceId(UUIDs.timeBased()); @@ -107,7 +126,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { entries.add(save(deviceId, 55000, 600)); List list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.NONE))).get(); + 60000, 20000, 3, Aggregation.NONE, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(55000, list.get(0).getTs()); assertEquals(java.util.Optional.of(600L), list.get(0).getLongValue()); @@ -119,7 +138,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(400L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.AVG))).get(); + 60000, 20000, 3, Aggregation.AVG, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); assertEquals(java.util.Optional.of(150L), list.get(0).getLongValue()); @@ -131,7 +150,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(550L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.SUM))).get(); + 60000, 20000, 3, Aggregation.SUM, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -144,7 +163,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(1100L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.MIN))).get(); + 60000, 20000, 3, Aggregation.MIN, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -157,7 +176,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(500L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.MAX))).get(); + 60000, 20000, 3, Aggregation.MAX, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -170,7 +189,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(600L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.COUNT))).get(); + 60000, 20000, 3, Aggregation.COUNT, DESC_ORDER))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java index 3ea754ae9c..0c7e3874ee 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java @@ -140,7 +140,7 @@ public class TelemetryRestMsgHandler extends DefaultRestMsgHandler { Aggregation agg = (interval.isPresent() && interval.get() == 0) ? Aggregation.valueOf(Aggregation.NONE.name()) : Aggregation.valueOf(request.getParameter("agg", Aggregation.NONE.name())); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs.get(), endTs.get(), interval.get(), limit.orElse(TelemetryWebsocketMsgHandler.DEFAULT_LIMIT), agg)) + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs.get(), endTs.get(), interval.get(), limit.orElse(TelemetryWebsocketMsgHandler.DEFAULT_LIMIT), agg, "DESC")) .collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, getTsKvListCallback(msg)); } else { diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java index 1374ef68ac..bf75c5de29 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java @@ -54,6 +54,7 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { public static final String FAILED_TO_FETCH_DATA = "Failed to fetch data!"; public static final String FAILED_TO_FETCH_ATTRIBUTES = "Failed to fetch attributes!"; public static final String SESSION_META_DATA_NOT_FOUND = "Session meta-data not found!"; + public static final String ORDER_BY = "DESC"; private final SubscriptionManager subscriptionManager; @@ -216,7 +217,7 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { log.debug("[{}] fetching timeseries data for last {} ms for keys: ({}) for device : {}", sessionId, cmd.getTimeWindow(), cmd.getKeys(), entityId); startTs = cmd.getStartTs(); long endTs = cmd.getStartTs() + cmd.getTimeWindow(); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs, endTs, cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()))).collect(Collectors.toList()); + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs, endTs, cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY)).collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, getSubscriptionCallback(sessionRef, cmd, sessionId, entityId, startTs, keys)); } else { List keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); @@ -300,7 +301,7 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { } EntityId entityId = EntityIdFactory.getByTypeAndId(cmd.getEntityType(), cmd.getEntityId()); List keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()))) + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY)) .collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, new PluginCallback>() { @Override From 36ad0f844c104287b2cfc38f37e594cf723f5a47 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Thu, 3 May 2018 16:29:32 +0300 Subject: [PATCH 2/7] temp fix --- .../dao/service/timeseries/BaseTimeseriesServiceTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index 2130ab97d0..181a19743b 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -93,7 +93,8 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { Assert.assertEquals(toTsEntry(TS, stringKvEntry), entries.get(0)); } - @Test + //TODO: sql delete implement + /*@Test public void testDeleteDeviceTsData() throws Exception { DeviceId deviceId = new DeviceId(UUIDs.timeBased()); @@ -109,7 +110,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { new BaseTsKvQuery(STRING_KEY, 0, 60000, 60000, 5, Aggregation.NONE, DESC_ORDER))).get(); Assert.assertEquals(2, list.size()); - } + }*/ @Test public void testFindDeviceTsData() throws Exception { From faf14d43a81cee77cd17577c41a46004f94511ac Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Wed, 30 May 2018 18:27:06 +0300 Subject: [PATCH 3/7] improved removing timeseries --- .../server/controller/DeviceController.java | 56 ++++++- .../src/main/resources/thingsboard.yml | 6 +- .../dao/sql/timeseries/JpaTimeseriesDao.java | 6 +- .../dao/timeseries/BaseTimeseriesService.java | 3 +- .../CassandraBaseTimeseriesDao.java | 139 +++++++++++++++--- .../server/dao/timeseries/TimeseriesDao.java | 2 + 6 files changed, 182 insertions(+), 30 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java index bceea54135..ed9bf77c46 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java @@ -16,6 +16,7 @@ package org.thingsboard.server.controller; import com.google.common.util.concurrent.ListenableFuture; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.*; @@ -23,23 +24,27 @@ import org.thingsboard.server.common.data.Customer; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; -import org.thingsboard.server.common.data.audit.ActionStatus; import org.thingsboard.server.common.data.audit.ActionType; import org.thingsboard.server.common.data.device.DeviceSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.BaseTsKvQuery; +import org.thingsboard.server.common.data.kv.BasicTsKvEntry; +import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.model.ModelConstants; +import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.exception.ThingsboardErrorCode; import org.thingsboard.server.exception.ThingsboardException; import org.thingsboard.server.service.security.model.SecurityUser; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.stream.Collectors; @@ -47,6 +52,9 @@ import java.util.stream.Collectors; @RequestMapping("/api") public class DeviceController extends BaseController { + @Autowired + protected TimeseriesService timeseriesService; + public static final String DEVICE_ID = "deviceId"; @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @@ -70,7 +78,7 @@ public class DeviceController extends BaseController { device.setTenantId(getCurrentUser().getTenantId()); if (getCurrentUser().getAuthority() == Authority.CUSTOMER_USER) { if (device.getId() == null || device.getId().isNullUid() || - device.getCustomerId() == null || device.getCustomerId().isNullUid()) { + device.getCustomerId() == null || device.getCustomerId().isNullUid()) { throw new ThingsboardException("You don't have permission to perform this operation!", ThingsboardErrorCode.PERMISSION_DENIED); } else { @@ -368,4 +376,48 @@ public class DeviceController extends BaseController { throw handleException(e); } } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/device/testSave", method = RequestMethod.GET) + @ResponseBody + public void testSave() throws ThingsboardException { + try { + SecurityUser user = getCurrentUser(); + TenantId tenantId = user.getTenantId(); + + Device device = deviceService.findDeviceByTenantIdAndName(tenantId, "Test"); + + timeseriesService.save(device.getId(), new BasicTsKvEntry(1516892633000L, + new LongDataEntry("test", 1L))).get(); + timeseriesService.save(device.getId(), new BasicTsKvEntry(1519571033000L, + new LongDataEntry("test", 2L))).get(); + timeseriesService.save(device.getId(), new BasicTsKvEntry(1521990233000L, + new LongDataEntry("test", 3L))).get(); + timeseriesService.save(device.getId(), new BasicTsKvEntry(1524668633000L, + new LongDataEntry("test", 4L))).get(); + timeseriesService.save(device.getId(), new BasicTsKvEntry(1527260633000L, + new LongDataEntry("test", 5L))).get(); + + } catch (Exception e) { + throw handleException(e); + } + } + + @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") + @RequestMapping(value = "/device/testDelete", method = RequestMethod.GET) + @ResponseBody + public void testDelete() throws ThingsboardException { + try { + SecurityUser user = getCurrentUser(); + TenantId tenantId = user.getTenantId(); + + Device device = deviceService.findDeviceByTenantIdAndName(tenantId, "Test"); + + timeseriesService.remove(device.getId(), Collections.singletonList(new BaseTsKvQuery("test", + 1519139033000L, 1524668633000L))).get(); + + } catch (Exception e) { + throw handleException(e); + } + } } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index c47fc28fbd..aef92804f1 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -66,8 +66,8 @@ plugins: # JWT Token parameters security.jwt: - tokenExpirationTime: "${JWT_TOKEN_EXPIRATION_TIME:900}" # Number of seconds (15 mins) - refreshTokenExpTime: "${JWT_REFRESH_TOKEN_EXPIRATION_TIME:3600}" # Seconds (1 hour) + tokenExpirationTime: "${JWT_TOKEN_EXPIRATION_TIME:9000000}" # Number of seconds (15 mins) + refreshTokenExpTime: "${JWT_REFRESH_TOKEN_EXPIRATION_TIME:36000000}" # Seconds (1 hour) tokenIssuer: "${JWT_TOKEN_ISSUER:thingsboard.io}" tokenSigningKey: "${JWT_TOKEN_SIGNING_KEY:thingsboardDefaultSigningKey}" @@ -133,7 +133,7 @@ quota: intervalMin: 2 database: - type: "${DATABASE_TYPE:sql}" # cassandra OR sql + type: "${DATABASE_TYPE:cassandra}" # cassandra OR sql # Cassandra driver configuration parameters cassandra: diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java index 5503f49262..bef82c3677 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java @@ -42,7 +42,6 @@ import java.util.ArrayList; import java.util.List; import java.util.Optional; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.stream.Collectors; @@ -311,6 +310,11 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp return null; } + @Override + public ListenableFuture removePartition(EntityId entityId, TsKvQuery query) { + return insertService.submit(() -> null); + } + @PreDestroy void onDestroy() { if (insertService != null) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index a075885308..f035fd8128 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -40,7 +40,7 @@ import static org.apache.commons.lang3.StringUtils.isBlank; public class BaseTimeseriesService implements TimeseriesService { public static final int INSERTS_PER_ENTRY = 3; - public static final int DELETES_PER_ENTRY = 2; + public static final int DELETES_PER_ENTRY = INSERTS_PER_ENTRY; @Autowired private TimeseriesDao timeseriesDao; @@ -110,6 +110,7 @@ public class BaseTimeseriesService implements TimeseriesService { private void deleteAndRegisterFutures(List> futures, EntityId entityId, TsKvQuery query) { futures.add(timeseriesDao.remove(entityId, query)); futures.add(timeseriesDao.removeLatest(entityId, query)); + futures.add(timeseriesDao.removePartition(entityId, query)); } private static void validate(EntityId entityId) { 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 7895b59f20..423e90efbf 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 @@ -41,10 +41,7 @@ import javax.annotation.PreDestroy; import java.time.Instant; import java.time.LocalDateTime; import java.time.ZoneOffset; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.Optional; +import java.util.*; import java.util.stream.Collectors; import static com.datastax.driver.core.querybuilder.QueryBuilder.eq; @@ -82,6 +79,8 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem private PreparedStatement[] fetchStmtsDesc; private PreparedStatement findLatestStmt; private PreparedStatement findAllLatestStmt; + private PreparedStatement deleteStmt; + private PreparedStatement deletePartitionStmt; private boolean isInstall() { return environment.acceptsProfiles("install"); @@ -374,35 +373,68 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } private PreparedStatement getDeleteStmt() { - return prepare("DELETE FROM " + ModelConstants.TS_KV_CF + - " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM - + "AND " + ModelConstants.TS_COLUMN + " > ? " - + "AND " + ModelConstants.TS_COLUMN + " <= ?"); + if (deleteStmt == null) { + deleteStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.TS_COLUMN + " > ? " + + "AND " + ModelConstants.TS_COLUMN + " <= ?"); + } + return deleteStmt; } @Override public ListenableFuture removeLatest(EntityId entityId, TsKvQuery query) { - ListenableFuture future = findLatest(entityId, query.getKey()); - return Futures.transform(future, new Function() { - @Nullable - @Override - public Void apply(@Nullable TsKvEntry latestEntry) { - if (latestEntry != null) { + ListenableFuture latestEntryFuture = findLatest(entityId, query.getKey()); + + ListenableFuture booleanFuture = Futures.transform(latestEntryFuture, + (AsyncFunction) latestEntry -> { long ts = latestEntry.getTs(); if (ts >= query.getStartTs() && ts <= query.getEndTs()) { - deleteLatest(entityId, latestEntry.getKey()); - - //TODO: save new latest entry(< query.getStartTs() - if present) to TS_KV_LATEST_CF + return Futures.immediateFuture(true); } else { log.trace("Won't be deleted latest value for [{}], key - {}", entityId, query.getKey()); } - } - return null; + return Futures.immediateFuture(false); + }, readResultsProcessingExecutor); + + + ListenableFuture savedLatestFuture = Futures.transform(booleanFuture, + (AsyncFunction) isRemove -> { + if (isRemove) { + return getNewLatestEntryFuture(entityId, query); + } + return Futures.immediateFuture(null); + }, readResultsProcessingExecutor); + + ListenableFuture removedLatestFuture = Futures.transform(booleanFuture, + (AsyncFunction) isRemove -> { + if (isRemove) { + return deleteLatest(entityId, query.getKey()); + } + return Futures.immediateFuture(null); + }, readResultsProcessingExecutor); + return Futures.transform(Futures.allAsList(Arrays.asList(savedLatestFuture, removedLatestFuture)), + (AsyncFunction, Void>) list -> Futures.immediateFuture(null), readResultsProcessingExecutor); + } + + private ListenableFuture getNewLatestEntryFuture(EntityId entityId, TsKvQuery query) { + long startTs = 0; + long endTs = query.getStartTs() - 1; + TsKvQuery findNewLatestQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, + Aggregation.NONE, DESC_ORDER); + ListenableFuture> future = findAllAsync(entityId, findNewLatestQuery); + + return Futures.transform(future, (AsyncFunction, Void>) entryList -> { + if (entryList.size() == 1) { + return saveLatest(entityId, entryList.get(0)); + } else { + log.trace("Could not find new latest value for [{}], key - {}", entityId, query.getKey()); } - }); + return Futures.immediateFuture(null); + }, readResultsProcessingExecutor); } private ListenableFuture deleteLatest(EntityId entityId, String key) { @@ -414,6 +446,67 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return getFuture(executeAsyncWrite(delete), rs -> null); } + @Override + public ListenableFuture removePartition(EntityId entityId, TsKvQuery query) { + long minPartition = toPartitionTs(query.getStartTs()); + long maxPartition = toPartitionTs(query.getEndTs()); + + ResultSetFuture partitionsFuture = fetchPartitions(entityId, query.getKey(), minPartition, maxPartition); + + final SimpleListenableFuture resultFuture = new SimpleListenableFuture<>(); + final ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor); + + Futures.addCallback(partitionsListFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List partitions) { + TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions); + deletePartitionAsync(cursor, resultFuture); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); + } + }, readResultsProcessingExecutor); + return resultFuture; + } + + private void deletePartitionAsync(final TsKvQueryCursor cursor, final SimpleListenableFuture resultFuture) { + if (!cursor.hasNextPartition()) { + resultFuture.set(null); + } else { + PreparedStatement proto = getDeletePartitionStmt(); + BoundStatement stmt = proto.bind(); + stmt.setString(0, cursor.getEntityType()); + stmt.setUUID(1, cursor.getEntityId()); + stmt.setLong(2, cursor.getNextPartition()); + stmt.setString(3, cursor.getKey()); + + Futures.addCallback(executeAsyncWrite(stmt), new FutureCallback() { + @Override + public void onSuccess(@Nullable ResultSet result) { + deletePartitionAsync(cursor, resultFuture); + } + + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to delete data for query {}-{}", stmt, t); + } + }, readResultsProcessingExecutor); + } + } + + private PreparedStatement getDeletePartitionStmt() { + if (deletePartitionStmt == null) { + deletePartitionStmt = prepare("DELETE FROM " + ModelConstants.TS_KV_PARTITIONS_CF + + " WHERE " + ModelConstants.ENTITY_TYPE_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.ENTITY_ID_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.PARTITION_COLUMN + EQUALS_PARAM + + "AND " + ModelConstants.KEY_COLUMN + EQUALS_PARAM); + } + return deletePartitionStmt; + } + private List convertResultToTsKvEntryList(List rows) { List entries = new ArrayList<>(rows.size()); if (!rows.isEmpty()) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java index 22bb166585..62dbd504c5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java @@ -42,4 +42,6 @@ public interface TimeseriesDao { ListenableFuture remove(EntityId entityId, TsKvQuery query); ListenableFuture removeLatest(EntityId entityId, TsKvQuery query); + + ListenableFuture removePartition(EntityId entityId, TsKvQuery query); } From e37e7242fdb01e4cd385c347f14e6706bd6a40b8 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Mon, 4 Jun 2018 19:38:08 +0300 Subject: [PATCH 4/7] find and save new latest if previous deleted --- .../server/controller/DeviceController.java | 5 +++- .../server/common/data/kv/BaseTsKvQuery.java | 7 +++-- .../server/common/data/kv/TsKvQuery.java | 2 ++ .../CassandraBaseTimeseriesDao.java | 29 ++++++++++--------- .../timeseries/BaseTimeseriesServiceTest.java | 12 ++++---- .../handlers/TelemetryRestMsgHandler.java | 10 +++---- .../TelemetryWebsocketMsgHandler.java | 6 ++-- 7 files changed, 42 insertions(+), 29 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java index ed9bf77c46..856ca8e399 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java @@ -29,6 +29,7 @@ import org.thingsboard.server.common.data.device.DeviceSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; +import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.BaseTsKvQuery; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; @@ -413,8 +414,10 @@ public class DeviceController extends BaseController { Device device = deviceService.findDeviceByTenantIdAndName(tenantId, "Test"); + long startTs = 1519561033000L; + long endTs = 1528260633000L; timeseriesService.remove(device.getId(), Collections.singletonList(new BaseTsKvQuery("test", - 1519139033000L, 1524668633000L))).get(); + startTs, endTs, endTs - startTs, 0, Aggregation.NONE, "DESC", true))).get(); } catch (Exception e) { throw handleException(e); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java index 51d4ad2004..55d279768e 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BaseTsKvQuery.java @@ -27,8 +27,10 @@ public class BaseTsKvQuery implements TsKvQuery { private final int limit; private final Aggregation aggregation; private final String orderBy; + private final Boolean rewriteLatestIfDeleted; - public BaseTsKvQuery(String key, long startTs, long endTs, long interval, int limit, Aggregation aggregation, String orderBy) { + public BaseTsKvQuery(String key, long startTs, long endTs, long interval, int limit, Aggregation aggregation, String orderBy, + boolean rewriteLatestIfDeleted) { this.key = key; this.startTs = startTs; this.endTs = endTs; @@ -36,10 +38,11 @@ public class BaseTsKvQuery implements TsKvQuery { this.limit = limit; this.aggregation = aggregation; this.orderBy = orderBy; + this.rewriteLatestIfDeleted = rewriteLatestIfDeleted; } public BaseTsKvQuery(String key, long startTs, long endTs) { - this(key, startTs, endTs, endTs - startTs, 1, Aggregation.AVG, "DESC"); + this(key, startTs, endTs, endTs - startTs, 1, Aggregation.AVG, "DESC", false); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java index 9b907c3440..825df6c176 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvQuery.java @@ -30,4 +30,6 @@ public interface TsKvQuery { Aggregation getAggregation(); String getOrderBy(); + + Boolean getRewriteLatestIfDeleted(); } 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 423e90efbf..3e4e2bdac5 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 @@ -134,7 +134,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem while (stepTs < query.getEndTs()) { long startTs = stepTs; long endTs = stepTs + step; - TsKvQuery subQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, step, 1, query.getAggregation(), query.getOrderBy()); + TsKvQuery subQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, step, 1, query.getAggregation(), query.getOrderBy(), false); futures.add(findAndAggregateAsync(entityId, subQuery, toPartitionTs(startTs), toPartitionTs(endTs))); stepTs = endTs; } @@ -400,15 +400,6 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem return Futures.immediateFuture(false); }, readResultsProcessingExecutor); - - ListenableFuture savedLatestFuture = Futures.transform(booleanFuture, - (AsyncFunction) isRemove -> { - if (isRemove) { - return getNewLatestEntryFuture(entityId, query); - } - return Futures.immediateFuture(null); - }, readResultsProcessingExecutor); - ListenableFuture removedLatestFuture = Futures.transform(booleanFuture, (AsyncFunction) isRemove -> { if (isRemove) { @@ -416,15 +407,27 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } return Futures.immediateFuture(null); }, readResultsProcessingExecutor); - return Futures.transform(Futures.allAsList(Arrays.asList(savedLatestFuture, removedLatestFuture)), - (AsyncFunction, Void>) list -> Futures.immediateFuture(null), readResultsProcessingExecutor); + + if (query.getRewriteLatestIfDeleted()) { + ListenableFuture savedLatestFuture = Futures.transform(booleanFuture, + (AsyncFunction) isRemove -> { + if (isRemove) { + return getNewLatestEntryFuture(entityId, query); + } + return Futures.immediateFuture(null); + }, readResultsProcessingExecutor); + + return Futures.transform(Futures.allAsList(Arrays.asList(savedLatestFuture, removedLatestFuture)), + (AsyncFunction, Void>) list -> Futures.immediateFuture(null), readResultsProcessingExecutor); + } + return removedLatestFuture; } private ListenableFuture getNewLatestEntryFuture(EntityId entityId, TsKvQuery query) { long startTs = 0; long endTs = query.getStartTs() - 1; TsKvQuery findNewLatestQuery = new BaseTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, - Aggregation.NONE, DESC_ORDER); + Aggregation.NONE, DESC_ORDER, false); ListenableFuture> future = findAllAsync(entityId, findNewLatestQuery); return Futures.transform(future, (AsyncFunction, Void>) entryList -> { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index 181a19743b..a8d022ba2f 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -127,7 +127,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { entries.add(save(deviceId, 55000, 600)); List list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.NONE, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.NONE, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(55000, list.get(0).getTs()); assertEquals(java.util.Optional.of(600L), list.get(0).getLongValue()); @@ -139,7 +139,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(400L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.AVG, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.AVG, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); assertEquals(java.util.Optional.of(150L), list.get(0).getLongValue()); @@ -151,7 +151,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(550L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.SUM, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.SUM, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -164,7 +164,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(1100L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.MIN, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.MIN, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -177,7 +177,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(500L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.MAX, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.MAX, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); @@ -190,7 +190,7 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { assertEquals(java.util.Optional.of(600L), list.get(2).getLongValue()); list = tsService.findAll(deviceId, Collections.singletonList(new BaseTsKvQuery(LONG_KEY, 0, - 60000, 20000, 3, Aggregation.COUNT, DESC_ORDER))).get(); + 60000, 20000, 3, Aggregation.COUNT, DESC_ORDER, false))).get(); assertEquals(3, list.size()); assertEquals(10000, list.get(0).getTs()); diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java index 0c7e3874ee..bb0813ff77 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryRestMsgHandler.java @@ -28,7 +28,6 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityIdFactory; -import org.thingsboard.server.common.data.id.UUIDBased; import org.thingsboard.server.common.data.kv.*; import org.thingsboard.server.common.msg.core.TelemetryUploadRequest; import org.thingsboard.server.common.transport.adaptor.JsonConverter; @@ -138,9 +137,10 @@ public class TelemetryRestMsgHandler extends DefaultRestMsgHandler { // If interval is 0, convert this to a NONE aggregation, which is probably what the user really wanted Aggregation agg = (interval.isPresent() && interval.get() == 0) ? Aggregation.valueOf(Aggregation.NONE.name()) : - Aggregation.valueOf(request.getParameter("agg", Aggregation.NONE.name())); + Aggregation.valueOf(request.getParameter("agg", Aggregation.NONE.name())); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs.get(), endTs.get(), interval.get(), limit.orElse(TelemetryWebsocketMsgHandler.DEFAULT_LIMIT), agg, "DESC")) + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs.get(), endTs.get(), + interval.get(), limit.orElse(TelemetryWebsocketMsgHandler.DEFAULT_LIMIT), agg, "DESC", false)) .collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, getTsKvListCallback(msg)); } else { @@ -218,7 +218,7 @@ public class TelemetryRestMsgHandler extends DefaultRestMsgHandler { } private boolean handleHttpPostAttributes(PluginContext ctx, PluginRestMsg msg, RestRequest request, - EntityId entityId, String scope) throws ServletException, IOException { + EntityId entityId, String scope) throws ServletException, IOException { if (DataConstants.SERVER_SCOPE.equals(scope) || DataConstants.SHARED_SCOPE.equals(scope)) { JsonNode jsonNode; @@ -274,7 +274,7 @@ public class TelemetryRestMsgHandler extends DefaultRestMsgHandler { } } }); - return attributes; + return attributes; } private void handleHttpPostTimeseries(PluginContext ctx, PluginRestMsg msg, RestRequest request, EntityId entityId, long ttl) { diff --git a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java index bf75c5de29..c024b673e3 100644 --- a/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java +++ b/extensions-core/src/main/java/org/thingsboard/server/extensions/core/plugin/telemetry/handlers/TelemetryWebsocketMsgHandler.java @@ -217,7 +217,8 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { log.debug("[{}] fetching timeseries data for last {} ms for keys: ({}) for device : {}", sessionId, cmd.getTimeWindow(), cmd.getKeys(), entityId); startTs = cmd.getStartTs(); long endTs = cmd.getStartTs() + cmd.getTimeWindow(); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs, endTs, cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY)).collect(Collectors.toList()); + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, startTs, endTs, cmd.getInterval(), + getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY, false)).collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, getSubscriptionCallback(sessionRef, cmd, sessionId, entityId, startTs, keys)); } else { List keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); @@ -301,7 +302,8 @@ public class TelemetryWebsocketMsgHandler extends DefaultWebsocketMsgHandler { } EntityId entityId = EntityIdFactory.getByTypeAndId(cmd.getEntityType(), cmd.getEntityId()); List keys = new ArrayList<>(getKeys(cmd).orElse(Collections.emptySet())); - List queries = keys.stream().map(key -> new BaseTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY)) + List queries = keys.stream().map(key -> new BaseTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), + cmd.getInterval(), getLimit(cmd.getLimit()), getAggregation(cmd.getAgg()), ORDER_BY, false)) .collect(Collectors.toList()); ctx.loadTimeseries(entityId, queries, new PluginCallback>() { @Override From ce9967488b72e28c0d79d252e57422624b971a16 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Tue, 5 Jun 2018 18:25:39 +0300 Subject: [PATCH 5/7] jpa delete timeseries implementation, partitions delete fix --- .../server/controller/DeviceController.java | 56 ------------------- .../dao/sql/timeseries/JpaTimeseriesDao.java | 21 +++++-- .../dao/sql/timeseries/TsKvRepository.java | 45 +++++++++------ .../CassandraBaseTimeseriesDao.java | 45 +++++++++------ .../timeseries/BaseTimeseriesServiceTest.java | 17 +++--- 5 files changed, 84 insertions(+), 100 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java index 856ca8e399..79557f6047 100644 --- a/application/src/main/java/org/thingsboard/server/controller/DeviceController.java +++ b/application/src/main/java/org/thingsboard/server/controller/DeviceController.java @@ -16,7 +16,6 @@ package org.thingsboard.server.controller; import com.google.common.util.concurrent.ListenableFuture; -import org.springframework.beans.factory.annotation.Autowired; import org.springframework.http.HttpStatus; import org.springframework.security.access.prepost.PreAuthorize; import org.springframework.web.bind.annotation.*; @@ -29,23 +28,17 @@ import org.thingsboard.server.common.data.device.DeviceSearchQuery; import org.thingsboard.server.common.data.id.CustomerId; import org.thingsboard.server.common.data.id.DeviceId; import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.data.kv.Aggregation; -import org.thingsboard.server.common.data.kv.BaseTsKvQuery; -import org.thingsboard.server.common.data.kv.BasicTsKvEntry; -import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.data.security.Authority; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.dao.exception.IncorrectParameterException; import org.thingsboard.server.dao.model.ModelConstants; -import org.thingsboard.server.dao.timeseries.TimeseriesService; import org.thingsboard.server.exception.ThingsboardErrorCode; import org.thingsboard.server.exception.ThingsboardException; import org.thingsboard.server.service.security.model.SecurityUser; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.stream.Collectors; @@ -53,9 +46,6 @@ import java.util.stream.Collectors; @RequestMapping("/api") public class DeviceController extends BaseController { - @Autowired - protected TimeseriesService timeseriesService; - public static final String DEVICE_ID = "deviceId"; @PreAuthorize("hasAnyAuthority('TENANT_ADMIN', 'CUSTOMER_USER')") @@ -377,50 +367,4 @@ public class DeviceController extends BaseController { throw handleException(e); } } - - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/device/testSave", method = RequestMethod.GET) - @ResponseBody - public void testSave() throws ThingsboardException { - try { - SecurityUser user = getCurrentUser(); - TenantId tenantId = user.getTenantId(); - - Device device = deviceService.findDeviceByTenantIdAndName(tenantId, "Test"); - - timeseriesService.save(device.getId(), new BasicTsKvEntry(1516892633000L, - new LongDataEntry("test", 1L))).get(); - timeseriesService.save(device.getId(), new BasicTsKvEntry(1519571033000L, - new LongDataEntry("test", 2L))).get(); - timeseriesService.save(device.getId(), new BasicTsKvEntry(1521990233000L, - new LongDataEntry("test", 3L))).get(); - timeseriesService.save(device.getId(), new BasicTsKvEntry(1524668633000L, - new LongDataEntry("test", 4L))).get(); - timeseriesService.save(device.getId(), new BasicTsKvEntry(1527260633000L, - new LongDataEntry("test", 5L))).get(); - - } catch (Exception e) { - throw handleException(e); - } - } - - @PreAuthorize("hasAnyAuthority('TENANT_ADMIN')") - @RequestMapping(value = "/device/testDelete", method = RequestMethod.GET) - @ResponseBody - public void testDelete() throws ThingsboardException { - try { - SecurityUser user = getCurrentUser(); - TenantId tenantId = user.getTenantId(); - - Device device = deviceService.findDeviceByTenantIdAndName(tenantId, "Test"); - - long startTs = 1519561033000L; - long endTs = 1528260633000L; - timeseriesService.remove(device.getId(), Collections.singletonList(new BaseTsKvQuery("test", - startTs, endTs, endTs - startTs, 0, Aggregation.NONE, "DESC", true))).get(); - - } catch (Exception e) { - throw handleException(e); - } - } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java index bef82c3677..7915e844e2 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java @@ -300,14 +300,27 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp @Override public ListenableFuture remove(EntityId entityId, TsKvQuery query) { - //TODO: implement - return null; + return insertService.submit(() -> { + tsKvRepository.delete( + fromTimeUUID(entityId.getId()), + entityId.getEntityType(), + query.getKey(), + query.getStartTs(), + query.getEndTs()); + return null; + }); } @Override public ListenableFuture removeLatest(EntityId entityId, TsKvQuery query) { - //TODO: implement - return null; + TsKvLatestEntity latestEntity = new TsKvLatestEntity(); + latestEntity.setEntityType(entityId.getEntityType()); + latestEntity.setEntityId(fromTimeUUID(entityId.getId())); + latestEntity.setKey(query.getKey()); + return insertService.submit(() -> { + tsKvLatestRepository.delete(latestEntity); + return null; + }); } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java index a1d19208b5..2b39d2596e 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/TsKvRepository.java @@ -16,10 +16,12 @@ package org.thingsboard.server.dao.sql.timeseries; import org.springframework.data.domain.Pageable; +import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Query; import org.springframework.data.repository.CrudRepository; import org.springframework.data.repository.query.Param; import org.springframework.scheduling.annotation.Async; +import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.dao.model.sql.TsKvCompositeKey; import org.thingsboard.server.dao.model.sql.TsKvEntity; @@ -41,6 +43,17 @@ public interface TsKvRepository extends CrudRepository :startTs AND tskv.ts < :endTs") + void delete(@Param("entityId") String entityId, + @Param("entityType") EntityType entityType, + @Param("entityKey") String key, + @Param("startTs") long startTs, + @Param("endTs") long endTs); + @Async @Query("SELECT new TsKvEntity(MAX(tskv.strValue), MAX(tskv.longValue), MAX(tskv.doubleValue)) FROM TsKvEntity tskv " + "WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " + @@ -56,30 +69,30 @@ public interface TsKvRepository extends CrudRepository :startTs AND tskv.ts < :endTs") CompletableFuture findMin(@Param("entityId") String entityId, - @Param("entityType") EntityType entityType, - @Param("entityKey") String entityKey, - @Param("startTs") long startTs, - @Param("endTs") long endTs); + @Param("entityType") EntityType entityType, + @Param("entityKey") String entityKey, + @Param("startTs") long startTs, + @Param("endTs") long endTs); @Async @Query("SELECT new TsKvEntity(COUNT(tskv.booleanValue), COUNT(tskv.strValue), COUNT(tskv.longValue), COUNT(tskv.doubleValue)) FROM TsKvEntity tskv " + "WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " + "AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts < :endTs") CompletableFuture findCount(@Param("entityId") String entityId, - @Param("entityType") EntityType entityType, - @Param("entityKey") String entityKey, - @Param("startTs") long startTs, - @Param("endTs") long endTs); + @Param("entityType") EntityType entityType, + @Param("entityKey") String entityKey, + @Param("startTs") long startTs, + @Param("endTs") long endTs); @Async @Query("SELECT new TsKvEntity(AVG(tskv.longValue), AVG(tskv.doubleValue)) FROM TsKvEntity tskv " + "WHERE tskv.entityId = :entityId AND tskv.entityType = :entityType " + "AND tskv.key = :entityKey AND tskv.ts > :startTs AND tskv.ts < :endTs") CompletableFuture findAvg(@Param("entityId") String entityId, - @Param("entityType") EntityType entityType, - @Param("entityKey") String entityKey, - @Param("startTs") long startTs, - @Param("endTs") long endTs); + @Param("entityType") EntityType entityType, + @Param("entityKey") String entityKey, + @Param("startTs") long startTs, + @Param("endTs") long endTs); @Async @@ -87,8 +100,8 @@ public interface TsKvRepository extends CrudRepository :startTs AND tskv.ts < :endTs") CompletableFuture findSum(@Param("entityId") String entityId, - @Param("entityType") EntityType entityType, - @Param("entityKey") String entityKey, - @Param("startTs") long startTs, - @Param("endTs") long endTs); + @Param("entityType") EntityType entityType, + @Param("entityKey") String entityKey, + @Param("startTs") long startTs, + @Param("endTs") long endTs); } 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 3e4e2bdac5..db61088b39 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 @@ -441,7 +441,7 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem } private ListenableFuture deleteLatest(EntityId entityId, String key) { - Statement delete = QueryBuilder.delete().from(ModelConstants.TS_KV_LATEST_CF) + Statement delete = QueryBuilder.delete().all().from(ModelConstants.TS_KV_LATEST_CF) .where(eq(ModelConstants.ENTITY_TYPE_COLUMN, entityId.getEntityType())) .and(eq(ModelConstants.ENTITY_ID_COLUMN, entityId.getId())) .and(eq(ModelConstants.KEY_COLUMN, key)); @@ -453,25 +453,36 @@ public class CassandraBaseTimeseriesDao extends CassandraAbstractAsyncDao implem public ListenableFuture removePartition(EntityId entityId, TsKvQuery query) { long minPartition = toPartitionTs(query.getStartTs()); long maxPartition = toPartitionTs(query.getEndTs()); + if (minPartition == maxPartition) { + return Futures.immediateFuture(null); + } else { + ResultSetFuture partitionsFuture = fetchPartitions(entityId, query.getKey(), minPartition, maxPartition); - ResultSetFuture partitionsFuture = fetchPartitions(entityId, query.getKey(), minPartition, maxPartition); + final SimpleListenableFuture resultFuture = new SimpleListenableFuture<>(); + final ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor); - final SimpleListenableFuture resultFuture = new SimpleListenableFuture<>(); - final ListenableFuture> partitionsListFuture = Futures.transform(partitionsFuture, getPartitionsArrayFunction(), readResultsProcessingExecutor); - - Futures.addCallback(partitionsListFuture, new FutureCallback>() { - @Override - public void onSuccess(@Nullable List partitions) { - TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions); - deletePartitionAsync(cursor, resultFuture); - } + Futures.addCallback(partitionsListFuture, new FutureCallback>() { + @Override + public void onSuccess(@Nullable List partitions) { + int index = 0; + if (minPartition != query.getStartTs()) { + index = 1; + } + List partitionsToDelete = new ArrayList<>(); + for (int i = index; i < partitions.size() - 1; i++) { + partitionsToDelete.add(partitions.get(i)); + } + TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitionsToDelete); + deletePartitionAsync(cursor, resultFuture); + } - @Override - public void onFailure(Throwable t) { - log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); - } - }, readResultsProcessingExecutor); - return resultFuture; + @Override + public void onFailure(Throwable t) { + log.error("[{}][{}] Failed to fetch partitions for interval {}-{}", entityId.getEntityType().name(), entityId.getId(), minPartition, maxPartition, t); + } + }, readResultsProcessingExecutor); + return resultFuture; + } } private void deletePartitionAsync(final TsKvQueryCursor cursor, final SimpleListenableFuture resultFuture) { diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index a8d022ba2f..b3a742cda9 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -93,24 +93,27 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { Assert.assertEquals(toTsEntry(TS, stringKvEntry), entries.get(0)); } - //TODO: sql delete implement - /*@Test + @Test public void testDeleteDeviceTsData() throws Exception { DeviceId deviceId = new DeviceId(UUIDs.timeBased()); + saveEntries(deviceId, TS - 4); saveEntries(deviceId, TS - 3); saveEntries(deviceId, TS - 2); saveEntries(deviceId, TS - 1); - saveEntries(deviceId, TS); tsService.remove(deviceId, Collections.singletonList( - new BaseTsKvQuery(STRING_KEY, TS - 4, TS - 2))).get(); + new BaseTsKvQuery(STRING_KEY, TS - 4, TS, 60000, 0, Aggregation.NONE, DESC_ORDER, + false))).get(); List list = tsService.findAll(deviceId, Collections.singletonList( - new BaseTsKvQuery(STRING_KEY, 0, 60000, 60000, 5, Aggregation.NONE, DESC_ORDER))).get(); + new BaseTsKvQuery(STRING_KEY, 0, 60000, 60000, 5, Aggregation.NONE, DESC_ORDER, + false))).get(); + Assert.assertEquals(1, list.size()); - Assert.assertEquals(2, list.size()); - }*/ + List latest = tsService.findLatest(deviceId, Collections.singletonList(STRING_KEY)).get(); + Assert.assertEquals(null, latest.get(0).getValueAsString()); + } @Test public void testFindDeviceTsData() throws Exception { From a58eb0f74ed23326ea8f954d1b504e93ab88c71f Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Tue, 5 Jun 2018 18:36:49 +0300 Subject: [PATCH 6/7] typo --- application/src/main/resources/thingsboard.yml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index aef92804f1..c47fc28fbd 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -66,8 +66,8 @@ plugins: # JWT Token parameters security.jwt: - tokenExpirationTime: "${JWT_TOKEN_EXPIRATION_TIME:9000000}" # Number of seconds (15 mins) - refreshTokenExpTime: "${JWT_REFRESH_TOKEN_EXPIRATION_TIME:36000000}" # Seconds (1 hour) + tokenExpirationTime: "${JWT_TOKEN_EXPIRATION_TIME:900}" # Number of seconds (15 mins) + refreshTokenExpTime: "${JWT_REFRESH_TOKEN_EXPIRATION_TIME:3600}" # Seconds (1 hour) tokenIssuer: "${JWT_TOKEN_ISSUER:thingsboard.io}" tokenSigningKey: "${JWT_TOKEN_SIGNING_KEY:thingsboardDefaultSigningKey}" @@ -133,7 +133,7 @@ quota: intervalMin: 2 database: - type: "${DATABASE_TYPE:cassandra}" # cassandra OR sql + type: "${DATABASE_TYPE:sql}" # cassandra OR sql # Cassandra driver configuration parameters cassandra: From 712e1ae9b78764c581d8164f009b862709fb2451 Mon Sep 17 00:00:00 2001 From: Dima Landiak Date: Wed, 6 Jun 2018 18:07:16 +0300 Subject: [PATCH 7/7] test fix --- .../server/dao/sql/timeseries/JpaTimeseriesDao.java | 6 +++--- .../timeseries/BaseTimeseriesServiceTest.java | 12 ++++++------ 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java index 7915e844e2..d87e7c5af5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/timeseries/JpaTimeseriesDao.java @@ -300,7 +300,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp @Override public ListenableFuture remove(EntityId entityId, TsKvQuery query) { - return insertService.submit(() -> { + return service.submit(() -> { tsKvRepository.delete( fromTimeUUID(entityId.getId()), entityId.getEntityType(), @@ -317,7 +317,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp latestEntity.setEntityType(entityId.getEntityType()); latestEntity.setEntityId(fromTimeUUID(entityId.getId())); latestEntity.setKey(query.getKey()); - return insertService.submit(() -> { + return service.submit(() -> { tsKvLatestRepository.delete(latestEntity); return null; }); @@ -325,7 +325,7 @@ public class JpaTimeseriesDao extends JpaAbstractDaoListeningExecutorService imp @Override public ListenableFuture removePartition(EntityId entityId, TsKvQuery query) { - return insertService.submit(() -> null); + return service.submit(() -> null); } @PreDestroy diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java index b3a742cda9..43f107c427 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/timeseries/BaseTimeseriesServiceTest.java @@ -97,17 +97,17 @@ public abstract class BaseTimeseriesServiceTest extends AbstractServiceTest { public void testDeleteDeviceTsData() throws Exception { DeviceId deviceId = new DeviceId(UUIDs.timeBased()); - saveEntries(deviceId, TS - 4); - saveEntries(deviceId, TS - 3); - saveEntries(deviceId, TS - 2); - saveEntries(deviceId, TS - 1); + saveEntries(deviceId, 10000); + saveEntries(deviceId, 20000); + saveEntries(deviceId, 30000); + saveEntries(deviceId, 40000); tsService.remove(deviceId, Collections.singletonList( - new BaseTsKvQuery(STRING_KEY, TS - 4, TS, 60000, 0, Aggregation.NONE, DESC_ORDER, + new BaseTsKvQuery(STRING_KEY, 15000, 45000, 10000, 0, Aggregation.NONE, DESC_ORDER, false))).get(); List list = tsService.findAll(deviceId, Collections.singletonList( - new BaseTsKvQuery(STRING_KEY, 0, 60000, 60000, 5, Aggregation.NONE, DESC_ORDER, + new BaseTsKvQuery(STRING_KEY, 5000, 45000, 10000, 10, Aggregation.NONE, DESC_ORDER, false))).get(); Assert.assertEquals(1, list.size());