Browse Source

Merge branch 'master' into feature/latestTsAggregation

pull/7288/head
Andrii Shvaika 4 years ago
parent
commit
74a857cbe2
  1. 43
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 15
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  3. 4
      common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java
  4. 37
      common/data/src/main/java/org/thingsboard/server/common/data/kv/AggTsKvEntry.java
  5. 2
      common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java
  6. 46
      common/data/src/main/java/org/thingsboard/server/common/data/kv/ReadTsKvQueryResult.java
  7. 8
      common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java
  8. 26
      common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntryAggWrapper.java
  9. 10
      common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java
  10. 6
      dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java
  11. 20
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java
  12. 10
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/ts/TimescaleTsKvEntity.java
  13. 13
      dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java
  14. 165
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java
  15. 13
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java
  16. 3
      dao/src/main/java/org/thingsboard/server/dao/sqlts/AggregationTimeseriesDao.java
  17. 22
      dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java
  18. 4
      dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java
  19. 98
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java
  20. 134
      dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java
  21. 37
      dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java
  22. 70
      dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java
  23. 19
      dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java
  24. 59
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java
  25. 6
      dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java
  26. 4
      dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java
  27. 5
      dao/src/test/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDaoTest.java

43
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -34,6 +34,7 @@ import org.thingsboard.server.common.data.id.CustomerId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.query.AlarmDataQuery; import org.thingsboard.server.common.data.query.AlarmDataQuery;
@ -83,6 +84,7 @@ import java.util.concurrent.ThreadFactory;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@SuppressWarnings("UnstableApiUsage")
@Slf4j @Slf4j
@TbCoreComponent @TbCoreComponent
@Service @Service
@ -429,23 +431,34 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
} else { } else {
finalTsKvQueryList = tsKvQueryList; finalTsKvQueryList = tsKvQueryList;
} }
Map<EntityData, ListenableFuture<List<TsKvEntry>>> fetchResultMap = new HashMap<>(); Map<EntityData, ListenableFuture<List<ReadTsKvQueryResult>>> fetchResultMap = new HashMap<>();
ctx.getData().getData().forEach(entityData -> fetchResultMap.put(entityData, List<EntityData> entityDataList = ctx.getData().getData();
tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList))); entityDataList.forEach(entityData -> fetchResultMap.put(entityData,
tsService.findAllByQueries(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList)));
return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> {
// Map that holds last ts for each key for each entity.
Map<EntityData, Map<String, Long>> lastTsEntityMap = new HashMap<>();
fetchResultMap.forEach((entityData, future) -> { fetchResultMap.forEach((entityData, future) -> {
Map<String, List<TsValue>> keyData = new LinkedHashMap<>();
cmd.getKeys().forEach(key -> keyData.put(key, new ArrayList<>()));
try { try {
List<TsKvEntry> entityTsData = future.get(); Map<String, Long> lastTsMap = new HashMap<>();
if (entityTsData != null) { lastTsEntityMap.put(entityData, lastTsMap);
entityTsData.forEach(entry -> keyData.get(entry.getKey()).add(new TsValue(entry.getTs(), entry.getValueAsString())));
List<ReadTsKvQueryResult> queryResults = future.get();
if (queryResults != null) {
for (ReadTsKvQueryResult queryResult : queryResults) {
entityData.getTimeseries().put(queryResult.getKey(), queryResult.toTsValues());
lastTsMap.put(queryResult.getKey(), queryResult.getLastEntryTs());
}
} }
keyData.forEach((k, v) -> entityData.getTimeseries().put(k, v.toArray(new TsValue[v.size()]))); // Populate with empty values if no data found.
cmd.getKeys().forEach(key -> {
if (!entityData.getTimeseries().containsKey(key)) {
entityData.getTimeseries().put(key, new TsValue[0]);
}
});
if (cmd.isFetchLatestPreviousPoint()) { if (cmd.isFetchLatestPreviousPoint()) {
entityData.getTimeseries().values().forEach(dataArray -> { entityData.getTimeseries().values().forEach(dataArray -> Arrays.sort(dataArray, (o1, o2) -> Long.compare(o2.getTs(), o1.getTs())));
Arrays.sort(dataArray, (o1, o2) -> Long.compare(o2.getTs(), o1.getTs()));
});
} }
} catch (InterruptedException | ExecutionException e) { } catch (InterruptedException | ExecutionException e) {
log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e); log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e);
@ -459,13 +472,13 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription());
ctx.setInitialDataSent(true); ctx.setInitialDataSent(true);
} else { } else {
update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); update = new EntityDataUpdate(ctx.getCmdId(), null, entityDataList, ctx.getMaxEntitiesPerDataSubscription());
} }
if (subscribe) { if (subscribe) {
ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs()); ctx.createTimeseriesSubscriptions(lastTsEntityMap, cmd.getStartTs(), cmd.getEndTs());
} }
ctx.sendWsMsg(update); ctx.sendWsMsg(update);
ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); entityDataList.forEach(ed -> ed.getTimeseries().clear());
} finally { } finally {
ctx.getWsLock().unlock(); ctx.getWsLock().unlock();
} }

15
application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java

@ -86,7 +86,7 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
} else { } else {
oldDataMap = Collections.emptyMap(); oldDataMap = Collections.emptyMap();
} }
Map<EntityId, EntityData> newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity(), (a,b)-> a)); Map<EntityId, EntityData> newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity(), (a, b) -> a));
if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) {
log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId);
} else { } else {
@ -122,8 +122,13 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
createSubscriptions(keys, true, 0, 0); createSubscriptions(keys, true, 0, 0);
} }
public void createTimeseriesSubscriptions(List<EntityKey> keys, long startTs, long endTs) { public void createTimeseriesSubscriptions(Map<EntityData, Map<String, Long>> entityKeyStates, long startTs, long endTs) {
createSubscriptions(keys, false, startTs, endTs); entityKeyStates.forEach((entityData, keyStates) -> {
int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet();
subToEntityIdMap.put(subIdx, entityData.getEntityId());
localSubscriptionService.addSubscription(
createTsSub(entityData, subIdx, false, startTs, endTs, keyStates));
});
} }
private void createSubscriptions(List<EntityKey> keys, boolean latestValues, long startTs, long endTs) { private void createSubscriptions(List<EntityKey> keys, boolean latestValues, long startTs, long endTs) {
@ -191,6 +196,10 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
keyStates.put(k, ts); keyStates.put(k, ts);
}); });
} }
return createTsSub(entityData, subIdx, latestValues, startTs, endTs, keyStates);
}
private TbTimeseriesSubscription createTsSub(EntityData entityData, int subIdx, boolean latestValues, long startTs, long endTs, Map<String, Long> keyStates) {
log.trace("[{}][{}][{}] Creating time-series subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); log.trace("[{}][{}][{}] Creating time-series subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates);
return TbTimeseriesSubscription.builder() return TbTimeseriesSubscription.builder()
.serviceId(serviceId) .serviceId(serviceId)

4
common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java

@ -21,17 +21,21 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.Collection; import java.util.Collection;
import java.util.List; import java.util.List;
import java.util.Map;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
public interface TimeseriesService { public interface TimeseriesService {
ListenableFuture<List<ReadTsKvQueryResult>> findAllByQueries(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries);
ListenableFuture<List<TsKvEntry>> findAll(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries); ListenableFuture<List<TsKvEntry>> findAll(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries);
ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys); ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys);

37
common/data/src/main/java/org/thingsboard/server/common/data/kv/AggTsKvEntry.java

@ -0,0 +1,37 @@
/**
* Copyright © 2016-2022 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.common.data.kv;
import lombok.ToString;
import org.thingsboard.server.common.data.query.TsValue;
@ToString
public class AggTsKvEntry extends BasicTsKvEntry {
private static final long serialVersionUID = -1933884317450255935L;
private final long count;
public AggTsKvEntry(long ts, KvEntry kv, long count) {
super(ts, kv);
this.count = count;
}
@Override
public TsValue toTsValue() {
return new TsValue(ts, getValueAsString(), count);
}
}

2
common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java

@ -20,7 +20,7 @@ import java.util.Optional;
public class BasicTsKvEntry implements TsKvEntry { public class BasicTsKvEntry implements TsKvEntry {
private static final int MAX_CHARS_PER_DATA_POINT = 512; private static final int MAX_CHARS_PER_DATA_POINT = 512;
private final long ts; protected final long ts;
private final KvEntry kv; private final KvEntry kv;
public BasicTsKvEntry(long ts, KvEntry kv) { public BasicTsKvEntry(long ts, KvEntry kv) {

46
common/data/src/main/java/org/thingsboard/server/common/data/kv/ReadTsKvQueryResult.java

@ -0,0 +1,46 @@
/**
* Copyright © 2016-2022 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.common.data.kv;
import lombok.Data;
import org.thingsboard.server.common.data.query.TsValue;
import java.util.ArrayList;
import java.util.List;
@Data
public class ReadTsKvQueryResult {
private final String key;
// Holds the data list;
private final List<TsKvEntry> data;
// Holds the max ts of the records that match aggregation intervals (not the ts of the aggregation window, but the ts of the last record among all the intervals)
private final long lastEntryTs;
public TsValue[] toTsValues() {
if (data != null && !data.isEmpty()) {
List<TsValue> queryValues = new ArrayList<>();
for (TsKvEntry v : data) {
queryValues.add(v.toTsValue()); // TODO: add count here.
}
return queryValues.toArray(new TsValue[queryValues.size()]);
} else {
return new TsValue[0];
}
}
}

8
common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java

@ -16,10 +16,11 @@
package org.thingsboard.server.common.data.kv; package org.thingsboard.server.common.data.kv;
import com.fasterxml.jackson.annotation.JsonIgnore; import com.fasterxml.jackson.annotation.JsonIgnore;
import org.thingsboard.server.common.data.query.TsValue;
/** /**
* Represents time series KV data entry * Represents time series KV data entry
* *
* @author ashvayka * @author ashvayka
* *
*/ */
@ -30,4 +31,9 @@ public interface TsKvEntry extends KvEntry {
@JsonIgnore @JsonIgnore
int getDataPoints(); int getDataPoints();
@JsonIgnore
default TsValue toTsValue() {
return new TsValue(getTs(), getValueAsString());
}
} }

