Browse Source

do not send update if count was not changed

pull/12652/head
dashevchenko 2 years ago
parent
commit
ca1185de54
  1. 22
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java
  2. 19
      application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmCountSubCtx.java

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

@ -33,6 +33,7 @@ import org.springframework.stereotype.Service;
import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.CloseStatus;
import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardExecutors;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.id.EntityId;
import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery; import org.thingsboard.server.common.data.kv.BaseReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
@ -59,6 +60,7 @@ import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggHistoryCmd;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggKey; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggKey;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggTimeSeriesCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AggTimeSeriesCmd;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountCmd;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmCountUpdate;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataCmd;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataUpdate; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmDataUpdate;
import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd; import org.thingsboard.server.service.ws.telemetry.cmd.v2.AlarmStatusCmd;
@ -426,15 +428,21 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
long end = System.currentTimeMillis(); long end = System.currentTimeMillis();
stats.getRegularQueryInvocationCnt().incrementAndGet(); stats.getRegularQueryInvocationCnt().incrementAndGet();
stats.getRegularQueryTimeSpent().addAndGet(end - start); stats.getRegularQueryTimeSpent().addAndGet(end - start);
Set<EntityId> entitiesIds = ctx.getEntitiesIds();
ctx.cancelTasks(); ctx.cancelTasks();
ctx.clearAlarmSubscriptions(); ctx.clearAlarmSubscriptions();
ctx.fetchAlarmCount(); if (entitiesIds != null && entitiesIds.isEmpty()) {
ctx.createAlarmSubscriptions(); AlarmCountUpdate update = new AlarmCountUpdate(cmd.getCmdId(), 0);
TbAlarmCountSubCtx finalCtx = ctx; ctx.sendWsMsg(update);
ScheduledFuture<?> task = scheduler.scheduleWithFixedDelay( } else {
() -> refreshDynamicQuery(finalCtx), ctx.doFetchAlarmCount();
dynamicPageLinkRefreshInterval, dynamicPageLinkRefreshInterval, TimeUnit.SECONDS); ctx.createAlarmSubscriptions();
finalCtx.setRefreshTask(task); TbAlarmCountSubCtx finalCtx = ctx;
ScheduledFuture<?> task = scheduler.scheduleWithFixedDelay(
() -> refreshDynamicQuery(finalCtx),
dynamicPageLinkRefreshInterval, dynamicPageLinkRefreshInterval, TimeUnit.SECONDS);
finalCtx.setRefreshTask(task);
}
} else { } else {
log.debug("[{}][{}] Received duplicate command: {}", session.getSessionId(), cmd.getCmdId(), cmd); log.debug("[{}][{}] Received duplicate command: {}", session.getSessionId(), cmd.getCmdId(), cmd);
} }

19
application/src/main/java/org/thingsboard/server/service/subscription/TbAlarmCountSubCtx.java

@ -48,6 +48,7 @@ public class TbAlarmCountSubCtx extends TbAbstractEntityQuerySubCtx<AlarmCountQu
protected final Map<Integer, EntityId> subToEntityIdMap; protected final Map<Integer, EntityId> subToEntityIdMap;
@Getter
private LinkedHashSet<EntityId> entitiesIds; private LinkedHashSet<EntityId> entitiesIds;
private final int maxEntitiesPerAlarmSubscription; private final int maxEntitiesPerAlarmSubscription;
@ -102,26 +103,30 @@ public class TbAlarmCountSubCtx extends TbAbstractEntityQuerySubCtx<AlarmCountQu
fetchAlarmCount(); fetchAlarmCount();
} }
@Override
public boolean isDynamic() {
return true;
}
public void fetchAlarmCount() { public void fetchAlarmCount() {
alarmCountInvocationAttempts++; alarmCountInvocationAttempts++;
log.trace("[{}] Fetching alarms: {}", cmdId, alarmCountInvocationAttempts); log.trace("[{}] Fetching alarms: {}", cmdId, alarmCountInvocationAttempts);
if (alarmCountInvocationAttempts <= maxAlarmQueriesPerRefreshInterval) { if (alarmCountInvocationAttempts <= maxAlarmQueriesPerRefreshInterval) {
doFetchAlarmCount(); int newCount = (int) alarmService.countAlarmsByQuery(getTenantId(), getCustomerId(), query, entitiesIds);
if (newCount != result) {
result = newCount;
sendWsMsg(new AlarmCountUpdate(cmdId, result));
}
} else { } else {
log.trace("[{}] Ignore alarm count fetch due to rate limit: [{}] of maximum [{}]", cmdId, alarmCountInvocationAttempts, maxAlarmQueriesPerRefreshInterval); log.trace("[{}] Ignore alarm count fetch due to rate limit: [{}] of maximum [{}]", cmdId, alarmCountInvocationAttempts, maxAlarmQueriesPerRefreshInterval);
} }
} }
private void doFetchAlarmCount() { public void doFetchAlarmCount() {
result = (int) alarmService.countAlarmsByQuery(getTenantId(), getCustomerId(), query, entitiesIds); result = (int) alarmService.countAlarmsByQuery(getTenantId(), getCustomerId(), query, entitiesIds);
sendWsMsg(new AlarmCountUpdate(cmdId, result)); sendWsMsg(new AlarmCountUpdate(cmdId, result));
} }
@Override
public boolean isDynamic() {
return true;
}
private EntityDataQuery buildEntityDataQuery() { private EntityDataQuery buildEntityDataQuery() {
EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null, EntityDataPageLink edpl = new EntityDataPageLink(maxEntitiesPerAlarmSubscription, 0, null,
new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, ModelConstants.CREATED_TIME_PROPERTY))); new EntityDataSortOrder(new EntityKey(EntityKeyType.ENTITY_FIELD, ModelConstants.CREATED_TIME_PROPERTY)));

Loading…
Cancel
Save