diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java index adb816a0ba..7599f94539 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java +++ b/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.kv.BaseReadTsKvQuery; 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.page.PageData; 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.stream.Collectors; +@SuppressWarnings("UnstableApiUsage") @Slf4j @TbCoreComponent @Service @@ -429,23 +431,34 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } else { finalTsKvQueryList = tsKvQueryList; } - Map>> fetchResultMap = new HashMap<>(); - ctx.getData().getData().forEach(entityData -> fetchResultMap.put(entityData, - tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList))); + Map>> fetchResultMap = new HashMap<>(); + List entityDataList = ctx.getData().getData(); + entityDataList.forEach(entityData -> fetchResultMap.put(entityData, + tsService.findAllByQueries(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList))); return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { + // Map that holds last ts for each key for each entity. + Map> lastTsEntityMap = new HashMap<>(); fetchResultMap.forEach((entityData, future) -> { - Map> keyData = new LinkedHashMap<>(); - cmd.getKeys().forEach(key -> keyData.put(key, new ArrayList<>())); try { - List entityTsData = future.get(); - if (entityTsData != null) { - entityTsData.forEach(entry -> keyData.get(entry.getKey()).add(new TsValue(entry.getTs(), entry.getValueAsString()))); + Map lastTsMap = new HashMap<>(); + lastTsEntityMap.put(entityData, lastTsMap); + + List 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()) { - entityData.getTimeseries().values().forEach(dataArray -> { - Arrays.sort(dataArray, (o1, o2) -> Long.compare(o2.getTs(), o1.getTs())); - }); + entityData.getTimeseries().values().forEach(dataArray -> Arrays.sort(dataArray, (o1, o2) -> Long.compare(o2.getTs(), o1.getTs()))); } } catch (InterruptedException | ExecutionException 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()); ctx.setInitialDataSent(true); } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); + update = new EntityDataUpdate(ctx.getCmdId(), null, entityDataList, ctx.getMaxEntitiesPerDataSubscription()); } 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.getData().getData().forEach(ed -> ed.getTimeseries().clear()); + entityDataList.forEach(ed -> ed.getTimeseries().clear()); } finally { ctx.getWsLock().unlock(); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java index a8283a1e4d..0238921783 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java @@ -86,7 +86,7 @@ public abstract class TbAbstractDataSubCtx newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity(), (a,b)-> a)); + Map newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity(), (a, b) -> a)); if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); } else { @@ -122,8 +122,13 @@ public abstract class TbAbstractDataSubCtx keys, long startTs, long endTs) { - createSubscriptions(keys, false, startTs, endTs); + public void createTimeseriesSubscriptions(Map> entityKeyStates, long startTs, long 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 keys, boolean latestValues, long startTs, long endTs) { @@ -191,6 +196,10 @@ public abstract class TbAbstractDataSubCtx keyStates) { log.trace("[{}][{}][{}] Creating time-series subscription for [{}] with keys: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), keyStates); return TbTimeseriesSubscription.builder() .serviceId(serviceId) diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java index b8688b1524..0015ee73be 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesService.java +++ b/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.kv.DeleteTsKvQuery; 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.TsKvEntry; import java.util.Collection; import java.util.List; +import java.util.Map; /** * @author Andrew Shvayka */ public interface TimeseriesService { + ListenableFuture> findAllByQueries(TenantId tenantId, EntityId entityId, List queries); + ListenableFuture> findAll(TenantId tenantId, EntityId entityId, List queries); ListenableFuture> findLatest(TenantId tenantId, EntityId entityId, Collection keys); diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/AggTsKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/AggTsKvEntry.java new file mode 100644 index 0000000000..85330d3577 --- /dev/null +++ b/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); + } +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java index 6fc32e627b..e5a7ac65e0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/BasicTsKvEntry.java +++ b/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 { private static final int MAX_CHARS_PER_DATA_POINT = 512; - private final long ts; + protected final long ts; private final KvEntry kv; public BasicTsKvEntry(long ts, KvEntry kv) { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/ReadTsKvQueryResult.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/ReadTsKvQueryResult.java new file mode 100644 index 0000000000..dd1b50bb40 --- /dev/null +++ b/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 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 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]; + } + } + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java index 20d25d5399..06ac0d84a3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntry.java @@ -16,10 +16,11 @@ package org.thingsboard.server.common.data.kv; import com.fasterxml.jackson.annotation.JsonIgnore; +import org.thingsboard.server.common.data.query.TsValue; /** * Represents time series KV data entry - * + * * @author ashvayka * */ @@ -30,4 +31,9 @@ public interface TsKvEntry extends KvEntry { @JsonIgnore int getDataPoints(); + @JsonIgnore + default TsValue toTsValue() { + return new TsValue(getTs(), getValueAsString()); + } + } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntryAggWrapper.java b/common/data/src/main/java/org/thingsboard/server/common/data/kv/TsKvEntryAggWrapper.java new file mode 100644 index 0000000000..2825008485 --- /dev/null +++ b/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; + +} diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java index a086f1c21f..b7e83521e8 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/TsValue.java @@ -15,11 +15,21 @@ */ package org.thingsboard.server.common.data.query; +import com.fasterxml.jackson.annotation.JsonInclude; import lombok.Data; +import lombok.RequiredArgsConstructor; @Data +@RequiredArgsConstructor +@JsonInclude(JsonInclude.Include.NON_NULL) public class TsValue { private final long ts; private final String value; + private final Long count; + + public TsValue(long ts, String value) { + this(ts, value, null); + } + } diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java b/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java index 171417402d..5218de60e4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/ModelConstants.java +++ b/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[] 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 = - 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 = - 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; public static String min(String s) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java index 1824c96247..5a5985d125 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java +++ b/dao/src/main/java/org/thingsboard/server/dao/model/sql/AbstractTsKvEntity.java @@ -16,6 +16,7 @@ package org.thingsboard.server.dao.model.sql; 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.BooleanDataEntry; import org.thingsboard.server.common.data.kv.DoubleDataEntry; @@ -80,6 +81,18 @@ public abstract class AbstractTsKvEntity implements ToData { @Transient protected String strKey; + @Transient + protected Long aggValuesLastTs; + @Transient + protected Long aggValuesCount; + + public AbstractTsKvEntity() { + } + + public AbstractTsKvEntity(Long aggValuesLastTs) { + this.aggValuesLastTs = aggValuesLastTs; + } + public abstract boolean isNotEmpty(); protected static boolean isAllNull(Object... args) { @@ -105,7 +118,12 @@ public abstract class AbstractTsKvEntity implements ToData { } else if (jsonValue != null) { 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); + } } } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/ts/TimescaleTsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/ts/TimescaleTsKvEntity.java index 88c1693e10..9e072e25e4 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/timescale/ts/TimescaleTsKvEntity.java +++ b/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 = "strValue", 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 = "doubleValueCount", 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(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)) { this.strValue = strValue; } @@ -135,6 +139,7 @@ public final class TimescaleTsKvEntity extends AbstractTsKvEntity { } else { this.doubleValue = 0.0; } + this.aggValuesCount = totalCount; break; case SUM: 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)) { this.ts = tsBucket + interval / 2; if (booleanValueCount != 0) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java b/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java index 1c0277f8b6..15fd30c742 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/model/sqlts/ts/TsKvEntity.java +++ b/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; import lombok.Data; +import lombok.EqualsAndHashCode; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; import javax.persistence.Entity; import javax.persistence.IdClass; import javax.persistence.Table; +import javax.persistence.Transient; +@EqualsAndHashCode(callSuper = true) @Data @Entity @Table(name = "ts_kv") @@ -31,11 +34,13 @@ public final class TsKvEntity extends AbstractTsKvEntity { public TsKvEntity() { } - public TsKvEntity(String strValue) { + public TsKvEntity(String strValue, Long aggValuesLastTs) { + super(aggValuesLastTs); 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)) { switch (aggType) { case AVG: @@ -52,6 +57,7 @@ public final class TsKvEntity extends AbstractTsKvEntity { } else { this.doubleValue = 0.0; } + this.aggValuesCount = totalCount; break; case SUM: 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 (booleanValueCount != 0) { this.longValue = booleanValueCount; diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java index 69f16d4bbc..365d1240b3 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDao.java +++ b/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.ListenableFuture; -import com.google.common.util.concurrent.MoreExecutors; -import com.google.common.util.concurrent.SettableFuture; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; 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.TenantId; 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.ReadTsKvQuery; +import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.common.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; @@ -47,10 +45,9 @@ import java.util.Comparator; import java.util.List; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.CompletableFuture; import java.util.function.Function; -import java.util.stream.Collectors; +@SuppressWarnings("UnstableApiUsage") @Slf4j public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSqlTimeseriesDao implements TimeseriesDao { @@ -81,7 +78,7 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq Comparator.comparing((Function) AbstractTsKvEntity::getEntityId) .thenComparing(AbstractTsKvEntity::getKey) .thenComparing(AbstractTsKvEntity::getTs) - ); + ); } @PreDestroy @@ -114,16 +111,16 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { + public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { return processFindAllAsync(tenantId, entityId, queries); } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + public ListenableFuture findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { if (query.getAggregation() == Aggregation.NONE) { - return findAllAsyncWithLimit(entityId, query); + return Futures.immediateFuture(findAllAsyncWithLimit(entityId, query)); } else { - List>> futures = new ArrayList<>(); + List>> futures = new ArrayList<>(); long endPeriod = query.getEndTs(); long startPeriod = query.getStartTs(); long step = query.getInterval(); @@ -131,15 +128,16 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq long startTs = startPeriod; long endTs = Math.min(startPeriod + step, endPeriod + 1); long ts = startTs + (endTs - startTs) / 2; - ListenableFuture> aggregateTsKvEntry = findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation()); + ListenableFuture> aggregateTsKvEntry = + service.submit(() -> findAndAggregateAsync(entityId, query.getKey(), startTs, endTs, ts, query.getAggregation())); futures.add(aggregateTsKvEntry); startPeriod = endTs; } - return getTskvEntriesFuture(Futures.allAsList(futures)); + return getReadTsKvQueryResultFuture(query, Futures.allAsList(futures)); } } - private ListenableFuture> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { + private ReadTsKvQueryResult findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { Integer keyId = getOrSaveKeyId(query.getKey()); List tsKvEntities = tsKvRepository.findAllWithLimit( entityId.getId(), @@ -147,125 +145,50 @@ public abstract class AbstractChunkedAggregationTimeseriesDao extends AbstractSq query.getStartTs(), query.getEndTs(), 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())); - return Futures.immediateFuture(DaoUtil.convertDataList(tsKvEntities)); - } - - ListenableFuture> findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { - List> entitiesFutures = new ArrayList<>(); - switchAggregation(entityId, key, startTs, endTs, aggregation, entitiesFutures); - return Futures.transform(setFutures(entitiesFutures), entity -> { - if (entity != null && entity.isNotEmpty()) { - entity.setEntityId(entityId.getId()); - entity.setStrKey(key); - entity.setTs(ts); - return Optional.of(DaoUtil.getData(entity)); - } else { - return Optional.empty(); - } - }, MoreExecutors.directExecutor()); + List tsKvEntries = DaoUtil.convertDataList(tsKvEntities); + long lastTs = tsKvEntries.stream().map(TsKvEntry::getTs).max(Long::compare).orElse(query.getStartTs()); + return new ReadTsKvQueryResult(query.getKey(), tsKvEntries, lastTs); + } + + Optional findAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long ts, Aggregation aggregation) { + TsKvEntity entity = switchAggregation(entityId, key, startTs, endTs, aggregation); + if (entity != null && entity.isNotEmpty()) { + entity.setEntityId(entityId.getId()); + entity.setStrKey(key); + entity.setTs(ts); + return Optional.of(entity); + } else { + return Optional.empty(); + } } - protected void switchAggregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation, List> entitiesFutures) { + protected TsKvEntity switchAggregation(EntityId entityId, String key, long startTs, long endTs, Aggregation aggregation) { + var keyId = getOrSaveKeyId(key); switch (aggregation) { case AVG: - findAvg(entityId, key, startTs, endTs, entitiesFutures); - break; + return tsKvRepository.findAvg(entityId.getId(), keyId, startTs, endTs); case MAX: - findMax(entityId, key, startTs, endTs, entitiesFutures); - break; + var max = tsKvRepository.findNumericMax(entityId.getId(), keyId, startTs, endTs); + if (max.isNotEmpty()) { + return max; + } else { + return tsKvRepository.findStringMax(entityId.getId(), keyId, startTs, endTs); + } case MIN: - findMin(entityId, key, startTs, endTs, entitiesFutures); - break; + var min = tsKvRepository.findNumericMin(entityId.getId(), keyId, startTs, endTs); + if (min.isNotEmpty()) { + return min; + } else { + return tsKvRepository.findStringMin(entityId.getId(), keyId, startTs, endTs); + } case SUM: - findSum(entityId, key, startTs, endTs, entitiesFutures); - break; + return tsKvRepository.findSum(entityId.getId(), keyId, startTs, endTs); case COUNT: - findCount(entityId, key, startTs, endTs, entitiesFutures); - break; + return tsKvRepository.findCount(entityId.getId(), keyId, startTs, endTs); default: throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); } } - - protected void findCount(EntityId entityId, String key, long startTs, long endTs, List> 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> 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> 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> 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> entitiesFutures) { - Integer keyId = getOrSaveKeyId(key); - entitiesFutures.add(tsKvRepository.findAvg( - entityId.getId(), - keyId, - startTs, - endTs)); - } - - protected SettableFuture setFutures(List> entitiesFutures) { - SettableFuture listenableFuture = SettableFuture.create(); - CompletableFuture> 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; - } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java index aa0d9b1c47..3de21e41bf 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AbstractSqlTimeseriesDao.java +++ b/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.TenantId; 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.dao.model.ModelConstants; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; @@ -38,6 +39,7 @@ import java.util.Objects; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +@SuppressWarnings("UnstableApiUsage") @Slf4j public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseriesDao implements AggregationTimeseriesDao { @@ -86,22 +88,19 @@ public abstract class AbstractSqlTimeseriesDao extends BaseAbstractSqlTimeseries } } - protected ListenableFuture> processFindAllAsync(TenantId tenantId, EntityId entityId, List queries) { - List>> futures = queries + protected ListenableFuture> processFindAllAsync(TenantId tenantId, EntityId entityId, List queries) { + List> futures = queries .stream() .map(query -> findAllAsync(tenantId, entityId, query)) .collect(Collectors.toList()); return Futures.transform(Futures.allAsList(futures), new Function<>() { @Nullable @Override - public List apply(@Nullable List> results) { + public List apply(@Nullable List results) { if (results == null || results.isEmpty()) { return null; } - return results.stream() - .filter(Objects::nonNull) - .flatMap(List::stream) - .collect(Collectors.toList()); + return results.stream().filter(Objects::nonNull).collect(Collectors.toList()); } }, service); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AggregationTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/AggregationTimeseriesDao.java index 31270bacc7..994555aaa5 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/AggregationTimeseriesDao.java +++ b/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.TenantId; 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 java.util.List; public interface AggregationTimeseriesDao { - ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); + ListenableFuture findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query); } \ No newline at end of file diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java index 82e54168bc..314117fe06 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java +++ b/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.springframework.beans.factory.annotation.Autowired; 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.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.TsKvDictionaryCompositeKey; +import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sqlts.dictionary.TsKvDictionaryRepository; import javax.annotation.Nullable; import java.util.List; +import java.util.Objects; import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; @@ -80,18 +86,20 @@ public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeni return keyId; } - protected ListenableFuture> getTskvEntriesFuture(ListenableFuture>> future) { - return Futures.transform(future, new Function>, List>() { + protected ListenableFuture getReadTsKvQueryResultFuture(ReadTsKvQuery query, ListenableFuture>> future) { + return Futures.transform(future, new Function<>() { @Nullable @Override - public List apply(@Nullable List> results) { + public ReadTsKvQueryResult apply(@Nullable List> results) { if (results == null || results.isEmpty()) { return null; } - return results.stream() - .filter(Optional::isPresent) - .map(Optional::get) - .collect(Collectors.toList()); + List data = results.stream().filter(Optional::isPresent).map(Optional::get).collect(Collectors.toList()); + var lastTs = data.stream().map(AbstractTsKvEntity::getAggValuesLastTs).filter(Objects::nonNull).max(Long::compare); + if (lastTs.isEmpty()) { + 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); } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java index db30adead2..50975f21de 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/SqlTimeseriesLatestDao.java +++ b/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.DeleteTsKvQuery; 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.TsKvEntry; import org.thingsboard.server.common.data.kv.TsKvLatestRemovingResult; @@ -190,7 +191,8 @@ public class SqlTimeseriesLatestDao extends BaseAbstractSqlTimeseriesDao impleme long endTs = query.getStartTs() - 1; ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, Aggregation.NONE, DESC_ORDER); - return aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); + return Futures.transform(aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery), + ReadTsKvQueryResult::getData, MoreExecutors.directExecutor()); } protected ListenableFuture getFindLatestFuture(EntityId entityId, String key) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java index 9cf796b0aa..db0a1753d8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/AggregationRepository.java +++ b/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_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 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 "; - - 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 "; - - 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_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 "; - - 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 "; + 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 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, 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 private EntityManager entityManager; - @Async - public CompletableFuture> findAvg(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { - @SuppressWarnings("unchecked") - List resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG); - return CompletableFuture.supplyAsync(() -> resultList); + @SuppressWarnings("unchecked") + public List findAvg(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { + return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_AVG); } - @Async - public CompletableFuture> findMax(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { - @SuppressWarnings("unchecked") - List resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX); - return CompletableFuture.supplyAsync(() -> resultList); + @SuppressWarnings("unchecked") + public List findMax(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { + return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MAX); } - @Async - public CompletableFuture> findMin(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { - @SuppressWarnings("unchecked") - List resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN); - return CompletableFuture.supplyAsync(() -> resultList); + @SuppressWarnings("unchecked") + public List findMin(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { + return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_MIN); } - @Async - public CompletableFuture> findSum(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { - @SuppressWarnings("unchecked") - List resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM); - return CompletableFuture.supplyAsync(() -> resultList); + @SuppressWarnings("unchecked") + public List findSum(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { + return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_SUM); } - @Async - public CompletableFuture> findCount(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { - @SuppressWarnings("unchecked") - List resultList = getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_COUNT); - return CompletableFuture.supplyAsync(() -> resultList); + @SuppressWarnings("unchecked") + public List findCount(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs) { + return getResultList(entityId, entityKey, timeBucket, startTs, endTs, FIND_COUNT); } private List getResultList(UUID entityId, int entityKey, long timeBucket, long startTs, long endTs, String query) { @@ -96,5 +121,4 @@ public class AggregationRepository { .getResultList(); } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java index 9c03a9e3dd..bfefbf4bd8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/timescale/TimescaleTimeseriesDao.java +++ b/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.DeleteTsKvQuery; 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.stats.StatsFactory; import org.thingsboard.server.dao.DaoUtil; 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.ts.TsKvEntity; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sqlts.AbstractSqlTimeseriesDao; @@ -101,13 +103,13 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { + public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { return processFindAllAsync(tenantId, entityId, queries); } @Override public ListenableFuture 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(); Integer keyId = getOrSaveKeyId(strKey); TimescaleTsKvEntity entity = new TimescaleTsKvEntity(); @@ -148,15 +150,15 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + public ListenableFuture findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { if (query.getAggregation() == Aggregation.NONE) { - return findAllAsyncWithLimit(entityId, query); + return Futures.immediateFuture(findAllAsyncWithLimit(entityId, query)); } else { long startTs = query.getStartTs(); long endTs = query.getEndTs(); long timeBucket = query.getInterval(); - ListenableFuture>> future = findAllAndAggregateAsync(entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation()); - return getTskvEntriesFuture(future); + List> data = findAllAndAggregateAsync(entityId, query.getKey(), startTs, endTs, timeBucket, query.getAggregation()); + return getReadTsKvQueryResultFuture(query, Futures.immediateFuture(data)); } } @@ -165,7 +167,7 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements super.cleanup(systemTtl); } - private ListenableFuture> findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { + private ReadTsKvQueryResult findAllAsyncWithLimit(EntityId entityId, ReadTsKvQuery query) { String strKey = query.getKey(); Integer keyId = getOrSaveKeyId(strKey); List timescaleTsKvEntities = tsKvRepository.findAllWithLimit( @@ -174,105 +176,49 @@ public class TimescaleTimeseriesDao extends AbstractSqlTimeseriesDao implements query.getStartTs(), query.getEndTs(), 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)); - return Futures.immediateFuture(DaoUtil.convertDataList(timescaleTsKvEntities)); - } - - private ListenableFuture>> findAllAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) { - CompletableFuture> listCompletableFuture = switchAggregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId()); - SettableFuture> listenableFuture = SettableFuture.create(); - listCompletableFuture.whenComplete((timescaleTsKvEntities, throwable) -> { - if (throwable != null) { - listenableFuture.setException(throwable); - } else { - listenableFuture.set(timescaleTsKvEntities); - } - }); - return Futures.transform(listenableFuture, timescaleTsKvEntities -> { - if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) { - List> result = new ArrayList<>(); - timescaleTsKvEntities.forEach(entity -> { - if (entity != null && entity.isNotEmpty()) { - entity.setEntityId(entityId.getId()); - entity.setStrKey(key); - result.add(Optional.of(DaoUtil.getData(entity))); - } else { - result.add(Optional.empty()); - } - }); - return result; - } else { - return Collections.emptyList(); - } - }, MoreExecutors.directExecutor()); + 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 List> findAllAndAggregateAsync(EntityId entityId, String key, long startTs, long endTs, long timeBucket, Aggregation aggregation) { + List timescaleTsKvEntities = switchAggregation(key, startTs, endTs, timeBucket, aggregation, entityId.getId()); + if (!CollectionUtils.isEmpty(timescaleTsKvEntities)) { + List> result = new ArrayList<>(); + timescaleTsKvEntities.forEach(entity -> { + if (entity != null && entity.isNotEmpty()) { + entity.setEntityId(entityId.getId()); + entity.setStrKey(key); + result.add(Optional.of(entity)); + } else { + result.add(Optional.empty()); + } + }); + return result; + } else { + return Collections.emptyList(); + } } - private CompletableFuture> switchAggregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId) { + private List switchAggregation(String key, long startTs, long endTs, long timeBucket, Aggregation aggregation, UUID entityId) { + Integer keyId = getOrSaveKeyId(key); switch (aggregation) { case AVG: - return findAvg(key, startTs, endTs, timeBucket, entityId); + return aggregationRepository.findAvg(entityId, keyId, timeBucket, startTs, endTs); case MAX: - return findMax(key, startTs, endTs, timeBucket, entityId); + return aggregationRepository.findMax(entityId, keyId, timeBucket, startTs, endTs); case MIN: - return findMin(key, startTs, endTs, timeBucket, entityId); + return aggregationRepository.findMin(entityId, keyId, timeBucket, startTs, endTs); case SUM: - return findSum(key, startTs, endTs, timeBucket, entityId); + return aggregationRepository.findSum(entityId, keyId, timeBucket, startTs, endTs); case COUNT: - return findCount(key, startTs, endTs, timeBucket, entityId); + return aggregationRepository.findCount(entityId, keyId, timeBucket, startTs, endTs); default: throw new IllegalArgumentException("Not supported aggregation type: " + aggregation); } } - private CompletableFuture> 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> 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> 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> 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> findAvg(String key, long startTs, long endTs, long timeBucket, UUID entityId) { - Integer keyId = getOrSaveKeyId(key); - return aggregationRepository.findAvg( - entityId, - keyId, - timeBucket, - startTs, - endTs); - } - } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java index c603d9d92b..1d9817a5bd 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/ts/TsKvRepository.java +++ b/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.Query; import org.springframework.data.repository.query.Param; -import org.springframework.scheduling.annotation.Async; import org.springframework.transaction.annotation.Transactional; import org.thingsboard.server.dao.model.sqlts.ts.TsKvCompositeKey; import org.thingsboard.server.dao.model.sqlts.ts.TsKvEntity; import java.util.List; import java.util.UUID; -import java.util.concurrent.CompletableFuture; public interface TsKvRepository extends JpaRepository { @@ -48,82 +46,75 @@ public interface TsKvRepository extends JpaRepository= :startTs AND tskv.ts < :endTs") - CompletableFuture findStringMax(@Param("entityId") UUID entityId, + TsKvEntity findStringMax(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async @Query("SELECT new TsKvEntity(MAX(COALESCE(tskv.longValue, -9223372036854775807)), " + "MAX(COALESCE(tskv.doubleValue, -1.79769E+308)), " + "SUM(CASE WHEN tskv.longValue 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") - CompletableFuture findNumericMax(@Param("entityId") UUID entityId, + TsKvEntity findNumericMax(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async - @Query("SELECT new TsKvEntity(MIN(tskv.strValue)) FROM TsKvEntity tskv " + + @Query("SELECT new TsKvEntity(MIN(tskv.strValue), MAX(tskv.ts)) FROM TsKvEntity tskv " + "WHERE tskv.strValue IS NOT NULL " + "AND tskv.entityId = :entityId AND tskv.key = :entityKey AND tskv.ts >= :startTs AND tskv.ts < :endTs") - CompletableFuture findStringMin(@Param("entityId") UUID entityId, + TsKvEntity findStringMin(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async @Query("SELECT new TsKvEntity(MIN(COALESCE(tskv.longValue, 9223372036854775807)), " + "MIN(COALESCE(tskv.doubleValue, 1.79769E+308)), " + "SUM(CASE WHEN tskv.longValue 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") - CompletableFuture findNumericMin( + TsKvEntity findNumericMin( @Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async @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.longValue 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") - CompletableFuture findCount(@Param("entityId") UUID entityId, + TsKvEntity findCount(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async @Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " + "SUM(COALESCE(tskv.doubleValue, 0.0)), " + "SUM(CASE WHEN tskv.longValue 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") - CompletableFuture findAvg(@Param("entityId") UUID entityId, + TsKvEntity findAvg(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); - @Async @Query("SELECT new TsKvEntity(SUM(COALESCE(tskv.longValue, 0)), " + "SUM(COALESCE(tskv.doubleValue, 0.0)), " + "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') 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") - CompletableFuture findSum(@Param("entityId") UUID entityId, + TsKvEntity findSum(@Param("entityId") UUID entityId, @Param("entityKey") int entityKey, @Param("startTs") long startTs, @Param("endTs") long endTs); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java index 975c132011..7a19e9fbf8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/AggregatePartitionsFunction.java +++ b/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.ListenableFuture; 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.BasicTsKvEntry; 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.StringDataEntry; import org.thingsboard.server.common.data.kv.TsKvEntry; +import org.thingsboard.server.common.data.kv.TsKvEntryAggWrapper; import org.thingsboard.server.dao.nosql.TbResultSet; import javax.annotation.Nullable; @@ -40,7 +42,7 @@ import java.util.stream.Collectors; * Created by ashvayka on 20.02.17. */ @Slf4j -public class AggregatePartitionsFunction implements com.google.common.util.concurrent.AsyncFunction, Optional> { +public class AggregatePartitionsFunction implements com.google.common.util.concurrent.AsyncFunction, Optional> { private static final int LONG_CNT_POS = 0; 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 STR_POS = 8; private static final int JSON_POS = 9; + private static final int MAX_TS_POS = 10; private final Aggregation aggregation; private final String key; @@ -66,29 +69,29 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu } @Override - public ListenableFuture> apply(@Nullable List rsList) { - log.trace("[{}][{}][{}] Going to aggregate data", key, ts, aggregation); - if (rsList == null || rsList.isEmpty()) { - return Futures.immediateFuture(Optional.empty()); - } - return Futures.transform( - Futures.allAsList( - rsList.stream().map(rs -> rs.allRows(this.executor)) - .collect(Collectors.toList())), - rowsList -> { - try { - AggregationResult aggResult = new AggregationResult(); - for (List rs : rowsList) { - for (Row row : rs) { - processResultSetRow(row, aggResult); + public ListenableFuture> apply(@Nullable List rsList) { + log.trace("[{}][{}][{}] Going to aggregate data", key, ts, aggregation); + if (rsList == null || rsList.isEmpty()) { + return Futures.immediateFuture(Optional.empty()); + } + return Futures.transform( + Futures.allAsList( + rsList.stream().map(rs -> rs.allRows(this.executor)) + .collect(Collectors.toList())), + rowsList -> { + try { + AggregationResult aggResult = new AggregationResult(); + for (List rs : rowsList) { + for (Row row : rs) { + processResultSetRow(row, aggResult); + } + } + return processAggregationResult(aggResult); + } catch (Exception e) { + log.error("[{}][{}][{}] Failed to aggregate data", key, ts, aggregation, e); + return Optional.empty(); } - } - return processAggregationResult(aggResult); - } catch (Exception e) { - log.error("[{}][{}][{}] Failed to aggregate data", key, ts, aggregation, e); - return Optional.empty(); - } - }, this.executor); + }, this.executor); } 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 strCount = row.getLong(STR_CNT_POS); long jsonCount = row.getLong(JSON_CNT_POS); + long aggValuesLastTs = row.getLong(MAX_TS_POS); if (longCount > 0 || doubleCount > 0) { if (longCount > 0) { @@ -134,6 +138,8 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu return; } + aggResult.aggValuesLastTs = Math.max(aggResult.aggValuesLastTs, aggValuesLastTs); + if (aggregation == Aggregation.COUNT) { aggResult.count += curCount; } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { @@ -231,34 +237,37 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu } } - private Optional processAggregationResult(AggregationResult aggResult) { + private Optional processAggregationResult(AggregationResult aggResult) { Optional result; if (aggResult.dataType == null) { result = Optional.empty(); } else if (aggregation == Aggregation.COUNT) { result = Optional.of(new BasicTsKvEntry(ts, new LongDataEntry(key, aggResult.count))); } else if (aggregation == Aggregation.AVG || aggregation == Aggregation.SUM) { - result = processAvgOrSumResult(aggResult); + result = processAvgOrSumResult(aggregation, aggResult); } else if (aggregation == Aggregation.MIN || aggregation == Aggregation.MAX) { result = processMinOrMaxResult(aggResult); } else { result = Optional.empty(); } - if (!result.isPresent()) { + if (result.isEmpty()) { log.trace("[{}][{}][{}] Aggregated data is empty.", key, ts, aggregation); } - return result; + return result.map(tsKvEntry -> new TsKvEntryAggWrapper(tsKvEntry, aggResult.aggValuesLastTs)); } - private Optional processAvgOrSumResult(AggregationResult aggResult) { + private Optional processAvgOrSumResult(Aggregation aggregation, AggregationResult aggResult) { if (aggResult.count == 0 || (aggResult.dataType == DataType.DOUBLE && aggResult.dValue == null) || (aggResult.dataType == DataType.LONG && aggResult.lValue == null)) { return Optional.empty(); } else if (aggResult.dataType == DataType.DOUBLE || aggResult.dataType == DataType.LONG) { if (aggregation == Aggregation.AVG || aggResult.hasDouble) { 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 { - 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(); @@ -291,5 +300,6 @@ public class AggregatePartitionsFunction implements com.google.common.util.concu Long lValue = null; long count = 0; boolean hasDouble = false; + long aggValuesLastTs = 0; } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java index 8254b2bee0..b785ae2f9d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/BaseTimeseriesService.java @@ -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.DeleteTsKvQuery; 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.TsKvLatestRemovingResult; import org.thingsboard.server.dao.entityview.EntityViewService; @@ -52,6 +53,7 @@ import static org.thingsboard.server.common.data.StringUtils.isBlank; /** * @author Andrew Shvayka */ +@SuppressWarnings("UnstableApiUsage") @Service @Slf4j 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_WITHOUT_LATEST = 2; private static final int DELETES_PER_ENTRY = INSERTS_PER_ENTRY; - public static final Function, Integer> SUM_ALL_INTEGERS = new Function, Integer>() { + public static final Function, Integer> SUM_ALL_INTEGERS = new Function<>() { @Override public @Nullable Integer apply(@Nullable List input) { int result = 0; @@ -87,7 +89,7 @@ public class BaseTimeseriesService implements TimeseriesService { private EntityViewService entityViewService; @Override - public ListenableFuture> findAll(TenantId tenantId, EntityId entityId, List queries) { + public ListenableFuture> findAllByQueries(TenantId tenantId, EntityId entityId, List queries) { validate(entityId); queries.forEach(this::validate); if (entityId.getEntityType().equals(EntityType.ENTITY_VIEW)) { @@ -103,6 +105,17 @@ public class BaseTimeseriesService implements TimeseriesService { return timeseriesDao.findAllAsync(tenantId, entityId, queries); } + @Override + public ListenableFuture> findAll(TenantId tenantId, EntityId entityId, List 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 public ListenableFuture> findLatest(TenantId tenantId, EntityId entityId, Collection keys) { validate(entityId); @@ -244,7 +257,7 @@ public class BaseTimeseriesService implements TimeseriesService { public ListenableFuture> removeAllLatest(TenantId tenantId, EntityId entityId) { validate(entityId); return Futures.transformAsync(this.findAllLatest(tenantId, entityId), latest -> { - if (!latest.isEmpty()) { + if (latest != null && !latest.isEmpty()) { Collection keys = latest.stream().map(TsKvEntry::getKey).collect(Collectors.toList()); return Futures.transform(this.removeLatest(tenantId, entityId, keys), res -> keys, MoreExecutors.directExecutor()); } else { diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java index 3b5e45f206..14a49027f7 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesDao.java @@ -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.KvEntry; 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.TsKvEntryAggWrapper; import org.thingsboard.server.dao.model.ModelConstants; import org.thingsboard.server.dao.nosql.TbResultSet; import org.thingsboard.server.dao.nosql.TbResultSetFuture; @@ -71,6 +73,7 @@ import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.literal; /** * @author Andrew Shvayka */ +@SuppressWarnings("UnstableApiUsage") @Component @Slf4j @NoSqlTsDao @@ -139,20 +142,10 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { - List>> futures = queries.stream().map(query -> findAllAsync(tenantId, entityId, query)).collect(Collectors.toList()); - return Futures.transform(Futures.allAsList(futures), new Function<>() { - @Nullable - @Override - public List apply(@Nullable List> results) { - if (results == null || results.isEmpty()) { - return null; - } - return results.stream() - .flatMap(List::stream) - .collect(Collectors.toList()); - } - }, readResultsProcessingExecutor); + public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries) { + List> futures = queries.stream() + .map(query -> findAllAsync(tenantId, entityId, query)).collect(Collectors.toList()); + return Futures.allAsList(futures); } @Override @@ -270,14 +263,14 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } @Override - public ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + public ListenableFuture findAllAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { if (query.getAggregation() == Aggregation.NONE) { return findAllAsyncWithLimit(tenantId, entityId, query); } else { long startPeriod = query.getStartTs(); long endPeriod = query.getEndTs(); long step = Math.max(query.getInterval(), MIN_AGGREGATION_STEP_MS); - List>> futures = new ArrayList<>(); + List>> futures = new ArrayList<>(); while (startPeriod <= endPeriod) { long startTs = startPeriod; 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))); startPeriod = endTs; } - ListenableFuture>> future = Futures.allAsList(futures); + ListenableFuture>> future = Futures.allAsList(futures); return Futures.transform(future, new Function<>() { @Nullable @Override - public List apply(@Nullable List> input) { - return input == null ? Collections.emptyList() : input.stream().filter(v -> v.isPresent()).map(v -> v.get()).collect(Collectors.toList()); + public ReadTsKvQueryResult apply(@Nullable List> input) { + if (input == null) { + return new ReadTsKvQueryResult(query.getKey(), Collections.emptyList(), query.getStartTs()); + } else { + long maxTs = query.getStartTs(); + List 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); } @@ -302,13 +309,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD //Cleanup by TTL is native for Cassandra } - private ListenableFuture> findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { + private ListenableFuture findAllAsyncWithLimit(TenantId tenantId, EntityId entityId, ReadTsKvQuery query) { long minPartition = toPartitionTs(query.getStartTs()); long maxPartition = toPartitionTs(query.getEndTs()); final ListenableFuture> partitionsListFuture = getPartitionsFuture(tenantId, query, entityId, minPartition, maxPartition); final SimpleListenableFuture> resultFuture = new SimpleListenableFuture<>(); - Futures.addCallback(partitionsListFuture, new FutureCallback>() { + Futures.addCallback(partitionsListFuture, new FutureCallback<>() { @Override public void onSuccess(@Nullable List partitions) { TsKvQueryCursor cursor = new TsKvQueryCursor(entityId.getEntityType().name(), entityId.getId(), query, partitions); @@ -321,7 +328,13 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } }, 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) { @@ -379,7 +392,7 @@ public class CassandraBaseTimeseriesDao extends AbstractCassandraBaseTimeseriesD } } - private ListenableFuture> findAndAggregateAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query, long minPartition, long maxPartition) { + private ListenableFuture> findAndAggregateAsync(TenantId tenantId, EntityId entityId, ReadTsKvQuery query, long minPartition, long maxPartition) { final Aggregation aggregation = query.getAggregation(); final String key = query.getKey(); final long startTs = query.getStartTs(); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java index 124af03a1a..2b3d62712f 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/CassandraBaseTimeseriesLatestDao.java +++ b/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.DeleteTsKvQuery; 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.TsKvLatestRemovingResult; import org.thingsboard.server.dao.model.ModelConstants; @@ -145,9 +146,10 @@ public class CassandraBaseTimeseriesLatestDao extends AbstractCassandraBaseTimes long endTs = query.getStartTs() - 1; ReadTsKvQuery findNewLatestQuery = new BaseReadTsKvQuery(query.getKey(), startTs, endTs, endTs - startTs, 1, Aggregation.NONE, DESC_ORDER); - ListenableFuture> future = aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); + ListenableFuture future = aggregationTimeseriesDao.findAllAsync(tenantId, entityId, findNewLatestQuery); - return Futures.transformAsync(future, entryList -> { + return Futures.transformAsync(future, result -> { + var entryList = result.getData(); if (entryList.size() == 1) { TsKvEntry entry = entryList.get(0); return Futures.transform(saveLatest(tenantId, entityId, entryList.get(0)), v -> new TsKvLatestRemovingResult(entry), MoreExecutors.directExecutor()); diff --git a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java index c075a51434..5fd26d400a 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/timeseries/TimeseriesDao.java @@ -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.kv.DeleteTsKvQuery; 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 java.util.List; +import java.util.Map; /** * @author Andrew Shvayka */ public interface TimeseriesDao { - ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries); + ListenableFuture> findAllAsync(TenantId tenantId, EntityId entityId, List queries); ListenableFuture save(TenantId tenantId, EntityId entityId, TsKvEntry tsKvEntry, long ttl); diff --git a/dao/src/test/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDaoTest.java b/dao/src/test/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDaoTest.java index db216f7f6e..40705809fd 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/sqlts/AbstractChunkedAggregationTimeseriesDaoTest.java +++ b/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.thingsboard.server.common.data.kv.BaseReadTsKvQuery; 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 java.util.Optional; @@ -48,9 +49,9 @@ public class AbstractChunkedAggregationTimeseriesDaoTest { @Before public void setUp() throws Exception { tsDao = spy(AbstractChunkedAggregationTimeseriesDao.class); - ListenableFuture> optionalListenableFuture = Futures.immediateFuture(Optional.of(mock(TsKvEntry.class))); + Optional optionalListenableFuture = Optional.of(mock(TsKvEntry.class)); 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