26
common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntryAggWrapper.java

@ -0,0 +1,26 @@
/**
* Copyright © 2016-2022 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.common.data.kv;
import lombok.Data;
@Data
public class TsKvEntryAggWrapper {
private final TsKvEntry entry;
private final long lastEntryTs;
}

10
common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java

@ -15,11 +15,21 @@
*/ */
package org.thingsboard.server.common.data.query; package org.thingsboard.server.common.data.query;
import com.fasterxml.jackson.annotation.JsonInclude;
import lombok.Data; import lombok.Data;
import lombok.RequiredArgsConstructor;
@Data @Data
@RequiredArgsConstructor
@JsonInclude(JsonInclude.Include.NON_NULL)
public class TsValue { public class TsValue {
private final long ts; private final long ts;
private final String value; private final String value;
private final Long count;
public TsValue(long ts, String value) {
this(ts, value, null);
}
} }

6
dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java

@ -634,11 +634,11 @@ public class ModelConstants {
protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN)}; protected static final String[] COUNT_AGGREGATION_COLUMNS = new String[]{count(LONG_VALUE_COLUMN), count(DOUBLE_VALUE_COLUMN), count(BOOLEAN_VALUE_COLUMN), count(STRING_VALUE_COLUMN), count(JSON_VALUE_COLUMN)};
protected static final String[] MIN_AGGREGATION_COLUMNS = protected static final String[] MIN_AGGREGATION_COLUMNS =
ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{min(LONG_VALUE_COLUMN), min(DOUBLE_VALUE_COLUMN), min(BOOLEAN_VALUE_COLUMN), min(STRING_VALUE_COLUMN), min(JSON_VALUE_COLUMN)}); ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{min(LONG_VALUE_COLUMN), min(DOUBLE_VALUE_COLUMN), min(BOOLEAN_VALUE_COLUMN), min(STRING_VALUE_COLUMN), min(JSON_VALUE_COLUMN), max(TS_COLUMN)});
protected static final String[] MAX_AGGREGATION_COLUMNS = protected static final String[] MAX_AGGREGATION_COLUMNS =
ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{max(LONG_VALUE_COLUMN), max(DOUBLE_VALUE_COLUMN), max(BOOLEAN_VALUE_COLUMN), max(STRING_VALUE_COLUMN), max(JSON_VALUE_COLUMN)}); ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{max(LONG_VALUE_COLUMN), max(DOUBLE_VALUE_COLUMN), max(BOOLEAN_VALUE_COLUMN), max(STRING_VALUE_COLUMN), max(JSON_VALUE_COLUMN), max(TS_COLUMN)});
protected static final String[] SUM_AGGREGATION_COLUMNS = protected static final String[] SUM_AGGREGATION_COLUMNS =
ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{sum(LONG_VALUE_COLUMN), sum(DOUBLE_VALUE_COLUMN)}); ArrayUtils.addAll(COUNT_AGGREGATION_COLUMNS, new String[]{sum(LONG_VALUE_COLUMN), sum(DOUBLE_VALUE_COLUMN), max(TS_COLUMN)});
protected static final String[] AVG_AGGREGATION_COLUMNS = SUM_AGGREGATION_COLUMNS; protected static final String[] AVG_AGGREGATION_COLUMNS = SUM_AGGREGATION_COLUMNS;
public static String min(String s) { public static String min(String s) {

20
dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java

@ -16,6 +16,7 @@
package org.thingsboard.server.dao.model.sql; package org.thingsboard.server.dao.model.sql;
import lombok.Data; import lombok.Data;
import org.thingsboard.server.common.data.kv.AggTsKvEntry;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry;
import org.thingsboard.server.common.data.kv.DoubleDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry;
@ -80,6 +81,18 @@ public abstract class AbstractTsKvEntity implements ToData<TsKvEntry> {
@Transient @Transient
protected String strKey; protected String strKey;
@Transient
protected Long aggValuesLastTs;
@Transient
protected Long aggValuesCount;
public AbstractTsKvEntity() {
}
public AbstractTsKvEntity(Long aggValuesLastTs) {
this.aggValuesLastTs = aggValuesLastTs;
}
public abstract boolean isNotEmpty(); public abstract boolean isNotEmpty();
protected static boolean isAllNull(Object... args) { protected static boolean isAllNull(Object... args) {
@ -105,7 +118,12 @@ public abstract class AbstractTsKvEntity implements ToData<TsKvEntry> {
} else if (jsonValue != null) { } else if (jsonValue != null) {
kvEntry = new JsonDataEntry(strKey, jsonValue); kvEntry = new JsonDataEntry(strKey, jsonValue);
} }
return new BasicTsKvEntry(ts, kvEntry);
if (aggValuesCount == null) {
return new BasicTsKvEntry(ts, kvEntry);
} else {
return new AggTsKvEntry(ts, kvEntry, aggValuesCount);
}
} }
} }

10
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/ts/TimescaleTsKvEntity.java

