From 9689e9c04144c4ed0df1325f6a57eadecdc93a41 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 13 Sep 2022 12:39:07 +0300 Subject: [PATCH] Fix bugs. Remove floating window support --- ...efaultTbEntityDataSubscriptionService.java | 71 +++++++++++-------- .../telemetry/cmd/v2/AggTimeSeriesCmd.java | 1 - .../server/common/data/query/EntityData.java | 3 +- .../server/common/data/query/EntityKey.java | 2 + .../dao/sql/query/EntityDataAdapter.java | 2 +- 5 files changed, 45 insertions(+), 34 deletions(-) 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 2b30b7422b..998938b5e2 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 @@ -70,6 +70,7 @@ import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; +import java.util.HashSet; import java.util.LinkedHashSet; import java.util.List; import java.util.Map; @@ -209,20 +210,20 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } if (cmd.getAggHistoryCmd() != null) { - handleAggHistoryCmd(session, ctx, cmd.getAggHistoryCmd()); + handleAggHistoryCmd(ctx, cmd.getAggHistoryCmd(), cmd.getLatestCmd()); } else if (cmd.getAggTsCmd() != null) { - handleAggTsCmd(session, ctx, cmd.getAggTsCmd()); + handleAggTsCmd(ctx, cmd.getAggTsCmd(), cmd.getLatestCmd()); } else if (cmd.hasRegularCmds()) { - handleRegularCommands(session, ctx, cmd); + handleRegularCommands(ctx, cmd); } else { checkAndSendInitialData(ctx); } } - private void handleRegularCommands(TelemetryWebSocketSessionRef session, TbEntityDataSubCtx ctx, EntityDataCmd cmd) { + private void handleRegularCommands(TbEntityDataSubCtx ctx, EntityDataCmd cmd) { ListenableFuture historyFuture; if (cmd.getHistoryCmd() != null) { - log.trace("[{}][{}] Going to process history command: {}", session.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); + log.trace("[{}][{}] Going to process history command: {}", ctx.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); try { historyFuture = handleHistoryCmd(ctx, cmd.getHistoryCmd()); } catch (RuntimeException e) { @@ -253,7 +254,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Override public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to process command", session.getSessionId(), cmd.getCmdId()); + log.warn("[{}][{}] Failed to process command", ctx.getSessionId(), cmd.getCmdId()); } }, wsCallBackExecutor); } @@ -266,34 +267,30 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - private void handleAggHistoryCmd(TelemetryWebSocketSessionRef session, TbEntityDataSubCtx ctx, AggHistoryCmd cmd) { + private void handleAggHistoryCmd(TbEntityDataSubCtx ctx, AggHistoryCmd cmd, LatestValueCmd latestCmd) { var keys = cmd.getKeys(); long interval = cmd.getEndTs() - cmd.getStartTs(); List queries = keys.stream().map(key -> new BaseReadTsKvQuery( key.getKey(), cmd.getStartTs(), cmd.getEndTs(), interval, 1, key.getAgg() )).distinct().collect(Collectors.toList()); - handleAggCmd(session, ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getEndTs(), false, false); + handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getEndTs(), latestCmd, false); } - private void handleAggTsCmd(TelemetryWebSocketSessionRef session, TbEntityDataSubCtx ctx, AggTimeSeriesCmd cmd) { + private void handleAggTsCmd(TbEntityDataSubCtx ctx, AggTimeSeriesCmd cmd, LatestValueCmd latestCmd) { long endTs = cmd.getStartTs() + cmd.getTimeWindow(); - List queries = cmd.getKeys().stream().map(key -> { - if (cmd.isFloating()) { - return new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), endTs, cmd.getTimeWindow(), getLimit(maxDatapointLimit), Aggregation.NONE); - } else { - return new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), endTs, cmd.getTimeWindow(), 1, key.getAgg()); - } - }).distinct().collect(Collectors.toList()); - handleAggCmd(session, ctx, cmd.getKeys(), queries, cmd.getStartTs(), endTs, cmd.isFloating(), true); + List queries = cmd.getKeys().stream() + .map(key -> new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), endTs, cmd.getTimeWindow(), 1, key.getAgg())) + .distinct().collect(Collectors.toList()); + handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), endTs, latestCmd, true); } - private void handleAggCmd(TelemetryWebSocketSessionRef session, TbEntityDataSubCtx ctx, List keys, List queries, - long startTs, long endTs, boolean floating, boolean subscribe) { + private void handleAggCmd(TbEntityDataSubCtx ctx, List keys, List queries, + long startTs, long endTs, LatestValueCmd latestCmd, boolean subscribe) { Map>> fetchResultMap = new HashMap<>(); List entityDataList = ctx.getData().getData(); entityDataList.forEach(entityData -> fetchResultMap.put(entityData, tsService.findAllByQueries(ctx.getTenantId(), entityData.getEntityId(), queries))); - Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { + var mainFuture = 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) -> { @@ -304,21 +301,13 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc List queryResults = future.get(); if (queryResults != null) { for (ReadTsKvQueryResult queryResult : queryResults) { - if (floating) { - entityData.getAggFloating().put(queryResult.getKey(), queryResult.toTsValues()); - } else { - entityData.getAggLatest().computeIfAbsent(queryResult.getAgg(), agg -> new HashMap<>()).put(queryResult.getKey(), queryResult.toTsValue()); - } + entityData.getAggLatest().computeIfAbsent(queryResult.getAgg(), agg -> new HashMap<>()).put(queryResult.getKey(), queryResult.toTsValue()); lastTsMap.put(queryResult.getKey(), queryResult.getLastEntryTs()); } } // Populate with empty values if no data found. keys.forEach(key -> { - if (floating) { - entityData.getAggFloating().putIfAbsent(key.getKey(), new TsValue[]{TsValue.EMPTY}); - } else { - entityData.getAggLatest().computeIfAbsent(key.getAgg(), agg -> new HashMap<>()).putIfAbsent(key.getKey(), TsValue.EMPTY); - } + entityData.getAggLatest().computeIfAbsent(key.getAgg(), agg -> new HashMap<>()).putIfAbsent(key.getKey(), TsValue.EMPTY); }); } catch (InterruptedException | ExecutionException e) { log.warn("[{}][{}][{}] Failed to fetch historical data", ctx.getSessionId(), ctx.getCmdId(), entityData.getEntityId(), e); @@ -344,6 +333,28 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } return ctx; }, wsCallBackExecutor); + + Futures.addCallback(mainFuture, new FutureCallback<>() { + + @Override + public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { + if (latestCmd != null) { + if (subscribe) { + var commandEntityKeys = new HashSet<>(latestCmd.getKeys()); + var alreadySubscribedKeys = keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key.getKey())).collect(Collectors.toList()); + if (commandEntityKeys.removeAll(alreadySubscribedKeys)) { + latestCmd.setKeys(new ArrayList<>(commandEntityKeys)); + } + } + handleLatestCmd(ctx, latestCmd); + } + } + + @Override + public void onFailure(Throwable t) { + log.warn("[{}][{}] Failed to process command", ctx.getSessionId(), ctx.getCmdId()); + } + }, wsCallBackExecutor); } private void handleWsCmdRuntimeException(String sessionId, RuntimeException e, EntityDataCmd cmd) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AggTimeSeriesCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AggTimeSeriesCmd.java index f5e4192af4..1f4d88d2c3 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AggTimeSeriesCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AggTimeSeriesCmd.java @@ -25,6 +25,5 @@ public class AggTimeSeriesCmd { private List keys; private long startTs; private long timeWindow; - private boolean floating; } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java index 26bbc94290..f1445faaab 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityData.java @@ -31,9 +31,8 @@ public class EntityData { private final Map> latest; private final Map timeseries; private final Map> aggLatest; - private final Map aggFloating; public EntityData(EntityId entityId, Map> latest, Map timeseries) { - this(entityId, latest, timeseries, null, null); + this(entityId, latest, timeseries, null); } } diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java index 5b4b5cd257..6aacc90060 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/query/EntityKey.java @@ -23,6 +23,8 @@ import java.io.Serializable; @ApiModel @Data public class EntityKey implements Serializable { + private static final long serialVersionUID = -6421575477523085543L; + private final EntityKeyType type; private final String key; } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityDataAdapter.java b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityDataAdapter.java index 561be612c7..e2d691f32b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityDataAdapter.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/query/EntityDataAdapter.java @@ -55,7 +55,7 @@ public class EntityDataAdapter { EntityId entityId = EntityIdFactory.getByTypeAndUuid(entityType, id); Map> latest = new HashMap<>(); //Maybe avoid empty hashmaps? - EntityData entityData = new EntityData(entityId, latest, new HashMap<>(), new HashMap<>(), new HashMap<>()); + EntityData entityData = new EntityData(entityId, latest, new HashMap<>(), new HashMap<>()); for (EntityKeyMapping mapping : selectionMapping) { if (!mapping.isIgnore()) { EntityKey entityKey = mapping.getEntityKey();