diff --git a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java index 581279d25f..e6aed81c4b 100644 --- a/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java +++ b/application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java @@ -155,7 +155,8 @@ public class EntityQueryController extends BaseController { @ApiOperation(value = "Find Aggregated Historical Entity Data by Query", notes = "Runs the entity data query and, for each matched entity, fetches a single aggregated value per key over [startTs, endTs] " + "with per-key aggregation function. Optional previousStartTs/previousEndTs per key add a comparison window. " + - "REST equivalent of the WebSocket AggHistoryCmd. The aggregated values are returned in the 'aggLatest' field of each EntityData.") + "REST equivalent of the WebSocket AggHistoryCmd. The aggregated values are returned in the 'aggLatest' field of each EntityData. " + + "Note: a separate timeseries aggregation query is issued per matched entity, so a large pageSize fans out proportionally - callers should page responsibly.") @PreAuthorize("hasAnyAuthority('SYS_ADMIN', 'TENANT_ADMIN', 'CUSTOMER_USER')") @PostMapping("/entitiesQuery/find/aggHistory") public DeferredResult> findEntityDataAggHistoryByQuery( diff --git a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java index e4596d0f13..4ea9bc7529 100644 --- a/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java +++ b/application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java @@ -36,14 +36,12 @@ import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; -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.page.PageData; import org.thingsboard.server.common.data.query.AlarmCountQuery; import org.thingsboard.server.common.data.query.AlarmData; import org.thingsboard.server.common.data.query.AlarmDataQuery; -import org.thingsboard.server.common.data.query.ComparisonTsValue; import org.thingsboard.server.common.data.query.ComplexFilterPredicate; import org.thingsboard.server.common.data.query.DynamicValue; import org.thingsboard.server.common.data.query.EntityCountQuery; @@ -57,7 +55,6 @@ import org.thingsboard.server.common.data.query.FilterPredicateType; import org.thingsboard.server.common.data.query.KeyFilter; import org.thingsboard.server.common.data.query.KeyFilterPredicate; import org.thingsboard.server.common.data.query.SimpleKeyFilterPredicate; -import org.thingsboard.server.common.data.query.TsValue; import org.thingsboard.server.dao.alarm.AlarmService; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.entity.EntityService; @@ -68,14 +65,13 @@ import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.security.AccessValidator; import org.thingsboard.server.service.security.model.SecurityUser; +import org.thingsboard.server.service.subscription.AggregationQueryUtils; import org.thingsboard.server.service.subscription.ReadTsKvQueryInfo; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggHistoryCmd; -import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggKey; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; -import java.util.HashMap; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -358,19 +354,7 @@ public class DefaultEntityQueryService implements EntityQueryService { return response; } TenantId tenantId = securityUser.getTenantId(); - Map queries = new HashMap<>(); - for (AggKey key : cmd.getKeys()) { - if (key.getPreviousValueOnly() == null || !key.getPreviousValueOnly()) { - var q = new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), cmd.getEndTs(), - cmd.getEndTs() - cmd.getStartTs(), 1, key.getAgg()); - queries.put(q.getId(), new ReadTsKvQueryInfo(key, q, false)); - } - if (key.getPreviousStartTs() != null && key.getPreviousEndTs() != null && key.getPreviousEndTs() >= key.getPreviousStartTs()) { - var q = new BaseReadTsKvQuery(key.getKey(), key.getPreviousStartTs(), key.getPreviousEndTs(), - key.getPreviousEndTs() - key.getPreviousStartTs(), 1, key.getAgg()); - queries.put(q.getId(), new ReadTsKvQueryInfo(key, q, true)); - } - } + Map queries = AggregationQueryUtils.buildAggHistoryQueries(cmd.getKeys(), cmd.getStartTs(), cmd.getEndTs()); List queryList = queries.values().stream().map(ReadTsKvQueryInfo::getQuery).collect(Collectors.toList()); Map>> fetchResultMap = new LinkedHashMap<>(); entityDataList.forEach(entityData -> fetchResultMap.put(entityData, @@ -381,24 +365,10 @@ public class DefaultEntityQueryService implements EntityQueryService { fetchResultMap.forEach((entityData, future) -> { try { List queryResults = future.get(); - if (queryResults != null) { - for (ReadTsKvQueryResult queryResult : queryResults) { - ReadTsKvQueryInfo info = queries.get(queryResult.getQueryId()); - ComparisonTsValue comparisonTsValue = entityData.getAggLatest() - .computeIfAbsent(info.getKey().getId(), k -> new ComparisonTsValue()); - if (info.isPrevious()) { - comparisonTsValue.setPrevious(queryResult.toTsValue(info.getQuery())); - } else { - comparisonTsValue.setCurrent(queryResult.toTsValue(info.getQuery())); - } - } - } - cmd.getKeys().forEach(key -> entityData.getAggLatest() - .putIfAbsent(key.getId(), new ComparisonTsValue(TsValue.EMPTY, TsValue.EMPTY))); + AggregationQueryUtils.populateAggLatest(entityData, queryResults, queries, cmd.getKeys(), null); } catch (InterruptedException | ExecutionException e) { log.warn("[{}] Failed to fetch aggregated historical data", entityData.getEntityId(), e); - cmd.getKeys().forEach(key -> entityData.getAggLatest() - .putIfAbsent(key.getId(), new ComparisonTsValue(TsValue.EMPTY, TsValue.EMPTY))); + AggregationQueryUtils.populateAggLatest(entityData, null, queries, cmd.getKeys(), null); } }); response.setResult(pageData); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/AggregationQueryUtils.java b/application/src/main/java/org/thingsboard/server/service/subscription/AggregationQueryUtils.java new file mode 100644 index 0000000000..05d3f48fd5 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/subscription/AggregationQueryUtils.java @@ -0,0 +1,90 @@ +/** + * Copyright © 2016-2026 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.service.subscription; + +import org.thingsboard.server.common.data.kv.Aggregation; +import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; +import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; +import org.thingsboard.server.common.data.query.ComparisonTsValue; +import org.thingsboard.server.common.data.query.EntityData; +import org.thingsboard.server.common.data.query.TsValue; +import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggKey; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** + * Shared aggregated-history logic used by both the WebSocket {@code AggHistoryCmd} handler + * ({@link DefaultTbEntityDataSubscriptionService}) and the REST {@code /find/aggHistory} endpoint. + */ +public class AggregationQueryUtils { + + /** + * Builds one current (and, when requested, one previous) timeseries query per {@link AggKey}, keyed by query id. + * A key with {@code previousValueOnly} skips the current window; a previous window is added only when both bounds + * are present and {@code previousEndTs >= previousStartTs}. + */ + public static Map buildAggHistoryQueries(List keys, long startTs, long endTs) { + Map queries = new HashMap<>(); + for (AggKey key : keys) { + if (key.getPreviousValueOnly() == null || !key.getPreviousValueOnly()) { + var query = singleBucketQuery(key.getKey(), startTs, endTs, key.getAgg()); + queries.put(query.getId(), new ReadTsKvQueryInfo(key, query, false)); + } + if (key.getPreviousStartTs() != null && key.getPreviousEndTs() != null && key.getPreviousEndTs() >= key.getPreviousStartTs()) { + var query = singleBucketQuery(key.getKey(), key.getPreviousStartTs(), key.getPreviousEndTs(), key.getAgg()); + queries.put(query.getId(), new ReadTsKvQueryInfo(key, query, true)); + } + } + return queries; + } + + /** + * A single bucket spanning the whole {@code [startTs, endTs]} window (interval = the full window, limit = 1), + * yielding exactly one aggregated point. + */ + private static BaseReadTsKvQuery singleBucketQuery(String key, long startTs, long endTs, Aggregation agg) { + return new BaseReadTsKvQuery(key, startTs, endTs, endTs - startTs, 1, agg); + } + + /** + * Maps the per-entity query results into {@code entityData.aggLatest} (current/previous per key) and fills any key + * without data with {@link TsValue#EMPTY}. When {@code lastTsCollector} is non-null, the last entry ts of each + * current value is recorded into it (used by the WS path to set up follow-up subscriptions). + */ + public static void populateAggLatest(EntityData entityData, List queryResults, + Map queries, List keys, + Map lastTsCollector) { + if (queryResults != null) { + for (ReadTsKvQueryResult queryResult : queryResults) { + ReadTsKvQueryInfo queryInfo = queries.get(queryResult.getQueryId()); + ComparisonTsValue comparisonTsValue = entityData.getAggLatest() + .computeIfAbsent(queryInfo.getKey().getId(), agg -> new ComparisonTsValue()); + if (queryInfo.isPrevious()) { + comparisonTsValue.setPrevious(queryResult.toTsValue(queryInfo.getQuery())); + } else { + comparisonTsValue.setCurrent(queryResult.toTsValue(queryInfo.getQuery())); + if (lastTsCollector != null) { + lastTsCollector.put(queryInfo.getQuery().getKey(), queryResult.getLastEntryTs()); + } + } + } + } + keys.forEach(key -> entityData.getAggLatest().putIfAbsent(key.getId(), new ComparisonTsValue(TsValue.EMPTY, TsValue.EMPTY))); + } + +} 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 fc8c5c3be2..d1e1413715 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 @@ -41,7 +41,6 @@ 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; -import org.thingsboard.server.common.data.query.ComparisonTsValue; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityKey; @@ -296,17 +295,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } private ListenableFuture handleAggHistoryCmd(TbEntityDataSubCtx ctx, AggHistoryCmd cmd) { - ConcurrentMap queries = new ConcurrentHashMap<>(); - for (AggKey key : cmd.getKeys()) { - if (key.getPreviousValueOnly() == null || !key.getPreviousValueOnly()) { - var query = new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), cmd.getEndTs(), cmd.getEndTs() - cmd.getStartTs(), 1, key.getAgg()); - queries.put(query.getId(), new ReadTsKvQueryInfo(key, query, false)); - } - if (key.getPreviousStartTs() != null && key.getPreviousEndTs() != null && key.getPreviousEndTs() >= key.getPreviousStartTs()) { - var query = new BaseReadTsKvQuery(key.getKey(), key.getPreviousStartTs(), key.getPreviousEndTs(), key.getPreviousEndTs() - key.getPreviousStartTs(), 1, key.getAgg()); - queries.put(query.getId(), new ReadTsKvQueryInfo(key, query, true)); - } - } + Map queries = AggregationQueryUtils.buildAggHistoryQueries(cmd.getKeys(), cmd.getStartTs(), cmd.getEndTs()); return handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getEndTs(), false); } @@ -319,7 +308,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc return handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getStartTs() + cmd.getTimeWindow(), true); } - private ListenableFuture handleAggCmd(TbEntityDataSubCtx ctx, List keys, ConcurrentMap queries, + private ListenableFuture handleAggCmd(TbEntityDataSubCtx ctx, List keys, Map queries, long startTs, long endTs, boolean subscribe) { Map>> fetchResultMap = new HashMap<>(); List entityDataList = ctx.getData().getData(); @@ -335,22 +324,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc lastTsEntityMap.put(entityData, lastTsMap); List queryResults = future.get(); - if (queryResults != null) { - for (ReadTsKvQueryResult queryResult : queryResults) { - ReadTsKvQueryInfo queryInfo = queries.get(queryResult.getQueryId()); - ComparisonTsValue comparisonTsValue = entityData.getAggLatest().computeIfAbsent(queryInfo.getKey().getId(), agg -> new ComparisonTsValue()); - if (queryInfo.isPrevious()) { - comparisonTsValue.setPrevious(queryResult.toTsValue(queryInfo.getQuery())); - } else { - comparisonTsValue.setCurrent(queryResult.toTsValue(queryInfo.getQuery())); - lastTsMap.put(queryInfo.getQuery().getKey(), queryResult.getLastEntryTs()); - } - } - } - // Populate with empty values if no data found. - keys.forEach(key -> { - entityData.getAggLatest().putIfAbsent(key.getId(), new ComparisonTsValue(TsValue.EMPTY, TsValue.EMPTY)); - }); + AggregationQueryUtils.populateAggLatest(entityData, queryResults, queries, keys, lastTsMap); } catch (InterruptedException | ExecutionException e) { log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e); ctx.sendWsMsg(new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!"));