@ -62,6 +62,7 @@ import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.F
@ColumnResult(name = "doubleCountValue", type = Long.class), @ColumnResult(name = "doubleCountValue", type = Long.class),
@ColumnResult(name = "strValue", type = String.class), @ColumnResult(name = "strValue", type = String.class),
@ColumnResult(name = "aggType", type = String.class), @ColumnResult(name = "aggType", type = String.class),
@ColumnResult(name = "maxAggTs", type = Long.class),
} }
), ),
}), }),
@ -78,6 +79,8 @@ import static org.thingsboard.server.dao.sqlts.timescale.AggregationRepository.F
@ColumnResult(name = "longValueCount", type = Long.class), @ColumnResult(name = "longValueCount", type = Long.class),
@ColumnResult(name = "doubleValueCount", type = Long.class), @ColumnResult(name = "doubleValueCount", type = Long.class),
@ColumnResult(name = "jsonValueCount", type = Long.class), @ColumnResult(name = "jsonValueCount", type = Long.class),
@ColumnResult(name = "jsonValueCount", type = Long.class),
@ColumnResult(name = "maxAggTs", type = Long.class),
} }
) )
}), }),
@ -114,7 +117,8 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity {
public TimescaleTsKvEntity() { public TimescaleTsKvEntity() {
} }
public TimescaleTsKvEntity(Long tsBucket, Long interval, Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String strValue, String aggType) { public TimescaleTsKvEntity(Long tsBucket, Long interval, Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String strValue, String aggType, Long aggValuesLastTs) {
super(aggValuesLastTs);
if (!StringUtils.isEmpty(strValue)) { if (!StringUtils.isEmpty(strValue)) {
this.strValue = strValue; this.strValue = strValue;
} }
@ -135,6 +139,7 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity {
} else { } else {
this.doubleValue = 0.0; this.doubleValue = 0.0;
} }
this.aggValuesCount = totalCount;
break; break;
case SUM: case SUM:
if (doubleCountValue > 0) { if (doubleCountValue > 0) {
@ -157,7 +162,8 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity {
} }
} }
public TimescaleTsKvEntity(Long tsBucket, Long interval, Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount) { public TimescaleTsKvEntity(Long tsBucket, Long interval, Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount, Long aggValuesLastTs) {
super(aggValuesLastTs);
if (!isAllNull(tsBucket, interval, booleanValueCount, strValueCount, longValueCount, doubleValueCount, jsonValueCount)) { if (!isAllNull(tsBucket, interval, booleanValueCount, strValueCount, longValueCount, doubleValueCount, jsonValueCount)) {
this.ts = tsBucket + interval / 2; this.ts = tsBucket + interval / 2;
if (booleanValueCount != 0) { if (booleanValueCount != 0) {

13
dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java

@ -16,12 +16,15 @@
package org.thingsboard.server.dao.model.sqlts.ts; package org.thingsboard.server.dao.model.sqlts.ts;
import lombok.Data; import lombok.Data;
import lombok.EqualsAndHashCode;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import javax.persistence.Entity; import javax.persistence.Entity;
import javax.persistence.IdClass; import javax.persistence.IdClass;
import javax.persistence.Table; import javax.persistence.Table;
import javax.persistence.Transient;
@EqualsAndHashCode(callSuper = true)
@Data @Data
@Entity @Entity
@Table(name = "ts_kv") @Table(name = "ts_kv")
@ -31,11 +34,13 @@ public final class TsKvEntity extends AbstractTsKvEntity {
public TsKvEntity() { public TsKvEntity() {
} }
public TsKvEntity(String strValue) { public TsKvEntity(String strValue, Long aggValuesLastTs) {
super(aggValuesLastTs);
this.strValue = strValue; this.strValue = strValue;
} }
public TsKvEntity(Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String aggType) { public TsKvEntity(Long longValue, Double doubleValue, Long longCountValue, Long doubleCountValue, String aggType, Long aggValuesLastTs) {
super(aggValuesLastTs);
if (!isAllNull(longValue, doubleValue, longCountValue, doubleCountValue)) { if (!isAllNull(longValue, doubleValue, longCountValue, doubleCountValue)) {
switch (aggType) { switch (aggType) {
case AVG: case AVG:
@ -52,6 +57,7 @@ public final class TsKvEntity extends AbstractTsKvEntity {
} else { } else {
this.doubleValue = 0.0; this.doubleValue = 0.0;
} }
this.aggValuesCount = totalCount;
break; break;
case SUM: case SUM:
if (doubleCountValue > 0) { if (doubleCountValue > 0) {
@ -74,7 +80,8 @@ public final class TsKvEntity extends AbstractTsKvEntity {
} }
} }
public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount) { public TsKvEntity(Long booleanValueCount, Long strValueCount, Long longValueCount, Long doubleValueCount, Long jsonValueCount, Long aggValuesLastTs) {
super(aggValuesLastTs);
if (!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) { if (!isAllNull(booleanValueCount, strValueCount, longValueCount, doubleValueCount)) {
if (booleanValueCount != 0) { if (booleanValueCount != 0) {
this.longValue = booleanValueCount; this.longValue = booleanValueCount;

165
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java

@ -17,8 +17,6 @@ package org.thingsboard.server.dao.sqlts;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.PageRequest;
@ -26,9 +24,9 @@ import org.springframework.data.domain.Sort;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
@ -47,10 +45,9 @@ import java.util.Comparator;
import java.util.List; import java.util.List;
import java.util.Optional; import java.util.Optional;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.function.Function; import java.util.function.Function;
import java.util.stream.Collectors;
@SuppressWarnings("UnstableApiUsage")
@Slf4j @Slf4j
public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao { public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao {
@ -81,7 +78,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
Comparator.comparing((Function<TsKvEntity, UUID>) AbstractTsKvEntity::getEntityId) Comparator.comparing((Function<TsKvEntity, UUID>) AbstractTsKvEntity::getEntityId)
.thenComparing(AbstractTsKvEntity::getKey) .thenComparing(AbstractTsKvEntity::getKey)
.thenComparing(AbstractTsKvEntity::getTs) .thenComparing(AbstractTsKvEntity::getTs)
); );
} }
@PreDestroy @PreDestroy
@ -114,16 +111,16 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { public ListenableFuture<List<ReadTsKvQueryResult>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return processFindAllAsync(tenantId, entityId, queries); return processFindAllAsync(tenantId, entityId, queries);
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { public ListenableFuture<ReadTsKvQueryResult> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) { if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(entityId, query); return Futures.immediateFuture(findAllAsyncWithLimit(entityId, query));
} else { } else {
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>(); List<ListenableFuture<Optional<TsKvEntity>>> futures = new ArrayList<>();
long endPeriod = query.getEndTs(); long endPeriod = query.getEndTs();
long startPeriod = query.getStartTs(); long startPeriod = query.getStartTs();
long step = query.getInterval(); long step = query.getInterval();
@ -131,15 +128,16 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
long startTs = startPeriod; long startTs = startPeriod;
long endTs = Math.min(startPeriod + step, endPeriod + 1); long endTs = Math.min(startPeriod + step, endPeriod + 1);
long ts = startTs + (endTs - startTs) / 2; long ts = startTs + (endTs - startTs) / 2;
ListenableFuture<Optional<TsKvEntry>> aggregateTsKvEntry = findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()); ListenableFuture<Optional<TsKvEntity>> aggregateTsKvEntry =
service.submit(() -> findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()));
futures.add(aggregateTsKvEntry); futures.add(aggregateTsKvEntry);
startPeriod = endTs; startPeriod = endTs;
} }
return getTskvEntriesFuture(Futures.allAsList(futures)); return getReadTsKvQueryResultFuture(query, Futures.allAsList(futures));
} }
} }
private ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { private ReadTsKvQueryResult findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) {
Integer keyId = getOrSaveKeyId(query.getKey()); Integer keyId = getOrSaveKeyId(query.getKey());
List<TsKvEntity> tsKvEntities = tsKvRepository.findAllWithLimit( List<TsKvEntity> tsKvEntities = tsKvRepository.findAllWithLimit(
entityId.getId(), entityId.getId(),
@ -147,125 +145,50 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq
query.getStartTs(), query.getStartTs(),
query.getEndTs(), query.getEndTs(),
PageRequest.of(0, query.getLimit(), PageRequest.of(0, query.getLimit(),
Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative()))); Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative())));
tsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(query.getKey())); tsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(query.getKey()));
return Futures.immediateFuture(DaoUtil.convertDataList(tsKvEntities)); List<TsKvEntry> tsKvEntries = DaoUtil.convertDataList(tsKvEntities);
} long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs());
return new ReadTsKvQueryResult(query.getKey(), tsKvEntries, lastTs);
ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { }
List<CompletableFuture<TsKvEntity>> entitiesFutures = new ArrayList<>();
switchAggregation(entityId, key, startTs, endTs, aggregation, entitiesFutures); Optional<TsKvEntity> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) {
return Futures.transform(setFutures(entitiesFutures), entity -> { TsKvEntity entity = switchAggregation(entityId, key, startTs, endTs, aggregation);
if (entity != null && entity.isNotEmpty()) { if (entity != null && entity.isNotEmpty()) {
entity.setEntityId(entityId.getId()); entity.setEntityId(entityId.getId());
entity.setStrKey(key); entity.setStrKey(key);
entity.setTs(ts); entity.setTs(ts);
return Optional.of(DaoUtil.getData(entity)); return Optional.of(entity);
} else { } else {
return Optional.empty(); return Optional.empty();
} }
}, MoreExecutors.directExecutor());
} }
protected void switchAggregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List<CompletableFuture<TsKvEntity>> entitiesFutures) { protected TsKvEntity switchAggregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation) {
var keyId = getOrSaveKeyId(key);
switch (aggregation) { switch (aggregation) {
case AVG: case AVG:
findAvg(entityId, key, startTs, endTs, entitiesFutures); return tsKvRepository.findAvg(entityId.getId(), keyId, startTs, endTs);
break;
case MAX: case MAX:
findMax(entityId, key, startTs, endTs, entitiesFutures); var max = tsKvRepository.findNumericMax(entityId.getId(), keyId, startTs, endTs);
break; if (max.isNotEmpty()) {
return max;
} else {
return tsKvRepository.findStringMax(entityId.getId(), keyId, startTs, endTs);
}
case MIN: case MIN:
findMin(entityId, key, startTs, endTs, entitiesFutures); var min = tsKvRepository.findNumericMin(entityId.getId(), keyId, startTs, endTs);
break; if (min.isNotEmpty()) {
return min;
} else {
return tsKvRepository.findStringMin(entityId.getId(), keyId, startTs, endTs);
}
case SUM: case SUM:
findSum(entityId, key, startTs, endTs, entitiesFutures); return tsKvRepository.findSum(entityId.getId(), keyId, startTs, endTs);
break;
case COUNT: case COUNT:
findCount(entityId, key, startTs, endTs, entitiesFutures); return tsKvRepository.findCount(entityId.getId(), keyId, startTs, endTs);
break;
default: default:
throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); throw new IllegalArgumentException("Not supported aggregation type: " + aggregation);
} }
} }
protected void findCount(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findCount(
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findSum(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findSum(
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findMin(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMin(
entityId.getId(),
keyId,
startTs,
endTs));
entitiesFutures.add(tsKvRepository.findNumericMin(
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findMax(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findStringMax(
entityId.getId(),
keyId,
startTs,
endTs));
entitiesFutures.add(tsKvRepository.findNumericMax(
entityId.getId(),
keyId,
startTs,
endTs));
}
protected void findAvg(EntityId entityId, String key, long startTs, long endTs, List<CompletableFuture<TsKvEntity>> entitiesFutures) {
Integer keyId = getOrSaveKeyId(key);
entitiesFutures.add(tsKvRepository.findAvg(
entityId.getId(),
keyId,
startTs,
endTs));
}
protected SettableFuture<TsKvEntity> setFutures(List<CompletableFuture<TsKvEntity>> entitiesFutures) {
SettableFuture<TsKvEntity> listenableFuture = SettableFuture.create();
CompletableFuture<List<TsKvEntity>> entities =
CompletableFuture.allOf(entitiesFutures.toArray(new CompletableFuture[entitiesFutures.size()]))
.thenApply(v -> entitiesFutures.stream()
.map(CompletableFuture::join)
.collect(Collectors.toList()));
entities.whenComplete((tsKvEntities, throwable) -> {
if (throwable != null) {
listenableFuture.setException(throwable);
} else {
TsKvEntity result = null;
for (TsKvEntity entity : tsKvEntities) {
if (entity.isNotEmpty()) {
result = entity;
break;
}
}
listenableFuture.set(result);
}
});
return listenableFuture;
}
} }

13
dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java

@ -24,6 +24,7 @@ import org.springframework.beans.factory.annotation.Value;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
@ -38,6 +39,7 @@ import java.util.Objects;
import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeUnit;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@SuppressWarnings("UnstableApiUsage")
@Slf4j @Slf4j
public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseriesDao implements AggregationTimeseriesDao { public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseriesDao implements AggregationTimeseriesDao {
@ -86,22 +88,19 @@ public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseries
} }
} }
protected ListenableFuture<List<TsKvEntry>> processFindAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { protected ListenableFuture<List<ReadTsKvQueryResult>> processFindAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
List<ListenableFuture<List<TsKvEntry>>> futures = queries List<ListenableFuture<ReadTsKvQueryResult>> futures = queries
.stream() .stream()
.map(query -> findAllAsync(tenantId, entityId, query)) .map(query -> findAllAsync(tenantId, entityId, query))
.collect(Collectors.toList()); .collect(Collectors.toList());
return Futures.transform(Futures.allAsList(futures), new Function<>() { return Futures.transform(Futures.allAsList(futures), new Function<>() {
@Nullable @Nullable
@Override @Override
public List<TsKvEntry> apply(@Nullable List<List<TsKvEntry>> results) { public List<ReadTsKvQueryResult> apply(@Nullable List<ReadTsKvQueryResult> results) {
if (results == null || results.isEmpty()) { if (results == null || results.isEmpty()) {
return null; return null;
} }
return results.stream() return results.stream().filter(Objects::nonNull).collect(Collectors.toList());
.filter(Objects::nonNull)
.flatMap(List::stream)
.collect(Collectors.toList());
} }
}, service); }, service);
} }

