From 86a7e9cfa2de29022faa84b212ac02c45022470c Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 10 Jul 2020 11:57:52 +0300 Subject: [PATCH] Alarm API --- ...efaultTbEntityDataSubscriptionService.java | 44 +++++---- .../SubscriptionServiceStatistics.java | 31 +++++++ .../subscription/TbAbstractDataSubCtx.java | 45 +++++++-- .../subscription/TbAlarmDataSubCtx.java | 87 +++++++++++++----- .../subscription/TbEntityDataSubCtx.java | 91 ++++++++----------- .../telemetry/cmd/v2/AlarmDataUpdate.java | 9 +- 6 files changed, 209 insertions(+), 98 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionServiceStatistics.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 4c75486383..d1b31b70ca 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 @@ -134,10 +134,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private ExecutorService wsCallBackExecutor; private boolean tsInSqlDB; private String serviceId; - private AtomicInteger regularQueryInvocationCnt = new AtomicInteger(); - private AtomicInteger dynamicQueryInvocationCnt = new AtomicInteger(); - private AtomicLong regularQueryTimeSpent = new AtomicLong(); - private AtomicLong dynamicQueryTimeSpent = new AtomicLong(); + private SubscriptionServiceStatistics stats = new SubscriptionServiceStatistics(); @PostConstruct public void initExecutor() { @@ -196,8 +193,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc long start = System.currentTimeMillis(); ctx.fetchData(); long end = System.currentTimeMillis(); - regularQueryInvocationCnt.incrementAndGet(); - regularQueryTimeSpent.addAndGet(end - start); + stats.getRegularQueryInvocationCnt().incrementAndGet(); + stats.getRegularQueryTimeSpent().addAndGet(end - start); ctx.cancelTasks(); if (ctx.getQuery().getPageLink().isDynamic()) { //TODO: validate number of dynamic page links against rate limits. Ignore dynamic flag if limit is reached. @@ -245,12 +242,18 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } ctx.setAndResolveQuery(cmd.getQuery()); AlarmDataQuery adq = ctx.getQuery(); + long start = System.currentTimeMillis(); ctx.fetchData(); + long end = System.currentTimeMillis(); + stats.getRegularQueryInvocationCnt().incrementAndGet(); + stats.getRegularQueryTimeSpent().addAndGet(end - start); List entities = ctx.getEntitiesData(); ctx.cancelTasks(); ctx.clearEntitySubscriptions(); if (entities.isEmpty()) { - AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), new PageData<>(Collections.emptyList(), 1, 0, false), null); + AlarmDataUpdate update = new AlarmDataUpdate(cmd.getCmdId(), + new PageData<>(Collections.emptyList(), 1, 0, false), + null, false); wsService.sendWsMsg(ctx.getSessionId(), update); } else { ctx.fetchAlarms(); @@ -269,8 +272,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc long start = System.currentTimeMillis(); finalCtx.update(); long end = System.currentTimeMillis(); - dynamicQueryInvocationCnt.incrementAndGet(); - dynamicQueryTimeSpent.addAndGet(end - start); + stats.getDynamicQueryInvocationCnt().incrementAndGet(); + stats.getDynamicQueryTimeSpent().addAndGet(end - start); } catch (Exception e) { log.warn("[{}][{}] Failed to refresh query", finalCtx.getSessionId(), finalCtx.getCmdId(), e); } @@ -278,20 +281,26 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc @Scheduled(fixedDelayString = "${server.ws.dynamic_page_link.stats:10000}") public void printStats() { - int regularQueryInvocationCntValue = regularQueryInvocationCnt.getAndSet(0); - long regularQueryInvocationTimeValue = regularQueryTimeSpent.getAndSet(0); - int dynamicQueryInvocationCntValue = dynamicQueryInvocationCnt.getAndSet(0); - long dynamicQueryInvocationTimeValue = dynamicQueryTimeSpent.getAndSet(0); + int alarmQueryInvocationCntValue = stats.getAlarmQueryInvocationCnt().getAndSet(0); + long alarmQueryInvocationTimeValue = stats.getAlarmQueryTimeSpent().getAndSet(0); + int regularQueryInvocationCntValue = stats.getRegularQueryInvocationCnt().getAndSet(0); + long regularQueryInvocationTimeValue = stats.getRegularQueryTimeSpent().getAndSet(0); + int dynamicQueryInvocationCntValue = stats.getDynamicQueryInvocationCnt().getAndSet(0); + long dynamicQueryInvocationTimeValue = stats.getDynamicQueryTimeSpent().getAndSet(0); long dynamicQueryCnt = subscriptionsBySessionId.values().stream().map(Map::values).count(); if (regularQueryInvocationCntValue > 0 || dynamicQueryInvocationCntValue > 0 || dynamicQueryCnt > 0) { - log.info("Stats: regularQueryInvocationCnt = [{}], regularQueryInvocationTime = [{}], dynamicQueryCnt = [{}] dynamicQueryInvocationCnt = [{}], dynamicQueryInvocationTime = [{}]", - regularQueryInvocationCntValue, regularQueryInvocationTimeValue, dynamicQueryCnt, dynamicQueryInvocationCntValue, dynamicQueryInvocationTimeValue); + log.info("Stats: regularQueryInvocationCnt = [{}], regularQueryInvocationTime = [{}], " + + "dynamicQueryCnt = [{}] dynamicQueryInvocationCnt = [{}], dynamicQueryInvocationTime = [{}], " + + "alarmQueryInvocationCnt = [{}], alarmQueryInvocationTime = [{}]", + regularQueryInvocationCntValue, regularQueryInvocationTimeValue, + dynamicQueryCnt, dynamicQueryInvocationCntValue, dynamicQueryInvocationTimeValue, + alarmQueryInvocationCntValue, alarmQueryInvocationTimeValue); } } private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); - TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmd.getCmdId()); + TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmd.getCmdId()); ctx.setAndResolveQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; @@ -299,7 +308,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { Map sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); - TbAlarmDataSubCtx ctx = new TbAlarmDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, alarmService, sessionRef, cmd.getCmdId(), maxEntitiesPerAlarmSubscription); + TbAlarmDataSubCtx ctx = new TbAlarmDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, alarmService, sessionRef, cmd.getCmdId(), maxEntitiesPerAlarmSubscription); ctx.setAndResolveQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; @@ -454,6 +463,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc if (ctx != null) { ctx.cancelTasks(); ctx.clearEntitySubscriptions(); + ctx.clearDynamicValueSubscriptions(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionServiceStatistics.java b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionServiceStatistics.java new file mode 100644 index 0000000000..80f3759c80 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionServiceStatistics.java @@ -0,0 +1,31 @@ +/** + * 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.subscription; + +import lombok.Data; + +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; + +@Data +public class SubscriptionServiceStatistics { + private AtomicInteger alarmQueryInvocationCnt = new AtomicInteger(); + private AtomicInteger regularQueryInvocationCnt = new AtomicInteger(); + private AtomicInteger dynamicQueryInvocationCnt = new AtomicInteger(); + private AtomicLong alarmQueryTimeSpent = new AtomicLong(); + private AtomicLong regularQueryTimeSpent = new AtomicLong(); + private AtomicLong dynamicQueryTimeSpent = new AtomicLong(); +} 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 f87220469b..6e1bb13557 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 @@ -5,7 +5,7 @@ * 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 + * 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, @@ -50,6 +50,7 @@ import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.HashMap; import java.util.HashSet; import java.util.List; @@ -58,12 +59,15 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.concurrent.ScheduledFuture; +import java.util.function.Function; +import java.util.stream.Collectors; @Slf4j @Data public abstract class TbAbstractDataSubCtx> { protected final String serviceId; + protected final SubscriptionServiceStatistics stats; protected final TelemetryWebSocketService wsService; protected final EntityService entityService; protected final TbLocalSubscriptionService localSubscriptionService; @@ -84,12 +88,14 @@ public abstract class TbAbstractDataSubCtx(); @@ -151,8 +157,6 @@ public abstract class TbAbstractDataSubCtx newData = findEntityData(); + long end = System.currentTimeMillis(); + stats.getRegularQueryInvocationCnt().incrementAndGet(); + stats.getRegularQueryTimeSpent().addAndGet(end - start); + Map oldDataMap; + if (data != null && !data.getData().isEmpty()) { + oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); + } else { + oldDataMap = Collections.emptyMap(); + } + Map newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); + if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { + log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); + } else { + this.data = newData; + doUpdate(newDataMap); + } + } + + protected abstract void doUpdate(Map newDataMap); protected abstract EntityDataQuery buildEntityDataQuery(); @@ -312,6 +337,15 @@ public abstract class TbAbstractDataSubCtx task) { this.refreshTask = task; } @@ -323,7 +357,6 @@ public abstract class TbAbstractDataSubCtx keys, boolean resultToLatestValues) { Map> keysByType = getEntityKeyByTypeMap(keys); for (EntityData entityData : data.getData()) { diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java index 66fad14770..cb425daf79 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java @@ -43,12 +43,15 @@ import org.thingsboard.server.service.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.telemetry.sub.AlarmSubscriptionUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; +import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.HashMap; +import java.util.HashSet; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.function.Function; import java.util.stream.Collectors; @@ -74,9 +77,9 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, TbLocalSubscriptionService localSubscriptionService, - AttributesService attributesService, AlarmService alarmService, + AttributesService attributesService, SubscriptionServiceStatistics stats, AlarmService alarmService, TelemetryWebSocketSessionRef sessionRef, int cmdId, int maxEntitiesPerAlarmSubscription) { - super(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmdId); + super(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmdId); this.maxEntitiesPerAlarmSubscription = maxEntitiesPerAlarmSubscription; this.alarmService = alarmService; this.entitiesMap = new LinkedHashMap<>(); @@ -84,10 +87,14 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } public void fetchAlarms() { + long start = System.currentTimeMillis(); PageData alarms = alarmService.findAlarmDataByQueryForEntities(getTenantId(), getCustomerId(), query, getOrderedEntityIds()); + long end = System.currentTimeMillis(); + stats.getAlarmQueryInvocationCnt().incrementAndGet(); + stats.getAlarmQueryTimeSpent().addAndGet(end - start); alarms = setAndMergeAlarmsData(alarms); - AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null); + AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null, tooManyEntities); wsService.sendWsMsg(getSessionId(), update); } @@ -130,23 +137,27 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { AlarmDataPageLink pageLink = query.getPageLink(); long startTs = System.currentTimeMillis() - pageLink.getTimeWindow(); for (EntityData entityData : entitiesMap.values()) { - int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); - subToEntityIdMap.put(subIdx, entityData.getEntityId()); - log.trace("[{}][{}][{}] Creating alarms subscription for [{}] with query: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), pageLink); - TbAlarmsSubscription subscription = TbAlarmsSubscription.builder() - .type(TbSubscriptionType.ALARMS) - .serviceId(serviceId) - .sessionId(sessionRef.getSessionId()) - .subscriptionId(subIdx) - .tenantId(sessionRef.getSecurityCtx().getTenantId()) - .entityId(entityData.getEntityId()) - .updateConsumer(this::sendWsMsg) - .ts(startTs) - .build(); - localSubscriptionService.addSubscription(subscription); + createAlarmSubscriptionForEntity(pageLink, startTs, entityData); } } + private void createAlarmSubscriptionForEntity(AlarmDataPageLink pageLink, long startTs, EntityData entityData) { + int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); + subToEntityIdMap.put(subIdx, entityData.getEntityId()); + log.trace("[{}][{}][{}] Creating alarms subscription for [{}] with query: {}", serviceId, cmdId, subIdx, entityData.getEntityId(), pageLink); + TbAlarmsSubscription subscription = TbAlarmsSubscription.builder() + .type(TbSubscriptionType.ALARMS) + .serviceId(serviceId) + .sessionId(sessionRef.getSessionId()) + .subscriptionId(subIdx) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityData.getEntityId()) + .updateConsumer(this::sendWsMsg) + .ts(startTs) + .build(); + localSubscriptionService.addSubscription(subscription); + } + @Override void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) { EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); @@ -163,7 +174,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { alarm.getLatest().computeIfAbsent(keyType, tmp -> new HashMap<>()).putAll(latestUpdate); return alarm; }).collect(Collectors.toList()); - wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update)); + wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update, tooManyEntities)); } else { log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); } @@ -186,7 +197,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId()); updated.getLatest().putAll(current.getLatest()); alarmsMap.put(alarmId, updated); - wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated))); + wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, Collections.singletonList(updated), tooManyEntities)); } else { fetchAlarms(); } @@ -241,8 +252,42 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } @Override - protected void update() { - + protected synchronized void doUpdate(Map newDataMap) { + entitiesMap.clear(); + tooManyEntities = data.hasNext(); + for (EntityData entityData : data.getData()) { + entitiesMap.put(entityData.getEntityId(), entityData); + } + fetchAlarms(); + List subIdsToCancel = new ArrayList<>(); + List subsToAdd = new ArrayList<>(); + Set currentSubs = new HashSet<>(); + subToEntityIdMap.forEach((subId, entityId) -> { + if (!newDataMap.containsKey(entityId)) { + subIdsToCancel.add(subId); + } else { + currentSubs.add(entityId); + } + }); + log.trace("[{}][{}] Subscriptions that are invalid: {}", sessionRef.getSessionId(), cmdId, subIdsToCancel); + subIdsToCancel.forEach(subToEntityIdMap::remove); + List newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList()); + if (!newSubsList.isEmpty()) { + List keys = query.getLatestValues(); + if (keys != null && !keys.isEmpty()) { + Map> keysByType = getEntityKeyByTypeMap(keys); + newSubsList.forEach( + entity -> { + log.trace("[{}][{}] Found new subscription for entity: {}", sessionRef.getSessionId(), cmdId, entity.getEntityId()); + subsToAdd.addAll(addSubscriptions(entity, keysByType, true)); + } + ); + } + long startTs = System.currentTimeMillis() - query.getPageLink().getTimeWindow(); + newSubsList.forEach(entity -> createAlarmSubscriptionForEntity(query.getPageLink(), startTs, entity)); + } + subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); + subsToAdd.forEach(localSubscriptionService::addSubscription); } @Override diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java index e600337cb9..0e58c07994 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java @@ -66,8 +66,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService, - TelemetryWebSocketSessionRef sessionRef, int cmdId) { - super(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmdId); + SubscriptionServiceStatistics stats, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + super(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmdId); } @Override @@ -171,58 +171,45 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); } - public synchronized void update() { - PageData newData = findEntityData(); - Map oldDataMap; - if (data != null && !data.getData().isEmpty()) { - oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); - } else { - oldDataMap = Collections.emptyMap(); - } - Map newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); - if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { - log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); - } else { - this.data = newData; - List subIdsToCancel = new ArrayList<>(); - List subsToAdd = new ArrayList<>(); - Set currentSubs = new HashSet<>(); - subToEntityIdMap.forEach((subId, entityId) -> { - if (!newDataMap.containsKey(entityId)) { - subIdsToCancel.add(subId); - } else { - currentSubs.add(entityId); - } - }); - log.trace("[{}][{}] Subscriptions that are invalid: {}", sessionRef.getSessionId(), cmdId, subIdsToCancel); - subIdsToCancel.forEach(subToEntityIdMap::remove); - List newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList()); - if (!newSubsList.isEmpty()) { - boolean resultToLatestValues; - List keys = null; - if (curTsCmd != null) { - resultToLatestValues = false; - keys = curTsCmd.getKeys().stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()); - } else if (latestValueCmd != null) { - resultToLatestValues = true; - keys = latestValueCmd.getKeys(); - } else { - resultToLatestValues = true; - } - if (keys != null && !keys.isEmpty()) { - Map> keysByType = getEntityKeyByTypeMap(keys); - newSubsList.forEach( - entity -> { - log.trace("[{}][{}] Found new subscription for entity: {}", sessionRef.getSessionId(), cmdId, entity.getEntityId()); - subsToAdd.addAll(addSubscriptions(entity, keysByType, resultToLatestValues)); - } - ); - } + public synchronized void doUpdate(Map newDataMap) { + List subIdsToCancel = new ArrayList<>(); + List subsToAdd = new ArrayList<>(); + Set currentSubs = new HashSet<>(); + subToEntityIdMap.forEach((subId, entityId) -> { + if (!newDataMap.containsKey(entityId)) { + subIdsToCancel.add(subId); + } else { + currentSubs.add(entityId); + } + }); + log.trace("[{}][{}] Subscriptions that are invalid: {}", sessionRef.getSessionId(), cmdId, subIdsToCancel); + subIdsToCancel.forEach(subToEntityIdMap::remove); + List newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList()); + if (!newSubsList.isEmpty()) { + boolean resultToLatestValues; + List keys = null; + if (curTsCmd != null) { + resultToLatestValues = false; + keys = curTsCmd.getKeys().stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()); + } else if (latestValueCmd != null) { + resultToLatestValues = true; + keys = latestValueCmd.getKeys(); + } else { + resultToLatestValues = true; + } + if (keys != null && !keys.isEmpty()) { + Map> keysByType = getEntityKeyByTypeMap(keys); + newSubsList.forEach( + entity -> { + log.trace("[{}][{}] Found new subscription for entity: {}", sessionRef.getSessionId(), cmdId, entity.getEntityId()); + subsToAdd.addAll(addSubscriptions(entity, keysByType, resultToLatestValues)); + } + ); } - wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null)); - subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); - subsToAdd.forEach(localSubscriptionService::addSubscription); } + wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null)); + subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); + subsToAdd.forEach(localSubscriptionService::addSubscription); } public void setCurrentCmd(EntityDataCmd cmd) { diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java index fd2a52dc02..e928068d35 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java @@ -27,8 +27,11 @@ import java.util.List; public class AlarmDataUpdate extends DataUpdate { - public AlarmDataUpdate(int cmdId, PageData data, List update) { + private boolean tooManyEntities; + + public AlarmDataUpdate(int cmdId, PageData data, List update, boolean tooManyEntities) { super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); + this.tooManyEntities = tooManyEntities; } public AlarmDataUpdate(int cmdId, int errorCode, String errorMsg) { @@ -45,7 +48,9 @@ public class AlarmDataUpdate extends DataUpdate { @JsonProperty("data") PageData data, @JsonProperty("update") List update, @JsonProperty("errorCode") int errorCode, - @JsonProperty("errorMsg") String errorMsg) { + @JsonProperty("errorMsg") String errorMsg, + @JsonProperty("tooManyEntities") boolean tooManyEntities) { super(cmdId, data, update, errorCode, errorMsg); + this.tooManyEntities = tooManyEntities; } }