Browse Source

Alarm API

pull/3089/head
Andrii Shvaika 6 years ago
parent
commit
86a7e9cfa2
  1. 44
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 31
      application/src/main/java/org/thingsboard/server/service/subscription/SubscriptionServiceStatistics.java
  3. 45
      application/src/main/java/org/thingsboard/server/service/subscription/TbAbstractDataSubCtx.java
  4. 87
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmDataSubCtx.java
  5. 91
      application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java
  6. 9
      application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/AlarmDataUpdate.java

44
application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

@ -134,10 +134,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
private ExecutorService wsCallBackExecutor; private ExecutorService wsCallBackExecutor;
private boolean tsInSqlDB; private boolean tsInSqlDB;
private String serviceId; private String serviceId;
private AtomicInteger regularQueryInvocationCnt = new AtomicInteger(); private SubscriptionServiceStatistics stats = new SubscriptionServiceStatistics();
private AtomicInteger dynamicQueryInvocationCnt = new AtomicInteger();
private AtomicLong regularQueryTimeSpent = new AtomicLong();
private AtomicLong dynamicQueryTimeSpent = new AtomicLong();
@PostConstruct @PostConstruct
public void initExecutor() { public void initExecutor() {
@ -196,8 +193,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
long start = System.currentTimeMillis(); long start = System.currentTimeMillis();
ctx.fetchData(); ctx.fetchData();
long end = System.currentTimeMillis(); long end = System.currentTimeMillis();
regularQueryInvocationCnt.incrementAndGet(); stats.getRegularQueryInvocationCnt().incrementAndGet();
regularQueryTimeSpent.addAndGet(end - start); stats.getRegularQueryTimeSpent().addAndGet(end - start);
ctx.cancelTasks(); ctx.cancelTasks();
if (ctx.getQuery().getPageLink().isDynamic()) { if (ctx.getQuery().getPageLink().isDynamic()) {
//TODO: validate number of dynamic page links against rate limits. Ignore dynamic flag if limit is reached. //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()); ctx.setAndResolveQuery(cmd.getQuery());
AlarmDataQuery adq = ctx.getQuery(); AlarmDataQuery adq = ctx.getQuery();
long start = System.currentTimeMillis();
ctx.fetchData(); ctx.fetchData();
long end = System.currentTimeMillis();
stats.getRegularQueryInvocationCnt().incrementAndGet();
stats.getRegularQueryTimeSpent().addAndGet(end - start);
List<EntityData> entities = ctx.getEntitiesData(); List<EntityData> entities = ctx.getEntitiesData();
ctx.cancelTasks(); ctx.cancelTasks();
ctx.clearEntitySubscriptions(); ctx.clearEntitySubscriptions();
if (entities.isEmpty()) { 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); wsService.sendWsMsg(ctx.getSessionId(), update);
} else { } else {
ctx.fetchAlarms(); ctx.fetchAlarms();
@ -269,8 +272,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
long start = System.currentTimeMillis(); long start = System.currentTimeMillis();
finalCtx.update(); finalCtx.update();
long end = System.currentTimeMillis(); long end = System.currentTimeMillis();
dynamicQueryInvocationCnt.incrementAndGet(); stats.getDynamicQueryInvocationCnt().incrementAndGet();
dynamicQueryTimeSpent.addAndGet(end - start); stats.getDynamicQueryTimeSpent().addAndGet(end - start);
} catch (Exception e) { } catch (Exception e) {
log.warn("[{}][{}] Failed to refresh query", finalCtx.getSessionId(), finalCtx.getCmdId(), 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}") @Scheduled(fixedDelayString = "${server.ws.dynamic_page_link.stats:10000}")
public void printStats() { public void printStats() {
int regularQueryInvocationCntValue = regularQueryInvocationCnt.getAndSet(0); int alarmQueryInvocationCntValue = stats.getAlarmQueryInvocationCnt().getAndSet(0);
long regularQueryInvocationTimeValue = regularQueryTimeSpent.getAndSet(0); long alarmQueryInvocationTimeValue = stats.getAlarmQueryTimeSpent().getAndSet(0);
int dynamicQueryInvocationCntValue = dynamicQueryInvocationCnt.getAndSet(0); int regularQueryInvocationCntValue = stats.getRegularQueryInvocationCnt().getAndSet(0);
long dynamicQueryInvocationTimeValue = dynamicQueryTimeSpent.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(); long dynamicQueryCnt = subscriptionsBySessionId.values().stream().map(Map::values).count();
if (regularQueryInvocationCntValue > 0 || dynamicQueryInvocationCntValue > 0 || dynamicQueryCnt > 0) { if (regularQueryInvocationCntValue > 0 || dynamicQueryInvocationCntValue > 0 || dynamicQueryCnt > 0) {
log.info("Stats: regularQueryInvocationCnt = [{}], regularQueryInvocationTime = [{}], dynamicQueryCnt = [{}] dynamicQueryInvocationCnt = [{}], dynamicQueryInvocationTime = [{}]", log.info("Stats: regularQueryInvocationCnt = [{}], regularQueryInvocationTime = [{}], " +
regularQueryInvocationCntValue, regularQueryInvocationTimeValue, dynamicQueryCnt, dynamicQueryInvocationCntValue, dynamicQueryInvocationTimeValue); "dynamicQueryCnt = [{}] dynamicQueryInvocationCnt = [{}], dynamicQueryInvocationTime = [{}], " +
"alarmQueryInvocationCnt = [{}], alarmQueryInvocationTime = [{}]",
regularQueryInvocationCntValue, regularQueryInvocationTimeValue,
dynamicQueryCnt, dynamicQueryInvocationCntValue, dynamicQueryInvocationTimeValue,
alarmQueryInvocationCntValue, alarmQueryInvocationTimeValue);
} }
} }
private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) { private TbEntityDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, EntityDataCmd cmd) {
Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); Map<Integer, TbAbstractDataSubCtx> 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()); ctx.setAndResolveQuery(cmd.getQuery());
sessionSubs.put(cmd.getCmdId(), ctx); sessionSubs.put(cmd.getCmdId(), ctx);
return ctx; return ctx;
@ -299,7 +308,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) { private TbAlarmDataSubCtx createSubCtx(TelemetryWebSocketSessionRef sessionRef, AlarmDataCmd cmd) {
Map<Integer, TbAbstractDataSubCtx> sessionSubs = subscriptionsBySessionId.computeIfAbsent(sessionRef.getSessionId(), k -> new HashMap<>()); Map<Integer, TbAbstractDataSubCtx> 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()); ctx.setAndResolveQuery(cmd.getQuery());
sessionSubs.put(cmd.getCmdId(), ctx); sessionSubs.put(cmd.getCmdId(), ctx);
return ctx; return ctx;
@ -454,6 +463,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
if (ctx != null) { if (ctx != null) {
ctx.cancelTasks(); ctx.cancelTasks();
ctx.clearEntitySubscriptions(); ctx.clearEntitySubscriptions();
ctx.clearDynamicValueSubscriptions();
} }
} }