3
dao/src/main/java/org/thingsboard/server/dao/sqlts/AggregationTimeseriesDao.java

@ -19,11 +19,12 @@ import com.google.common.util.concurrent.ListenableFuture;
import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.List; import java.util.List;
public interface AggregationTimeseriesDao { public interface AggregationTimeseriesDao {
ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); ListenableFuture<ReadTsKvQueryResult> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query);
} }

22
dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java

@ -22,14 +22,20 @@ import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException; import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.dao.DataIntegrityViolationException; import org.springframework.dao.DataIntegrityViolationException;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary; import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionary;
import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey; import org.thingsboard.server.dao.model.sqlts.dictionary.TsKvDictionaryCompositeKey;
import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity;
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService;
import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository; import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository;
import javax.annotation.Nullable; import javax.annotation.Nullable;
import java.util.List; import java.util.List;
import java.util.Objects;
import java.util.Optional; import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentMap;
@ -80,18 +86,20 @@ public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeni
return keyId; return keyId;
} }
protected ListenableFuture<List<TsKvEntry>> getTskvEntriesFuture(ListenableFuture<List<Optional<TsKvEntry>>> future) { protected ListenableFuture<ReadTsKvQueryResult> getReadTsKvQueryResultFuture(ReadTsKvQuery query, ListenableFuture<List<Optional<? extends AbstractTsKvEntity>>> future) {
return Futures.transform(future, new Function<List<Optional<TsKvEntry>>, List<TsKvEntry>>() { return Futures.transform(future, new Function<>() {
@Nullable @Nullable
@Override @Override
public List<TsKvEntry> apply(@Nullable List<Optional<TsKvEntry>> results) { public ReadTsKvQueryResult apply(@Nullable List<Optional<? extends AbstractTsKvEntity>> results) {
if (results == null || results.isEmpty()) { if (results == null || results.isEmpty()) {
return null; return null;
} }
return results.stream() List<? extends AbstractTsKvEntity> data = results.stream().filter(Optional::isPresent).map(Optional::get).collect(Collectors.toList());
.filter(Optional::isPresent) var lastTs = data.stream().map(AbstractTsKvEntity::getAggValuesLastTs).filter(Objects::nonNull).max(Long::compare);
.map(Optional::get) if (lastTs.isEmpty()) {
.collect(Collectors.toList()); lastTs = data.stream().map(AbstractTsKvEntity::getTs).filter(Objects::nonNull).max(Long::compare);
}
return new ReadTsKvQueryResult(query.getKey(), DaoUtil.convertDataList(data), lastTs.orElse(query.getStartTs()));
} }
}, service); }, service);
} }

4
dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java

@ -31,6 +31,7 @@ import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult;
@ -190,7 +191,8 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme
long endTs = query.getStartTs() - 1; long endTs = query.getStartTs() - 1;
ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1,
Aggregation.NONE, DESC_ORDER); Aggregation.NONE, DESC_ORDER);
return aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); return Futures.transform(aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery),
ReadTsKvQueryResult::getData, MoreExecutors.directExecutor());
} }
protected ListenableFuture<TsKvEntry> getFindLatestFuture(EntityId entityId, String key) { protected ListenableFuture<TsKvEntry> getFindLatestFuture(EntityId entityId, String key) {

98
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java

@ -36,54 +36,79 @@ public class AggregationRepository {
public static final String FIND_SUM = "findSum"; public static final String FIND_SUM = "findSum";
public static final String FIND_COUNT = "findCount"; public static final String FIND_COUNT = "findCount";
public static final String FROM_WHERE_CLAUSE = "FROM ts_kv tskv WHERE tskv.entity_id = cast(:entityId AS uuid) AND tskv.key= cast(:entityKey AS int) AND tskv.ts > :startTs AND tskv.ts <= :endTs GROUP BY tskv.entity_id, tskv.key, tsBucket ORDER BY tskv.entity_id, tskv.key, tsBucket"; public static final String FROM_WHERE_CLAUSE = "FROM ts_kv tskv WHERE " +
"tskv.entity_id = cast(:entityId AS uuid) " +
public static final String FIND_AVG_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, 'AVG' AS aggType "; "AND tskv.key= cast(:entityKey AS int) " +
"AND tskv.ts > :startTs AND tskv.ts <= :endTs " +
public static final String FIND_MAX_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, MAX(COALESCE(tskv.long_v, -9223372036854775807)) AS longValue, MAX(COALESCE(tskv.dbl_v, -1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MAX(tskv.str_v) AS strValue, 'MAX' AS aggType "; "GROUP BY tskv.entity_id, tskv.key, tsBucket " +
"ORDER BY tskv.entity_id, tskv.key, tsBucket";
public static final String FIND_MIN_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, MIN(COALESCE(tskv.long_v, 9223372036854775807)) AS longValue, MIN(COALESCE(tskv.dbl_v, 1.79769E+308)) as doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, MIN(tskv.str_v) AS strValue, 'MIN' AS aggType ";
public static final String FIND_AVG_QUERY = "SELECT " +
public static final String FIND_SUM_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, null AS strValue, null AS jsonValue, 'SUM' AS aggType "; "time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, " +
"SUM(COALESCE(tskv.long_v, 0)) AS longValue, " +
public static final String FIND_COUNT_QUERY = "SELECT time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, SUM(CASE WHEN tskv.bool_v IS NULL THEN 0 ELSE 1 END) AS booleanValueCount, SUM(CASE WHEN tskv.str_v IS NULL THEN 0 ELSE 1 END) AS strValueCount, SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longValueCount, SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleValueCount, SUM(CASE WHEN tskv.json_v IS NULL THEN 0 ELSE 1 END) AS jsonValueCount "; "SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, " +
"SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, " +
"SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, " +
"null AS strValue, 'AVG' AS aggType, MAX(tskv.ts) AS maxAggTs ";
public static final String FIND_MAX_QUERY = "SELECT " +
"time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, " +
"MAX(COALESCE(tskv.long_v, -9223372036854775807)) AS longValue, " +
"MAX(COALESCE(tskv.dbl_v, -1.79769E+308)) as doubleValue, " +
"SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, " +
"SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, " +
"MAX(tskv.str_v) AS strValue, 'MAX' AS aggType, MAX(tskv.ts) AS maxAggTs ";
public static final String FIND_MIN_QUERY = "SELECT " +
"time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, " +
"MIN(COALESCE(tskv.long_v, 9223372036854775807)) AS longValue, " +
"MIN(COALESCE(tskv.dbl_v, 1.79769E+308)) as doubleValue, " +
"SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, " +
"SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, " +
"MIN(tskv.str_v) AS strValue, 'MIN' AS aggType, MAX(tskv.ts) AS maxAggTs ";
public static final String FIND_SUM_QUERY = "SELECT " +
"time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, " +
"SUM(COALESCE(tskv.long_v, 0)) AS longValue, SUM(COALESCE(tskv.dbl_v, 0.0)) AS doubleValue, " +
"SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longCountValue, " +
"SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleCountValue, " +
"null AS strValue, null AS jsonValue, 'SUM' AS aggType, MAX(tskv.ts) AS maxAggTs ";
public static final String FIND_COUNT_QUERY = "SELECT " +
"time_bucket(:timeBucket, tskv.ts, :startTs) AS tsBucket, :timeBucket AS interval, " +
"SUM(CASE WHEN tskv.bool_v IS NULL THEN 0 ELSE 1 END) AS booleanValueCount, " +
"SUM(CASE WHEN tskv.str_v IS NULL THEN 0 ELSE 1 END) AS strValueCount, " +
"SUM(CASE WHEN tskv.long_v IS NULL THEN 0 ELSE 1 END) AS longValueCount, " +
"SUM(CASE WHEN tskv.dbl_v IS NULL THEN 0 ELSE 1 END) AS doubleValueCount, " +
"SUM(CASE WHEN tskv.json_v IS NULL THEN 0 ELSE 1 END) AS jsonValueCount, " +
"MAX(tskv.ts) AS maxAggTs ";
@PersistenceContext @PersistenceContext
private EntityManager entityManager; private EntityManager entityManager;
@Async @SuppressWarnings("unchecked")
public CompletableFuture<List<TimescaleTsKvEntity>> findAvg(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { public List<TimescaleTsKvEntity> findAvg(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) {
@SuppressWarnings("unchecked") return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG);
List<TimescaleTsKvEntity> resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG);
return CompletableFuture.supplyAsync(() -> resultList);
} }
@Async @SuppressWarnings("unchecked")
public CompletableFuture<List<TimescaleTsKvEntity>> findMax(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { public List<TimescaleTsKvEntity> findMax(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) {
@SuppressWarnings("unchecked") return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX);
List<TimescaleTsKvEntity> resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX);
return CompletableFuture.supplyAsync(() -> resultList);
} }
@Async @SuppressWarnings("unchecked")
public CompletableFuture<List<TimescaleTsKvEntity>> findMin(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { public List<TimescaleTsKvEntity> findMin(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) {
@SuppressWarnings("unchecked") return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN);
List<TimescaleTsKvEntity> resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN);
return CompletableFuture.supplyAsync(() -> resultList);
} }
@Async @SuppressWarnings("unchecked")
public CompletableFuture<List<TimescaleTsKvEntity>> findSum(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { public List<TimescaleTsKvEntity> findSum(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) {
@SuppressWarnings("unchecked") return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM);
List<TimescaleTsKvEntity> resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM);
return CompletableFuture.supplyAsync(() -> resultList);
} }
@Async @SuppressWarnings("unchecked")
public CompletableFuture<List<TimescaleTsKvEntity>> findCount(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { public List<TimescaleTsKvEntity> findCount(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) {
@SuppressWarnings("unchecked") return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_COUNT);
List<TimescaleTsKvEntity> resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_COUNT);
return CompletableFuture.supplyAsync(() -> resultList);
} }
private List getResultList(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs, String query) { private List getResultList(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs, String query) {
@ -96,5 +121,4 @@ public class AggregationRepository {
.getResultList(); .getResultList();
} }
} }

