From f97d74d6be12bfbfafe1e6e48f26520c612fb745 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Wed, 24 Jun 2020 18:05:06 +0300 Subject: [PATCH] Improved WS API --- ...efaultTbEntityDataSubscriptionService.java | 157 ++++++++---------- .../telemetry/cmd/v2/EntityHistoryCmd.java | 2 +- .../service/telemetry/cmd/v2/GetTsCmd.java | 38 +++++ .../telemetry/cmd/v2/TimeSeriesCmd.java | 12 +- .../controller/BaseWebsocketApiTest.java | 10 +- 5 files changed, 124 insertions(+), 95 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/GetTsCmd.java 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 1772735f97..9cfd8cf80a 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 @@ -52,6 +52,7 @@ import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUnsubscribeCmd; import org.thingsboard.server.service.telemetry.cmd.v2.EntityDataUpdate; import org.thingsboard.server.service.telemetry.cmd.v2.EntityHistoryCmd; +import org.thingsboard.server.service.telemetry.cmd.v2.GetTsCmd; import org.thingsboard.server.service.telemetry.cmd.v2.LatestValueCmd; import org.thingsboard.server.service.telemetry.cmd.v2.TimeSeriesCmd; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; @@ -272,50 +273,86 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } } - private void handleTimeSeriesCmd(TbEntityDataSubCtx ctx, TimeSeriesCmd cmd) { - List keys = cmd.getKeys(); + private ListenableFuture handleTimeSeriesCmd(TbEntityDataSubCtx ctx, TimeSeriesCmd cmd) { log.debug("[{}][{}] Fetching time-series data for last {} ms for keys: ({})", ctx.getSessionId(), ctx.getCmdId(), cmd.getTimeWindow(), cmd.getKeys()); - long startTs = cmd.getStartTs(); - long endTs = cmd.getStartTs() + cmd.getTimeWindow(); - - Map>>> tsFutures = new HashMap<>(); - for (EntityData entityData : ctx.getData().getData()) { - List queries = keys.stream().map(key -> new BaseReadTsKvQuery(key, startTs, endTs, cmd.getInterval(), - getLimit(cmd.getLimit()), DefaultTelemetryWebSocketService.getAggregation(cmd.getAgg()))).collect(Collectors.toList()); - ListenableFuture> tsDataFutures = tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), queries); - tsFutures.put(entityData, Futures.transform(tsDataFutures, this::toTsValues, MoreExecutors.directExecutor())); + return handleGetTsCmd(ctx, cmd, true); + } + + + private ListenableFuture handleHistoryCmd(TbEntityDataSubCtx ctx, EntityHistoryCmd cmd) { + log.debug("[{}][{}] Fetching history data for start {} and end {} ms for keys: ({})", ctx.getSessionId(), ctx.getCmdId(), cmd.getStartTs(), cmd.getEndTs(), cmd.getKeys()); + return handleGetTsCmd(ctx, cmd, false); + } + + private ListenableFuture handleGetTsCmd(TbEntityDataSubCtx ctx, GetTsCmd cmd, boolean subscribe) { + List keys = cmd.getKeys(); + List finalTsKvQueryList; + List tsKvQueryList = cmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( + key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), getLimit(cmd.getLimit()), cmd.getAgg() + )).collect(Collectors.toList()); + if (cmd.isFetchLatestPreviousPoint()) { + finalTsKvQueryList = new ArrayList<>(tsKvQueryList); + tsKvQueryList.addAll(cmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( + key, cmd.getStartTs() - TimeUnit.DAYS.toMillis(365), cmd.getStartTs(), cmd.getInterval(), 1, cmd.getAgg() + )).collect(Collectors.toList())); + } else { + finalTsKvQueryList = tsKvQueryList; } - Futures.addCallback(Futures.allAsList(tsFutures.values()), new FutureCallback>>>() { - @Override - public void onSuccess(@Nullable List>> result) { - tsFutures.forEach((key, value) -> { - try { - value.get().forEach((k, v) -> key.getTimeseries().put(k, v.toArray(new TsValue[v.size()]))); - } catch (InterruptedException | ExecutionException e) { - log.warn("[{}][{}] Failed to lookup time-series data: {}:{}", ctx.getSessionId(), ctx.getCmdId(), key.getEntityId(), keys, e); + Map>> fetchResultMap = new HashMap<>(); + ctx.getData().getData().forEach(entityData -> fetchResultMap.put(entityData, + tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList))); + return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { + 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()))); } - }); - EntityDataUpdate update; - if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); - ctx.setInitialDataSent(true); - } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); + keyData.forEach((k, v) -> entityData.getTimeseries().put(k, v.toArray(new TsValue[v.size()]))); + if (cmd.isFetchLatestPreviousPoint()) { + 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); + wsService.sendWsMsg(ctx.getSessionId(), + new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!")); } - wsService.sendWsMsg(ctx.getSessionId(), update); - createSubscriptions(ctx, keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), false); - ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); + }); + EntityDataUpdate update; + if (!ctx.isInitialDataSent()) { + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + ctx.setInitialDataSent(true); + } else { + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); } - - @Override - public void onFailure(Throwable t) { - log.warn("[{}][{}] Failed to process websocket command: {}:{}", ctx.getSessionId(), ctx.getCmdId(), ctx.getQuery(), cmd, t); - wsService.sendWsMsg(ctx.getSessionId(), - new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to process websocket command!")); + wsService.sendWsMsg(ctx.getSessionId(), update); + if (subscribe) { + createSubscriptions(ctx, keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), false); } + ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); + return ctx; }, wsCallBackExecutor); } + private List getReadTsKvQueries(GetTsCmd cmd) { + List finalTsKvQueryList; + List queries = cmd.getKeys().stream().map(key -> new BaseReadTsKvQuery(key, cmd.getStartTs(), cmd.getEndTs(), cmd.getInterval(), + getLimit(cmd.getLimit()), cmd.getAgg())).collect(Collectors.toList()); + if (cmd.isFetchLatestPreviousPoint()) { + finalTsKvQueryList = new ArrayList<>(queries); + finalTsKvQueryList.addAll(cmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( + key, cmd.getStartTs() - TimeUnit.DAYS.toMillis(365), cmd.getStartTs(), cmd.getInterval(), 1, cmd.getAgg() + )).collect(Collectors.toList())); + } else { + finalTsKvQueryList = queries; + } + return finalTsKvQueryList; + } + private void handleLatestCmd(TbEntityDataSubCtx ctx, LatestValueCmd latestCmd) { log.trace("[{}][{}] Going to process latest command: {}", ctx.getSessionId(), ctx.getCmdId(), latestCmd); //Fetch the latest values for telemetry keys (in case they are not copied from NoSQL to SQL DB in hybrid mode. @@ -400,56 +437,6 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc return results; } - private ListenableFuture handleHistoryCmd(TbEntityDataSubCtx ctx, EntityHistoryCmd historyCmd) { - List finalTsKvQueryList; - List tsKvQueryList = historyCmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( - key, historyCmd.getStartTs(), historyCmd.getEndTs(), historyCmd.getInterval(), getLimit(historyCmd.getLimit()), historyCmd.getAgg() - )).collect(Collectors.toList()); - if (historyCmd.isFetchLatestPreviousPoint()) { - finalTsKvQueryList = new ArrayList<>(tsKvQueryList); - tsKvQueryList.addAll(historyCmd.getKeys().stream().map(key -> new BaseReadTsKvQuery( - key, historyCmd.getStartTs() - TimeUnit.DAYS.toMillis(365), historyCmd.getStartTs(), historyCmd.getInterval(), 1, historyCmd.getAgg() - )).collect(Collectors.toList())); - } else { - finalTsKvQueryList = tsKvQueryList; - } - Map>> fetchResultMap = new HashMap<>(); - ctx.getData().getData().forEach(entityData -> fetchResultMap.put(entityData, - tsService.findAll(ctx.getTenantId(), entityData.getEntityId(), finalTsKvQueryList))); - return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { - fetchResultMap.forEach((entityData, future) -> { - Map> keyData = new LinkedHashMap<>(); - historyCmd.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()))); - } - keyData.forEach((k, v) -> entityData.getTimeseries().put(k, v.toArray(new TsValue[v.size()]))); - if (historyCmd.isFetchLatestPreviousPoint()) { - 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); - wsService.sendWsMsg(ctx.getSessionId(), - new EntityDataUpdate(ctx.getCmdId(), SubscriptionErrorCode.INTERNAL_ERROR.getCode(), "Failed to fetch historical data!")); - } - }); - EntityDataUpdate update; - if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); - ctx.setInitialDataSent(true); - } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); - } - wsService.sendWsMsg(ctx.getSessionId(), update); - ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); - return ctx; - }, wsCallBackExecutor); - } - @Override public void cancelSubscription(String sessionId, EntityDataUnsubscribeCmd cmd) { cleanupAndCancel(getSubCtx(sessionId, cmd.getCmdId())); diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityHistoryCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityHistoryCmd.java index 8e3aafac25..b78cef98d7 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityHistoryCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityHistoryCmd.java @@ -21,7 +21,7 @@ import org.thingsboard.server.common.data.kv.Aggregation; import java.util.List; @Data -public class EntityHistoryCmd { +public class EntityHistoryCmd implements GetTsCmd { private List keys; private long startTs; diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/GetTsCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/GetTsCmd.java new file mode 100644 index 0000000000..25cc1a8e0f --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/GetTsCmd.java @@ -0,0 +1,38 @@ +/** + * Copyright © 2016-2020 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.telemetry.cmd.v2; + +import org.thingsboard.server.common.data.kv.Aggregation; + +import java.util.List; + +public interface GetTsCmd { + + long getStartTs(); + + long getEndTs(); + + List getKeys(); + + long getInterval(); + + int getLimit(); + + Aggregation getAgg(); + + boolean isFetchLatestPreviousPoint(); + +} diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/TimeSeriesCmd.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/TimeSeriesCmd.java index 336c40ac97..8762ee3438 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/TimeSeriesCmd.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/TimeSeriesCmd.java @@ -15,18 +15,26 @@ */ package org.thingsboard.server.service.telemetry.cmd.v2; +import com.fasterxml.jackson.annotation.JsonIgnore; import lombok.Data; +import org.thingsboard.server.common.data.kv.Aggregation; import java.util.List; @Data -public class TimeSeriesCmd { +public class TimeSeriesCmd implements GetTsCmd { private List keys; private long startTs; private long timeWindow; private long interval; private int limit; - private String agg; + private Aggregation agg; + private boolean fetchLatestPreviousPoint; + @JsonIgnore + @Override + public long getEndTs() { + return startTs + timeWindow; + } } diff --git a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java index 20c2bac13d..0a2e798b08 100644 --- a/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java +++ b/application/src/test/java/org/thingsboard/server/controller/BaseWebsocketApiTest.java @@ -194,7 +194,7 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { TimeSeriesCmd tsCmd = new TimeSeriesCmd(); tsCmd.setKeys(Arrays.asList("temperature")); - tsCmd.setAgg(Aggregation.NONE.name()); + tsCmd.setAgg(Aggregation.NONE); tsCmd.setLimit(1000); tsCmd.setStartTs(now - TimeUnit.HOURS.toMillis(1)); tsCmd.setTimeWindow(TimeUnit.HOURS.toMillis(1)); @@ -561,16 +561,12 @@ public class BaseWebsocketApiTest extends AbstractWebsocketTest { wsClient.registerWaitForUpdate(); AttributeKvEntry dataPoint1 = new BaseAttributeKvEntry(now - TimeUnit.MINUTES.toMillis(1), new LongDataEntry("serverAttributeKey", 42L)); List tsData = Arrays.asList(dataPoint1); - sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); - Thread.sleep(100); - cmd = new EntityDataCmd(1, edq, null, latestCmd, null); - wrapper = new TelemetryPluginCmdsWrapper(); - wrapper.setEntityDataCmds(Collections.singletonList(cmd)); + sendAttributes(device, TbAttributeSubscriptionScope.SERVER_SCOPE, tsData); msg = wsClient.waitForUpdate(); - + Assert.assertNotNull(msg); update = mapper.readValue(msg, EntityDataUpdate.class); Assert.assertEquals(1, update.getCmdId()); List eData = update.getUpdate();