Browse Source

Share aggregated-history logic between WebSocket and REST entity-data paths

Extract the query-building and aggLatest-population blocks into
AggregationQueryUtils; the WebSocket AggHistoryCmd handler and the REST
/find/aggHistory endpoint now call the shared methods. Document the
per-entity fan-out in the endpoint API notes.
pull/15429/head
Oleksandra Matviienko 2 months ago
parent
commit
2696dec738
  1. 3
      application/src/main/java/org/thingsboard/server/controller/EntityQueryController.java
  2. 38
      application/src/main/java/org/thingsboard/server/service/query/DefaultEntityQueryService.java
  3. 90
      application/src/main/java/org/thingsboard/server/service/subscription/AggregationQueryUtils.java
  4. 32
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

3
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<PageData<EntityData>> findEntityDataAggHistoryByQuery(

38
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<Integer, ReadTsKvQueryInfo> 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<Integer, ReadTsKvQueryInfo> queries = AggregationQueryUtils.buildAggHistoryQueries(cmd.getKeys(), cmd.getStartTs(), cmd.getEndTs());
List<ReadTsKvQuery> queryList = queries.values().stream().map(ReadTsKvQueryInfo::getQuery).collect(Collectors.toList());
Map<EntityData, ListenableFuture<List<ReadTsKvQueryResult>>> fetchResultMap = new LinkedHashMap<>();
entityDataList.forEach(entityData -> fetchResultMap.put(entityData,
@ -381,24 +365,10 @@ public class DefaultEntityQueryService implements EntityQueryService {
fetchResultMap.forEach((entityData, future) -> {
try {
List<ReadTsKvQueryResult> 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);

90
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<Integer, ReadTsKvQueryInfo> buildAggHistoryQueries(List<AggKey> keys, long startTs, long endTs) {
Map<Integer, ReadTsKvQueryInfo> 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<ReadTsKvQueryResult> queryResults,
Map<Integer, ReadTsKvQueryInfo> queries, List<AggKey> keys,
Map<String, Long> 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)));
}
}

32
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<TbEntityDataSubCtx> handleAggHistoryCmd(TbEntityDataSubCtx ctx, AggHistoryCmd cmd) {
ConcurrentMap<Integer, ReadTsKvQueryInfo> 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<Integer, ReadTsKvQueryInfo> 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<TbEntityDataSubCtx> handleAggCmd(TbEntityDataSubCtx ctx, List<AggKey> keys, ConcurrentMap<Integer, ReadTsKvQueryInfo> queries,
private ListenableFuture<TbEntityDataSubCtx> handleAggCmd(TbEntityDataSubCtx ctx, List<AggKey> keys, Map<Integer, ReadTsKvQueryInfo> queries,
long startTs, long endTs, boolean subscribe) {
Map<EntityData, ListenableFuture<List<ReadTsKvQueryResult>>> fetchResultMap = new HashMap<>();
List<EntityData> entityDataList = ctx.getData().getData();
@ -335,22 +324,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
lastTsEntityMap.put(entityData, lastTsMap);
List<ReadTsKvQueryResult> 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!"));

Loading…
Cancel
Save