134
dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java

@ -30,11 +30,13 @@ import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity; import org.thingsboard.server.dao.model.sqlts.timescale.ts.TimescaleTsKvEntity;
import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao; import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao;
@ -101,13 +103,13 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { public ListenableFuture<List<ReadTsKvQueryResult>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return processFindAllAsync(tenantId, entityId, queries); return processFindAllAsync(tenantId, entityId, queries);
} }
@Override @Override
public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl) { public ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl) {
int dataPointDays = getDataPointDays(tsKvEntry, computeTtl(ttl)); int dataPointDays = getDataPointDays(tsKvEntry, computeTtl(ttl));
String strKey = tsKvEntry.getKey(); String strKey = tsKvEntry.getKey();
Integer keyId = getOrSaveKeyId(strKey); Integer keyId = getOrSaveKeyId(strKey);
TimescaleTsKvEntity entity = new TimescaleTsKvEntity(); TimescaleTsKvEntity entity = new TimescaleTsKvEntity();
@ -148,15 +150,15 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { public ListenableFuture<ReadTsKvQueryResult> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) { if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(entityId, query); return Futures.immediateFuture(findAllAsyncWithLimit(entityId, query));
} else { } else {
long startTs = query.getStartTs(); long startTs = query.getStartTs();
long endTs = query.getEndTs(); long endTs = query.getEndTs();
long timeBucket = query.getInterval(); long timeBucket = query.getInterval();
ListenableFuture<List<Optional<TsKvEntry>>> future = findAllAndAggregateAsync(entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation()); List<Optional<? extends AbstractTsKvEntity>> data = findAllAndAggregateAsync(entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation());
return getTskvEntriesFuture(future); return getReadTsKvQueryResultFuture(query, Futures.immediateFuture(data));
} }
} }
@ -165,7 +167,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
super.cleanup(systemTtl); super.cleanup(systemTtl);
} }
private ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { private ReadTsKvQueryResult findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) {
String strKey = query.getKey(); String strKey = query.getKey();
Integer keyId = getOrSaveKeyId(strKey); Integer keyId = getOrSaveKeyId(strKey);
List<TimescaleTsKvEntity> timescaleTsKvEntities = tsKvRepository.findAllWithLimit( List<TimescaleTsKvEntity> timescaleTsKvEntities = tsKvRepository.findAllWithLimit(
@ -174,105 +176,49 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements
query.getStartTs(), query.getStartTs(),
query.getEndTs(), query.getEndTs(),
PageRequest.of(0, query.getLimit(), PageRequest.of(0, query.getLimit(),
Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative())));; Sort.by(new Sort.Order(Sort.Direction.fromString(query.getOrder()), "ts").nullsNative())));
;
timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey)); timescaleTsKvEntities.forEach(tsKvEntity -> tsKvEntity.setStrKey(strKey));
return Futures.immediateFuture(DaoUtil.convertDataList(timescaleTsKvEntities)); var tsKvEntries = DaoUtil.convertDataList(timescaleTsKvEntities);
} long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs());
return new ReadTsKvQueryResult(query.getKey(), tsKvEntries, lastTs);
private ListenableFuture<List<Optional<TsKvEntry>>> findAllAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) { }
CompletableFuture<List<TimescaleTsKvEntity>> listCompletableFuture = switchAggregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId());
SettableFuture<List<TimescaleTsKvEntity>> listenableFuture = SettableFuture.create(); private List<Optional<? extends AbstractTsKvEntity>> findAllAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) {
listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> { List<TimescaleTsKvEntity> timescaleTsKvEntities = switchAggregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId());
if (throwable != null) { if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) {
listenableFuture.setException(throwable); List<Optional<? extends AbstractTsKvEntity>> result = new ArrayList<>();
} else { timescaleTsKvEntities.forEach(entity -> {
listenableFuture.set(timescaleTsKvEntities); if (entity != null && entity.isNotEmpty()) {
} entity.setEntityId(entityId.getId());
}); entity.setStrKey(key);
return Futures.transform(listenableFuture, timescaleTsKvEntities -> { result.add(Optional.of(entity));
if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) { } else {
List<Optional<TsKvEntry>> result = new ArrayList<>(); result.add(Optional.empty());
timescaleTsKvEntities.forEach(entity -> { }
if (entity != null && entity.isNotEmpty()) { });
entity.setEntityId(entityId.getId()); return result;
entity.setStrKey(key); } else {
result.add(Optional.of(DaoUtil.getData(entity))); return Collections.emptyList();
} else { }
result.add(Optional.empty());
}
});
return result;
} else {
return Collections.emptyList();
}
}, MoreExecutors.directExecutor());
} }
private CompletableFuture<List<TimescaleTsKvEntity>> switchAggregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId) { private List<TimescaleTsKvEntity> switchAggregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
switch (aggregation) { switch (aggregation) {
case AVG: case AVG:
return findAvg(key, startTs, endTs, timeBucket, entityId); return aggregationRepository.findAvg(entityId, keyId, timeBucket, startTs, endTs);
case MAX: case MAX:
return findMax(key, startTs, endTs, timeBucket, entityId); return aggregationRepository.findMax(entityId, keyId, timeBucket, startTs, endTs);
case MIN: case MIN:
return findMin(key, startTs, endTs, timeBucket, entityId); return aggregationRepository.findMin(entityId, keyId, timeBucket, startTs, endTs);
case SUM: case SUM:
return findSum(key, startTs, endTs, timeBucket, entityId); return aggregationRepository.findSum(entityId, keyId, timeBucket, startTs, endTs);
case COUNT: case COUNT:
return findCount(key, startTs, endTs, timeBucket, entityId); return aggregationRepository.findCount(entityId, keyId, timeBucket, startTs, endTs);
default: default:
throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); throw new IllegalArgumentException("Not supported aggregation type: " + aggregation);
} }
} }
private CompletableFuture<List<TimescaleTsKvEntity>> findCount(String key, long startTs, long endTs, long timeBucket, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findCount(
entityId,
keyId,
timeBucket,
startTs,
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findSum(String key, long startTs, long endTs, long timeBucket, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findSum(
entityId,
keyId,
timeBucket,
startTs,
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findMin(String key, long startTs, long endTs, long timeBucket, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findMin(
entityId,
keyId,
timeBucket,
startTs,
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findMax(String key, long startTs, long endTs, long timeBucket, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findMax(
entityId,
keyId,
timeBucket,
startTs,
endTs);
}
private CompletableFuture<List<TimescaleTsKvEntity>> findAvg(String key, long startTs, long endTs, long timeBucket, UUID entityId) {
Integer keyId = getOrSaveKeyId(key);
return aggregationRepository.findAvg(
entityId,
keyId,
timeBucket,
startTs,
endTs);
}
} }

37
dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java

@ -20,14 +20,12 @@ import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Modifying; import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query; import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param; import org.springframework.data.repository.query.Param;
import org.springframework.scheduling.annotation.Async;
import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.dao.model.sqlts.ts.TsKvCompositeKey; import org.thingsboard.server.dao.model.sqlts.ts.TsKvCompositeKey;
import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity;
import java.util.List; import java.util.List;
import java.util.UUID; import java.util.UUID;
import java.util.concurrent.CompletableFuture;
public interface TsKvRepository extends JpaRepository<TsKvEntity, TsKvCompositeKey> { public interface TsKvRepository extends JpaRepository<TsKvEntity, TsKvCompositeKey> {
@ -48,82 +46,75 @@ public interface TsKvRepository extends JpaRepository<TsKvEntity, TsKvCompositeK
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async @Query("SELECT new TsKvEntity(MAX(tskv.strValue), MAX(tskv.ts)) FROM TsKvEntity tskv " +
@Query("SELECT new TsKvEntity(MAX(tskv.strValue)) FROM TsKvEntity tskv " +
"WHERE tskv.strValue IS NOT NULL " + "WHERE tskv.strValue IS NOT NULL " +
"AND tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "AND tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findStringMax(@Param("entityId") UUID entityId, TsKvEntity findStringMax(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(MAX(COALESCE(tskv.longValue, -9223372036854775807)), " + @Query("SELECT new TsKvEntity(MAX(COALESCE(tskv.longValue, -9223372036854775807)), " +
"MAX(COALESCE(tskv.doubleValue, -1.79769E+308)), " + "MAX(COALESCE(tskv.doubleValue, -1.79769E+308)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'MAX') FROM TsKvEntity tskv " + "'MAX', MAX(tskv.ts)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findNumericMax(@Param("entityId") UUID entityId, TsKvEntity findNumericMax(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async @Query("SELECT new TsKvEntity(MIN(tskv.strValue), MAX(tskv.ts)) FROM TsKvEntity tskv " +
@Query("SELECT new TsKvEntity(MIN(tskv.strValue)) FROM TsKvEntity tskv " +
"WHERE tskv.strValue IS NOT NULL " + "WHERE tskv.strValue IS NOT NULL " +
"AND tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "AND tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findStringMin(@Param("entityId") UUID entityId, TsKvEntity findStringMin(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(MIN(COALESCE(tskv.longValue, 9223372036854775807)), " + @Query("SELECT new TsKvEntity(MIN(COALESCE(tskv.longValue, 9223372036854775807)), " +
"MIN(COALESCE(tskv.doubleValue, 1.79769E+308)), " + "MIN(COALESCE(tskv.doubleValue, 1.79769E+308)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'MIN') FROM TsKvEntity tskv " + "'MIN', MAX(tskv.ts)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findNumericMin( TsKvEntity findNumericMin(
@Param("entityId") UUID entityId, @Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(SUM(CASE WHEN tskv.booleanValue IS NULL THEN 0 ELSE 1 END), " + @Query("SELECT new TsKvEntity(SUM(CASE WHEN tskv.booleanValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.strValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.strValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.jsonValue IS NULL THEN 0 ELSE 1 END)) FROM TsKvEntity tskv " + "SUM(CASE WHEN tskv.jsonValue IS NULL THEN 0 ELSE 1 END), MAX(tskv.ts)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findCount(@Param("entityId") UUID entityId, TsKvEntity findCount(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " + @Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " +
"SUM(COALESCE(tskv.doubleValue, 0.0)), " + "SUM(COALESCE(tskv.doubleValue, 0.0)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'AVG') FROM TsKvEntity tskv " + "'AVG', MAX(tskv.ts)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findAvg(@Param("entityId") UUID entityId, TsKvEntity findAvg(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);
@Async
@Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " + @Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " +
"SUM(COALESCE(tskv.doubleValue, 0.0)), " + "SUM(COALESCE(tskv.doubleValue, 0.0)), " +
"SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.longValue IS NULL THEN 0 ELSE 1 END), " +
"SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " + "SUM(CASE WHEN tskv.doubleValue IS NULL THEN 0 ELSE 1 END), " +
"'SUM') FROM TsKvEntity tskv " + "'SUM', MAX(tskv.ts)) FROM TsKvEntity tskv " +
"WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") "WHERE tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs")
CompletableFuture<TsKvEntity> findSum(@Param("entityId") UUID entityId, TsKvEntity findSum(@Param("entityId") UUID entityId,
@Param("entityKey") int entityKey, @Param("entityKey") int entityKey,
@Param("startTs") long startTs, @Param("startTs") long startTs,
@Param("endTs") long endTs); @Param("endTs") long endTs);

70
dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java

@ -19,6 +19,7 @@ import com.datastax.oss.driver.api.core.cql.Row;
import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.thingsboard.server.common.data.kv.AggTsKvEntry;
import org.thingsboard.server.common.data.kv.Aggregation; import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry;
import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry;
@ -28,6 +29,7 @@ import org.thingsboard.server.common.data.kv.JsonDataEntry;
import org.thingsboard.server.common.data.kv.LongDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry;
import org.thingsboard.server.common.data.kv.StringDataEntry; import org.thingsboard.server.common.data.kv.StringDataEntry;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper;
import org.thingsboard.server.dao.nosql.TbResultSet; import org.thingsboard.server.dao.nosql.TbResultSet;
import javax.annotation.Nullable; import javax.annotation.Nullable;
@ -40,7 +42,7 @@ import java.util.stream.Collectors;
* Created by ashvayka on 20.02.17. * Created by ashvayka on 20.02.17.
*/ */
@Slf4j @Slf4j
public class AggregatePartitionsFunction implements com.google.common.util.concurrent.AsyncFunction<List<TbResultSet>, Optional<TsKvEntry>> { public class AggregatePartitionsFunction implements com.google.common.util.concurrent.AsyncFunction<List<TbResultSet>, Optional<TsKvEntryAggWrapper>> {
private static final int LONG_CNT_POS = 0; private static final int LONG_CNT_POS = 0;
private static final int DOUBLE_CNT_POS = 1; private static final int DOUBLE_CNT_POS = 1;
@ -52,6 +54,7 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
private static final int BOOL_POS = 7; private static final int BOOL_POS = 7;
private static final int STR_POS = 8; private static final int STR_POS = 8;
private static final int JSON_POS = 9; private static final int JSON_POS = 9;
private static final int MAX_TS_POS = 10;
private final Aggregation aggregation; private final Aggregation aggregation;
private final String key; private final String key;
@ -66,29 +69,29 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
} }
@Override @Override
public ListenableFuture<Optional<TsKvEntry>> apply(@Nullable List<TbResultSet> rsList) { public ListenableFuture<Optional<TsKvEntryAggWrapper>> apply(@Nullable List<TbResultSet> rsList) {
log.trace("[{}][{}][{}] Going to aggregate data", key, ts, aggregation); log.trace("[{}][{}][{}] Going to aggregate data", key, ts, aggregation);
if (rsList == null || rsList.isEmpty()) { if (rsList == null || rsList.isEmpty()) {
return Futures.immediateFuture(Optional.empty()); return Futures.immediateFuture(Optional.empty());
} }
return Futures.transform( return Futures.transform(
Futures.allAsList( Futures.allAsList(
rsList.stream().map(rs -> rs.allRows(this.executor)) rsList.stream().map(rs -> rs.allRows(this.executor))
.collect(Collectors.toList())), .collect(Collectors.toList())),
rowsList -> { rowsList -> {
try { try {
AggregationResult aggResult = new AggregationResult(); AggregationResult aggResult = new AggregationResult();
for (List<Row> rs : rowsList) { for (List<Row> rs : rowsList) {
for (Row row : rs) { for (Row row : rs) {
processResultSetRow(row, aggResult); processResultSetRow(row, aggResult);
}
}
return processAggregationResult(aggResult);
} catch (Exception e) {
log.error("[{}][{}][{}] Failed to aggregate data", key, ts, aggregation, e);
return Optional.empty();
} }
} }, this.executor);
return processAggregationResult(aggResult);
} catch (Exception e) {
log.error("[{}][{}][{}] Failed to aggregate data", key, ts, aggregation, e);
return Optional.empty();
}
}, this.executor);
} }
private void processResultSetRow(Row row, AggregationResult aggResult) { private void processResultSetRow(Row row, AggregationResult aggResult) {
@ -105,6 +108,7 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
long boolCount = row.getLong(BOOL_CNT_POS); long boolCount = row.getLong(BOOL_CNT_POS);
long strCount = row.getLong(STR_CNT_POS); long strCount = row.getLong(STR_CNT_POS);
long jsonCount = row.getLong(JSON_CNT_POS); long jsonCount = row.getLong(JSON_CNT_POS);
long aggValuesLastTs = row.getLong(MAX_TS_POS);
if (longCount > 0 || doubleCount > 0) { if (longCount > 0 || doubleCount > 0) {
if (longCount > 0) { if (longCount > 0) {
@ -134,6 +138,8 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
return; return;
} }
aggResult.aggValuesLastTs = Math.max(aggResult.aggValuesLastTs, aggValuesLastTs);
if (aggregation == Aggregation.COUNT) { if (aggregation == Aggregation.COUNT) {
aggResult.count += curCount; aggResult.count += curCount;
} else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) {
@ -231,34 +237,37 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
} }
} }
private Optional<TsKvEntry> processAggregationResult(AggregationResult aggResult) { private Optional<TsKvEntryAggWrapper> processAggregationResult(AggregationResult aggResult) {
Optional<TsKvEntry> result; Optional<TsKvEntry> result;
if (aggResult.dataType == null) { if (aggResult.dataType == null) {
result = Optional.empty(); result = Optional.empty();
} else if (aggregation == Aggregation.COUNT) { } else if (aggregation == Aggregation.COUNT) {
result = Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggResult.count))); result = Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggResult.count)));
} else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) {
result = processAvgOrSumResult(aggResult); result = processAvgOrSumResult(aggregation, aggResult);
} else if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { } else if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) {
result = processMinOrMaxResult(aggResult); result = processMinOrMaxResult(aggResult);
} else { } else {
result = Optional.empty(); result = Optional.empty();
} }
if (!result.isPresent()) { if (result.isEmpty()) {
log.trace("[{}][{}][{}] Aggregated data is empty.", key, ts, aggregation); log.trace("[{}][{}][{}] Aggregated data is empty.", key, ts, aggregation);
} }
return result; return result.map(tsKvEntry -> new TsKvEntryAggWrapper(tsKvEntry, aggResult.aggValuesLastTs));
} }
private Optional<TsKvEntry> processAvgOrSumResult(AggregationResult aggResult) { private Optional<TsKvEntry> processAvgOrSumResult(Aggregation aggregation, AggregationResult aggResult) {
if (aggResult.count == 0 || (aggResult.dataType == DataType.DOUBLE && aggResult.dValue == null) || (aggResult.dataType == DataType.LONG && aggResult.lValue == null)) { if (aggResult.count == 0 || (aggResult.dataType == DataType.DOUBLE && aggResult.dValue == null) || (aggResult.dataType == DataType.LONG && aggResult.lValue == null)) {
return Optional.empty(); return Optional.empty();
} else if (aggResult.dataType == DataType.DOUBLE || aggResult.dataType == DataType.LONG) { } else if (aggResult.dataType == DataType.DOUBLE || aggResult.dataType == DataType.LONG) {
if (aggregation == Aggregation.AVG || aggResult.hasDouble) { if (aggregation == Aggregation.AVG || aggResult.hasDouble) {
double sum = Optional.ofNullable(aggResult.dValue).orElse(0.0d) + Optional.ofNullable(aggResult.lValue).orElse(0L); double sum = Optional.ofNullable(aggResult.dValue).orElse(0.0d) + Optional.ofNullable(aggResult.lValue).orElse(0L);
return Optional.of(new BasicTsKvEntry(ts, new DoubleDataEntry(key, aggregation == Aggregation.SUM ? sum : (sum / aggResult.count)))); DoubleDataEntry doubleDataEntry = new DoubleDataEntry(key, aggregation == Aggregation.SUM ? sum : (sum / aggResult.count));
TsKvEntry result = aggregation == Aggregation.AVG ? new AggTsKvEntry(ts, doubleDataEntry, aggResult.count) : new BasicTsKvEntry(ts, doubleDataEntry);
return Optional.of(result);
} else { } else {
return Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggregation == Aggregation.SUM ? aggResult.lValue : (aggResult.lValue / aggResult.count)))); LongDataEntry longDataEntry = new LongDataEntry(key, aggregation == Aggregation.SUM ? aggResult.lValue : (aggResult.lValue / aggResult.count));
return Optional.of(new BasicTsKvEntry(ts, longDataEntry));
} }
} }
return Optional.empty(); return Optional.empty();
@ -291,5 +300,6 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu
Long lValue = null; Long lValue = null;
long count = 0; long count = 0;
boolean hasDouble = false; boolean hasDouble = false;
long aggValuesLastTs = 0;
} }
} }

19
dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java

@ -36,6 +36,7 @@ import org.thingsboard.server.common.data.kv.BaseDeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult;
import org.thingsboard.server.dao.entityview.EntityViewService; import org.thingsboard.server.dao.entityview.EntityViewService;
@ -52,6 +53,7 @@ import static org.thingsboard.server.common.data.StringUtils.isBlank;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
@SuppressWarnings("UnstableApiUsage")
@Service @Service
@Slf4j @Slf4j
public class BaseTimeseriesService implements TimeseriesService { public class BaseTimeseriesService implements TimeseriesService {
@ -59,7 +61,7 @@ public class BaseTimeseriesService implements TimeseriesService {
private static final int INSERTS_PER_ENTRY = 3; private static final int INSERTS_PER_ENTRY = 3;
private static final int INSERTS_PER_ENTRY_WITHOUT_LATEST = 2; private static final int INSERTS_PER_ENTRY_WITHOUT_LATEST = 2;
private static final int DELETES_PER_ENTRY = INSERTS_PER_ENTRY; private static final int DELETES_PER_ENTRY = INSERTS_PER_ENTRY;
public static final Function<List<Integer>, Integer> SUM_ALL_INTEGERS = new Function<List<Integer>, Integer>() { public static final Function<List<Integer>, Integer> SUM_ALL_INTEGERS = new Function<>() {
@Override @Override
public @Nullable Integer apply(@Nullable List<Integer> input) { public @Nullable Integer apply(@Nullable List<Integer> input) {
int result = 0; int result = 0;
@ -87,7 +89,7 @@ public class BaseTimeseriesService implements TimeseriesService {
private EntityViewService entityViewService; private EntityViewService entityViewService;
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAll(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { public ListenableFuture<List<ReadTsKvQueryResult>> findAllByQueries(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
validate(entityId); validate(entityId);
queries.forEach(this::validate); queries.forEach(this::validate);
if (entityId.getEntityType().equals(EntityType.ENTITY_VIEW)) { if (entityId.getEntityType().equals(EntityType.ENTITY_VIEW)) {
@ -103,6 +105,17 @@ public class BaseTimeseriesService implements TimeseriesService {
return timeseriesDao.findAllAsync(tenantId, entityId, queries); return timeseriesDao.findAllAsync(tenantId, entityId, queries);
} }
@Override
public ListenableFuture<List<TsKvEntry>> findAll(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
return Futures.transform(findAllByQueries(tenantId, entityId, queries),
result -> {
if (result != null && !result.isEmpty()) {
return result.stream().map(ReadTsKvQueryResult::getData).flatMap(Collection::stream).collect(Collectors.toList());
}
return Collections.emptyList();
}, MoreExecutors.directExecutor());
}
@Override @Override
public ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) { public ListenableFuture<List<TsKvEntry>> findLatest(TenantId tenantId, EntityId entityId, Collection<String> keys) {
validate(entityId); validate(entityId);
@ -244,7 +257,7 @@ public class BaseTimeseriesService implements TimeseriesService {
public ListenableFuture<Collection<String>> removeAllLatest(TenantId tenantId, EntityId entityId) { public ListenableFuture<Collection<String>> removeAllLatest(TenantId tenantId, EntityId entityId) {
validate(entityId); validate(entityId);
return Futures.transformAsync(this.findAllLatest(tenantId, entityId), latest -> { return Futures.transformAsync(this.findAllLatest(tenantId, entityId), latest -> {
if (!latest.isEmpty()) { if (latest != null && !latest.isEmpty()) {
Collection<String> keys = latest.stream().map(TsKvEntry::getKey).collect(Collectors.toList()); Collection<String> keys = latest.stream().map(TsKvEntry::getKey).collect(Collectors.toList());
return Futures.transform(this.removeLatest(tenantId, entityId, keys), res -> keys, MoreExecutors.directExecutor()); return Futures.transform(this.removeLatest(tenantId, entityId, keys), res -> keys, MoreExecutors.directExecutor());
} else { } else {

59
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java

@ -42,7 +42,9 @@ import org.thingsboard.server.common.data.kv.DataType;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.KvEntry; import org.thingsboard.server.common.data.kv.KvEntry;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.nosql.TbResultSet; import org.thingsboard.server.dao.nosql.TbResultSet;
import org.thingsboard.server.dao.nosql.TbResultSetFuture; import org.thingsboard.server.dao.nosql.TbResultSetFuture;
@ -71,6 +73,7 @@ import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
@SuppressWarnings("UnstableApiUsage")
@Component @Component
@Slf4j @Slf4j
@NoSqlTsDao @NoSqlTsDao
@ -139,20 +142,10 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) { public ListenableFuture<List<ReadTsKvQueryResult>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries) {
List<ListenableFuture<List<TsKvEntry>>> futures = queries.stream().map(query -> findAllAsync(tenantId, entityId, query)).collect(Collectors.toList()); List<ListenableFuture<ReadTsKvQueryResult>> futures = queries.stream()
return Futures.transform(Futures.allAsList(futures), new Function<>() { .map(query -> findAllAsync(tenantId, entityId, query)).collect(Collectors.toList());
@Nullable return Futures.allAsList(futures);
@Override
public List<TsKvEntry> apply(@Nullable List<List<TsKvEntry>> results) {
if (results == null || results.isEmpty()) {
return null;
}
return results.stream()
.flatMap(List::stream)
.collect(Collectors.toList());
}
}, readResultsProcessingExecutor);
} }
@Override @Override
@ -270,14 +263,14 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
@Override @Override
public ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { public ListenableFuture<ReadTsKvQueryResult> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
if (query.getAggregation() == Aggregation.NONE) { if (query.getAggregation() == Aggregation.NONE) {
return findAllAsyncWithLimit(tenantId, entityId, query); return findAllAsyncWithLimit(tenantId, entityId, query);
} else { } else {
long startPeriod = query.getStartTs(); long startPeriod = query.getStartTs();
long endPeriod = query.getEndTs(); long endPeriod = query.getEndTs();
long step = Math.max(query.getInterval(), MIN_AGGREGATION_STEP_MS); long step = Math.max(query.getInterval(), MIN_AGGREGATION_STEP_MS);
List<ListenableFuture<Optional<TsKvEntry>>> futures = new ArrayList<>(); List<ListenableFuture<Optional<TsKvEntryAggWrapper>>> futures = new ArrayList<>();
while (startPeriod <= endPeriod) { while (startPeriod <= endPeriod) {
long startTs = startPeriod; long startTs = startPeriod;
long endTs = Math.min(startPeriod + step, endPeriod + 1); long endTs = Math.min(startPeriod + step, endPeriod + 1);
@ -286,12 +279,26 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
futures.add(findAndAggregateAsync(tenantId, entityId, subQuery, toPartitionTs(startTs), toPartitionTs(endTs))); futures.add(findAndAggregateAsync(tenantId, entityId, subQuery, toPartitionTs(startTs), toPartitionTs(endTs)));
startPeriod = endTs; startPeriod = endTs;
} }
ListenableFuture<List<Optional<TsKvEntry>>> future = Futures.allAsList(futures); ListenableFuture<List<Optional<TsKvEntryAggWrapper>>> future = Futures.allAsList(futures);
return Futures.transform(future, new Function<>() { return Futures.transform(future, new Function<>() {
@Nullable @Nullable
@Override @Override
public List<TsKvEntry> apply(@Nullable List<Optional<TsKvEntry>> input) { public ReadTsKvQueryResult apply(@Nullable List<Optional<TsKvEntryAggWrapper>> input) {
return input == null ? Collections.emptyList() : input.stream().filter(v -> v.isPresent()).map(v -> v.get()).collect(Collectors.toList()); if (input == null) {
return new ReadTsKvQueryResult(query.getKey(), Collections.emptyList(), query.getStartTs());
} else {
long maxTs = query.getStartTs();
List<TsKvEntry> data = new ArrayList<>();
for (var opt : input) {
if (opt.isPresent()) {
TsKvEntryAggWrapper tsKvEntryAggWrapper = opt.get();
maxTs = Math.max(maxTs, tsKvEntryAggWrapper.getLastEntryTs());
data.add(tsKvEntryAggWrapper.getEntry());
}
}
return new ReadTsKvQueryResult(query.getKey(), data, maxTs);
}
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
} }
@ -302,13 +309,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
//Cleanup by TTL is native for Cassandra //Cleanup by TTL is native for Cassandra
} }
private ListenableFuture<List<TsKvEntry>> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { private ListenableFuture<ReadTsKvQueryResult> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) {
long minPartition = toPartitionTs(query.getStartTs()); long minPartition = toPartitionTs(query.getStartTs());
long maxPartition = toPartitionTs(query.getEndTs()); long maxPartition = toPartitionTs(query.getEndTs());
final ListenableFuture<List<Long>> partitionsListFuture = getPartitionsFuture(tenantId, query, entityId, minPartition, maxPartition); final ListenableFuture<List<Long>> partitionsListFuture = getPartitionsFuture(tenantId, query, entityId, minPartition, maxPartition);
final SimpleListenableFuture<List<TsKvEntry>> resultFuture = new SimpleListenableFuture<>(); final SimpleListenableFuture<List<TsKvEntry>> resultFuture = new SimpleListenableFuture<>();
Futures.addCallback(partitionsListFuture, new FutureCallback<List<Long>>() { Futures.addCallback(partitionsListFuture, new FutureCallback<>() {
@Override @Override
public void onSuccess(@Nullable List<Long> partitions) { public void onSuccess(@Nullable List<Long> partitions) {
TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions); TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions);
@ -321,7 +328,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
}, readResultsProcessingExecutor); }, readResultsProcessingExecutor);
return resultFuture; return Futures.transform(resultFuture, tsKvEntries -> {
long lastTs = query.getStartTs();
if (tsKvEntries != null) {
lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs());
}
return new ReadTsKvQueryResult(query.getKey(), tsKvEntries, lastTs);
}, MoreExecutors.directExecutor());
} }
private long toPartitionTs(long ts) { private long toPartitionTs(long ts) {
@ -379,7 +392,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD
} }
} }
private ListenableFuture<Optional<TsKvEntry>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query, long minPartition, long maxPartition) { private ListenableFuture<Optional<TsKvEntryAggWrapper>> findAndAggregateAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query, long minPartition, long maxPartition) {
final Aggregation aggregation = query.getAggregation(); final Aggregation aggregation = query.getAggregation();
final String key = query.getKey(); final String key = query.getKey();
final long startTs = query.getStartTs(); final long startTs = query.getStartTs();

6
dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java

@ -33,6 +33,7 @@ import org.thingsboard.server.common.data.kv.Aggregation;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult;
import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.model.ModelConstants;
@ -145,9 +146,10 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes
long endTs = query.getStartTs() - 1; long endTs = query.getStartTs() - 1;
ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1,
Aggregation.NONE, DESC_ORDER); Aggregation.NONE, DESC_ORDER);
ListenableFuture<List<TsKvEntry>> future = aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); ListenableFuture<ReadTsKvQueryResult> future = aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery);
return Futures.transformAsync(future, entryList -> { return Futures.transformAsync(future, result -> {
var entryList = result.getData();
if (entryList.size() == 1) { if (entryList.size() == 1) {
TsKvEntry entry = entryList.get(0); TsKvEntry entry = entryList.get(0);
return Futures.transform(saveLatest(tenantId, entityId, entryList.get(0)), v -> new TsKvLatestRemovingResult(entry), MoreExecutors.directExecutor()); return Futures.transform(saveLatest(tenantId, entityId, entryList.get(0)), v -> new TsKvLatestRemovingResult(entry), MoreExecutors.directExecutor());

4
dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java

@ -20,16 +20,18 @@ import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.common.data.kv.DeleteTsKvQuery; import org.thingsboard.server.common.data.kv.DeleteTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.List; import java.util.List;
import java.util.Map;
/** /**
* @author Andrew Shvayka * @author Andrew Shvayka
*/ */
public interface TimeseriesDao { public interface TimeseriesDao {
ListenableFuture<List<TsKvEntry>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries); ListenableFuture<List<ReadTsKvQueryResult>> findAllAsync(TenantId tenantId, EntityId entityId, List<ReadTsKvQuery> queries);
ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl); ListenableFuture<Integer> save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl);

5
dao/src/test/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDaoTest.java

@ -21,6 +21,7 @@ import org.junit.Before;
import org.junit.Test; import org.junit.Test;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvEntry;
import java.util.Optional; import java.util.Optional;
@ -48,9 +49,9 @@ public class AbstractChunkedAggregationTimeseriesDaoTest {
@Before @Before
public void setUp() throws Exception { public void setUp() throws Exception {
tsDao = spy(AbstractChunkedAggregationTimeseriesDao.class); tsDao = spy(AbstractChunkedAggregationTimeseriesDao.class);
ListenableFuture<Optional<TsKvEntry>> optionalListenableFuture = Futures.immediateFuture(Optional.of(mock(TsKvEntry.class))); Optional<TsKvEntry> optionalListenableFuture = Optional.of(mock(TsKvEntry.class));
willReturn(optionalListenableFuture).given(tsDao).findAndAggregateAsync(any(), anyString(), anyLong(), anyLong(), anyLong(), any()); willReturn(optionalListenableFuture).given(tsDao).findAndAggregateAsync(any(), anyString(), anyLong(), anyLong(), anyLong(), any());
willReturn(Futures.immediateFuture(mock(TsKvEntry.class))).given(tsDao).getTskvEntriesFuture(any()); willReturn(Futures.immediateFuture(mock(ReadTsKvQueryResult.class))).given(tsDao).getReadTsKvQueryResultFuture(any(), any());
} }
@Test @Test

Loading…
Cancel
Save