31
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();
}

45
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 not use this file except in compliance with the License.
* You may obtain a copy of the License at * 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 * Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, * 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.ArrayList;
import java.util.Arrays; import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.HashSet; import java.util.HashSet;
import java.util.List; import java.util.List;
@ -58,12 +59,15 @@ import java.util.Optional;
import java.util.Set; import java.util.Set;
import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutionException;
import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledFuture;
import java.util.function.Function;
import java.util.stream.Collectors;
@Slf4j @Slf4j
@Data @Data
public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends EntityDataPageLink>> { public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends EntityDataPageLink>> {
protected final String serviceId; protected final String serviceId;
protected final SubscriptionServiceStatistics stats;
protected final TelemetryWebSocketService wsService; protected final TelemetryWebSocketService wsService;
protected final EntityService entityService; protected final EntityService entityService;
protected final TbLocalSubscriptionService localSubscriptionService; protected final TbLocalSubscriptionService localSubscriptionService;
@ -84,12 +88,14 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService, public TbAbstractDataSubCtx(String serviceId, TelemetryWebSocketService wsService,
EntityService entityService, TbLocalSubscriptionService localSubscriptionService, EntityService entityService, TbLocalSubscriptionService localSubscriptionService,
AttributesService attributesService, TelemetryWebSocketSessionRef sessionRef, int cmdId) { AttributesService attributesService, SubscriptionServiceStatistics stats,
TelemetryWebSocketSessionRef sessionRef, int cmdId) {
this.serviceId = serviceId; this.serviceId = serviceId;
this.wsService = wsService; this.wsService = wsService;
this.entityService = entityService; this.entityService = entityService;
this.localSubscriptionService = localSubscriptionService; this.localSubscriptionService = localSubscriptionService;
this.attributesService = attributesService; this.attributesService = attributesService;
this.stats = stats;
this.sessionRef = sessionRef; this.sessionRef = sessionRef;
this.cmdId = cmdId; this.cmdId = cmdId;
this.subToEntityIdMap = new HashMap<>(); this.subToEntityIdMap = new HashMap<>();
@ -151,8 +157,6 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
subToDynamicValueKeySet.add(subIdx); subToDynamicValueKeySet.add(subIdx);
localSubscriptionService.addSubscription(sub); localSubscriptionService.addSubscription(sub);
} }
} catch (InterruptedException | ExecutionException e) { } catch (InterruptedException | ExecutionException e) {
log.info("[{}][{}][{}] Failed to resolve dynamic values: {}", tenantId, customerId, userId, dynamicValues.keySet()); log.info("[{}][{}][{}] Failed to resolve dynamic values: {}", tenantId, customerId, userId, dynamicValues.keySet());
} }
@ -197,7 +201,28 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
return result; return result;
} }
protected abstract void update(); protected synchronized void update(){
long start = System.currentTimeMillis();
PageData<EntityData> newData = findEntityData();
long end = System.currentTimeMillis();
stats.getRegularQueryInvocationCnt().incrementAndGet();
stats.getRegularQueryTimeSpent().addAndGet(end - start);
Map<EntityId, EntityData> oldDataMap;
if (data != null && !data.getData().isEmpty()) {
oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity()));
} else {
oldDataMap = Collections.emptyMap();
}
Map<EntityId, EntityData> 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<EntityId, EntityData> newDataMap);
protected abstract EntityDataQuery buildEntityDataQuery(); protected abstract EntityDataQuery buildEntityDataQuery();
@ -312,6 +337,15 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
} }
} }
public void clearDynamicValueSubscriptions(){
if (subToDynamicValueKeySet != null) {
for (Integer subId : subToDynamicValueKeySet) {
localSubscriptionService.cancelSubscription(sessionRef.getSessionId(), subId);
}
subToDynamicValueKeySet.clear();
}
}
public void setRefreshTask(ScheduledFuture<?> task) { public void setRefreshTask(ScheduledFuture<?> task) {
this.refreshTask = task; this.refreshTask = task;
} }
@ -323,7 +357,6 @@ public abstract class TbAbstractDataSubCtx<T extends AbstractDataQuery<? extends
} }
} }
public void createSubscriptions(List<EntityKey> keys, boolean resultToLatestValues) { public void createSubscriptions(List<EntityKey> keys, boolean resultToLatestValues) {
Map<EntityKeyType, List<EntityKey>> keysByType = getEntityKeyByTypeMap(keys); Map<EntityKeyType, List<EntityKey>> keysByType = getEntityKeyByTypeMap(keys);
for (EntityData entityData : data.getData()) { for (EntityData entityData : data.getData()) {

87
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.AlarmSubscriptionUpdate;
import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate;
import java.util.ArrayList;
import java.util.Collection; import java.util.Collection;
import java.util.Collections; import java.util.Collections;
import java.util.HashMap; import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashMap; import java.util.LinkedHashMap;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
import java.util.Set;
import java.util.function.Function; import java.util.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -74,9 +77,9 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService, public TbAlarmDataSubCtx(String serviceId, TelemetryWebSocketService wsService,
EntityService entityService, TbLocalSubscriptionService localSubscriptionService, EntityService entityService, TbLocalSubscriptionService localSubscriptionService,
AttributesService attributesService, AlarmService alarmService, AttributesService attributesService, SubscriptionServiceStatistics stats, AlarmService alarmService,
TelemetryWebSocketSessionRef sessionRef, int cmdId, int maxEntitiesPerAlarmSubscription) { 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.maxEntitiesPerAlarmSubscription = maxEntitiesPerAlarmSubscription;
this.alarmService = alarmService; this.alarmService = alarmService;
this.entitiesMap = new LinkedHashMap<>(); this.entitiesMap = new LinkedHashMap<>();
@ -84,10 +87,14 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
} }
public void fetchAlarms() { public void fetchAlarms() {
long start = System.currentTimeMillis();
PageData<AlarmData> alarms = alarmService.findAlarmDataByQueryForEntities(getTenantId(), getCustomerId(), PageData<AlarmData> alarms = alarmService.findAlarmDataByQueryForEntities(getTenantId(), getCustomerId(),
query, getOrderedEntityIds()); query, getOrderedEntityIds());
long end = System.currentTimeMillis();
stats.getAlarmQueryInvocationCnt().incrementAndGet();
stats.getAlarmQueryTimeSpent().addAndGet(end - start);
alarms = setAndMergeAlarmsData(alarms); alarms = setAndMergeAlarmsData(alarms);
AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null); AlarmDataUpdate update = new AlarmDataUpdate(cmdId, alarms, null, tooManyEntities);
wsService.sendWsMsg(getSessionId(), update); wsService.sendWsMsg(getSessionId(), update);
} }
@ -130,23 +137,27 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
AlarmDataPageLink pageLink = query.getPageLink(); AlarmDataPageLink pageLink = query.getPageLink();
long startTs = System.currentTimeMillis() - pageLink.getTimeWindow(); long startTs = System.currentTimeMillis() - pageLink.getTimeWindow();
for (EntityData entityData : entitiesMap.values()) { for (EntityData entityData : entitiesMap.values()) {
int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); createAlarmSubscriptionForEntity(pageLink, startTs, entityData);
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);
} }
} }
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 @Override
void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) { void sendWsMsg(String sessionId, TelemetrySubscriptionUpdate subscriptionUpdate, EntityKeyType keyType, boolean resultToLatestValues) {
EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId()); EntityId entityId = subToEntityIdMap.get(subscriptionUpdate.getSubscriptionId());
@ -163,7 +174,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
alarm.getLatest().computeIfAbsent(keyType, tmp -> new HashMap<>()).putAll(latestUpdate); alarm.getLatest().computeIfAbsent(keyType, tmp -> new HashMap<>()).putAll(latestUpdate);
return alarm; return alarm;
}).collect(Collectors.toList()); }).collect(Collectors.toList());
wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update)); wsService.sendWsMsg(sessionId, new AlarmDataUpdate(cmdId, null, update, tooManyEntities));
} else { } else {
log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate); log.trace("[{}][{}][{}][{}] Received stale subscription update: {}", sessionId, cmdId, subscriptionUpdate.getSubscriptionId(), keyType, subscriptionUpdate);
} }
@ -186,7 +197,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId()); AlarmData updated = new AlarmData(alarm, current.getOriginatorName(), current.getEntityId());
updated.getLatest().putAll(current.getLatest()); updated.getLatest().putAll(current.getLatest());
alarmsMap.put(alarmId, updated); 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 { } else {
fetchAlarms(); fetchAlarms();
} }
@ -241,8 +252,42 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx<AlarmDataQuery> {
} }
@Override @Override
protected void update() { protected synchronized void doUpdate(Map<EntityId, EntityData> newDataMap) {
entitiesMap.clear();
tooManyEntities = data.hasNext();
for (EntityData entityData : data.getData()) {
entitiesMap.put(entityData.getEntityId(), entityData);
}
fetchAlarms();
List<Integer> subIdsToCancel = new ArrayList<>();
List<TbSubscription> subsToAdd = new ArrayList<>();
Set<EntityId> 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<EntityData> newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList());
if (!newSubsList.isEmpty()) {
List<EntityKey> keys = query.getLatestValues();
if (keys != null && !keys.isEmpty()) {
Map<EntityKeyType, List<EntityKey>> 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 @Override

91
application/src/main/java/org/thingsboard/server/service/subscription/TbEntityDataSubCtx.java

@ -66,8 +66,8 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService,
TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService, TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService,
TelemetryWebSocketSessionRef sessionRef, int cmdId) { SubscriptionServiceStatistics stats, TelemetryWebSocketSessionRef sessionRef, int cmdId) {
super(serviceId, wsService, entityService, localSubscriptionService, attributesService, sessionRef, cmdId); super(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmdId);
} }
@Override @Override
@ -171,58 +171,45 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx<EntityDataQuery> {
return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null); return data.getData().stream().filter(item -> item.getEntityId().equals(entityId)).findFirst().orElse(null);
} }
public synchronized void update() { public synchronized void doUpdate(Map<EntityId, EntityData> newDataMap) {
PageData<EntityData> newData = findEntityData(); List<Integer> subIdsToCancel = new ArrayList<>();
Map<EntityId, EntityData> oldDataMap; List<TbSubscription> subsToAdd = new ArrayList<>();
if (data != null && !data.getData().isEmpty()) { Set<EntityId> currentSubs = new HashSet<>();
oldDataMap = data.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); subToEntityIdMap.forEach((subId, entityId) -> {
} else { if (!newDataMap.containsKey(entityId)) {
oldDataMap = Collections.emptyMap(); subIdsToCancel.add(subId);
} } else {
Map<EntityId, EntityData> newDataMap = newData.getData().stream().collect(Collectors.toMap(EntityData::getEntityId, Function.identity())); currentSubs.add(entityId);
if (oldDataMap.size() == newDataMap.size() && oldDataMap.keySet().equals(newDataMap.keySet())) { }
log.trace("[{}][{}] No updates to entity data found", sessionRef.getSessionId(), cmdId); });
} else { log.trace("[{}][{}] Subscriptions that are invalid: {}", sessionRef.getSessionId(), cmdId, subIdsToCancel);
this.data = newData; subIdsToCancel.forEach(subToEntityIdMap::remove);
List<Integer> subIdsToCancel = new ArrayList<>(); List<EntityData> newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList());
List<TbSubscription> subsToAdd = new ArrayList<>(); if (!newSubsList.isEmpty()) {
Set<EntityId> currentSubs = new HashSet<>(); boolean resultToLatestValues;
subToEntityIdMap.forEach((subId, entityId) -> { List<EntityKey> keys = null;
if (!newDataMap.containsKey(entityId)) { if (curTsCmd != null) {
subIdsToCancel.add(subId); resultToLatestValues = false;
} else { keys = curTsCmd.getKeys().stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList());
currentSubs.add(entityId); } else if (latestValueCmd != null) {
} resultToLatestValues = true;
}); keys = latestValueCmd.getKeys();
log.trace("[{}][{}] Subscriptions that are invalid: {}", sessionRef.getSessionId(), cmdId, subIdsToCancel); } else {
subIdsToCancel.forEach(subToEntityIdMap::remove); resultToLatestValues = true;
List<EntityData> newSubsList = newDataMap.entrySet().stream().filter(entry -> !currentSubs.contains(entry.getKey())).map(Map.Entry::getValue).collect(Collectors.toList()); }
if (!newSubsList.isEmpty()) { if (keys != null && !keys.isEmpty()) {
boolean resultToLatestValues; Map<EntityKeyType, List<EntityKey>> keysByType = getEntityKeyByTypeMap(keys);
List<EntityKey> keys = null; newSubsList.forEach(
if (curTsCmd != null) { entity -> {
resultToLatestValues = false; log.trace("[{}][{}] Found new subscription for entity: {}", sessionRef.getSessionId(), cmdId, entity.getEntityId());
keys = curTsCmd.getKeys().stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()); subsToAdd.addAll(addSubscriptions(entity, keysByType, resultToLatestValues));
} else if (latestValueCmd != null) { }
resultToLatestValues = true; );
keys = latestValueCmd.getKeys();
} else {
resultToLatestValues = true;
}
if (keys != null && !keys.isEmpty()) {
Map<EntityKeyType, List<EntityKey>> 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) { public void setCurrentCmd(EntityDataCmd cmd) {

9
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<AlarmData> { public class AlarmDataUpdate extends DataUpdate<AlarmData> {
public AlarmDataUpdate(int cmdId, PageData<AlarmData> data, List<AlarmData> update) { private boolean tooManyEntities;
public AlarmDataUpdate(int cmdId, PageData<AlarmData> data, List<AlarmData> update, boolean tooManyEntities) {
super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null);
this.tooManyEntities = tooManyEntities;
} }
public AlarmDataUpdate(int cmdId, int errorCode, String errorMsg) { public AlarmDataUpdate(int cmdId, int errorCode, String errorMsg) {
@ -45,7 +48,9 @@ public class AlarmDataUpdate extends DataUpdate<AlarmData> {
@JsonProperty("data") PageData<AlarmData> data, @JsonProperty("data") PageData<AlarmData> data,
@JsonProperty("update") List<AlarmData> update, @JsonProperty("update") List<AlarmData> update,
@JsonProperty("errorCode") int errorCode, @JsonProperty("errorCode") int errorCode,
@JsonProperty("errorMsg") String errorMsg) { @JsonProperty("errorMsg") String errorMsg,
@JsonProperty("tooManyEntities") boolean tooManyEntities) {
super(cmdId, data, update, errorCode, errorMsg); super(cmdId, data, update, errorCode, errorMsg);
this.tooManyEntities = tooManyEntities;
} }
} }

Loading…
Cancel
Save