From 10dd5c352f107cda96d2b85691724372fffdded1 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 11 Aug 2020 10:39:13 +0300 Subject: [PATCH] Max entities for data subscription over WS --- ...efaultTbEntityDataSubscriptionService.java | 20 +++++++++++-------- .../subscription/TbAbstractDataSubCtx.java | 7 ++++--- .../subscription/TbEntityDataSubCtx.java | 18 +++++++---------- .../telemetry/cmd/v2/EntityDataUpdate.java | 7 ++++++- .../src/main/resources/thingsboard.yml | 3 ++- 5 files changed, 31 insertions(+), 24 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 39c9c66a98..a968eb47d8 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 @@ -128,6 +128,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc private long dynamicPageLinkRefreshInterval; @Value("${server.ws.dynamic_page_link.refresh_pool_size:1}") private int dynamicPageLinkRefreshPoolSize; + @Value("${server.ws.max_entities_per_data_subscription:1000}") + private int maxEntitiesPerDataSubscription; @Value("${server.ws.max_entities_per_alarm_subscription:1000}") private int maxEntitiesPerAlarmSubscription; @@ -220,7 +222,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } else if (cmd.getTsCmd() != null) { handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); } else if (!theCtx.isInitialDataSent()) { - EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null); + EntityDataUpdate update = new EntityDataUpdate(theCtx.getCmdId(), theCtx.getData(), null, theCtx.getMaxEntitiesPerDataSubscription()); wsService.sendWsMsg(theCtx.getSessionId(), update); theCtx.setInitialDataSent(true); } @@ -298,7 +300,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc 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, stats, sessionRef, cmd.getCmdId()); + TbEntityDataSubCtx ctx = new TbEntityDataSubCtx(serviceId, wsService, entityService, localSubscriptionService, + attributesService, stats, sessionRef, cmd.getCmdId(), maxEntitiesPerDataSubscription); ctx.setAndResolveQuery(cmd.getQuery()); sessionSubs.put(cmd.getCmdId(), ctx); return ctx; @@ -306,7 +309,8 @@ 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, stats, 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; @@ -372,10 +376,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }); EntityDataUpdate update; if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); ctx.setInitialDataSent(true); } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); } wsService.sendWsMsg(ctx.getSessionId(), update); if (subscribe) { @@ -422,10 +426,10 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }); EntityDataUpdate update; if (!ctx.isInitialDataSent()) { - update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); ctx.setInitialDataSent(true); } else { - update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData()); + update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); } wsService.sendWsMsg(ctx.getSessionId(), update); ctx.createSubscriptions(latestCmd.getKeys(), true); @@ -440,7 +444,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc }, wsCallBackExecutor); } else { if (!ctx.isInitialDataSent()) { - EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null); + EntityDataUpdate update = new EntityDataUpdate(ctx.getCmdId(), ctx.getData(), null, ctx.getMaxEntitiesPerDataSubscription()); wsService.sendWsMsg(ctx.getSessionId(), update); ctx.setInitialDataSent(true); } 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 cd09ea9d44..41fbfd7ab2 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 @@ -57,6 +57,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ScheduledFuture; import java.util.function.Function; @@ -98,9 +99,9 @@ public abstract class TbAbstractDataSubCtx(); - this.subToDynamicValueKeySet = new HashSet<>(); - this.dynamicValues = new HashMap<>(); + this.subToEntityIdMap = new ConcurrentHashMap<>(); + this.subToDynamicValueKeySet = ConcurrentHashMap.newKeySet(); + this.dynamicValues = new ConcurrentHashMap<>(); } public void setAndResolveQuery(T query) { 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 0e58c07994..c43ee664fc 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 @@ -15,13 +15,10 @@ */ package org.thingsboard.server.service.subscription; -import lombok.AllArgsConstructor; -import lombok.Data; import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.data.id.EntityId; -import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.common.data.query.EntityDataQuery; import org.thingsboard.server.common.data.query.EntityKey; @@ -38,8 +35,6 @@ import org.thingsboard.server.service.telemetry.cmd.v2.TimeSeriesCmd; import org.thingsboard.server.service.telemetry.sub.TelemetrySubscriptionUpdate; import java.util.ArrayList; -import java.util.Arrays; -import java.util.Collection; import java.util.Collections; import java.util.Comparator; import java.util.HashMap; @@ -48,8 +43,6 @@ import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; -import java.util.concurrent.ScheduledFuture; -import java.util.function.Function; import java.util.stream.Collectors; @Slf4j @@ -63,11 +56,14 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { private boolean initialDataSent; private TimeSeriesCmd curTsCmd; private LatestValueCmd latestValueCmd; + @Getter + private final int maxEntitiesPerDataSubscription; public TbEntityDataSubCtx(String serviceId, TelemetryWebSocketService wsService, EntityService entityService, TbLocalSubscriptionService localSubscriptionService, AttributesService attributesService, - SubscriptionServiceStatistics stats, TelemetryWebSocketSessionRef sessionRef, int cmdId) { + SubscriptionServiceStatistics stats, TelemetryWebSocketSessionRef sessionRef, int cmdId, int maxEntitiesPerDataSubscription) { super(serviceId, wsService, entityService, localSubscriptionService, attributesService, stats, sessionRef, cmdId); + this.maxEntitiesPerDataSubscription = maxEntitiesPerDataSubscription; } @Override @@ -120,7 +116,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { if (!latestUpdate.isEmpty()) { Map> latestMap = Collections.singletonMap(keyType, latestUpdate); entityData = new EntityData(entityId, latestMap, null); - wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData))); + wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); } } @@ -163,7 +159,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { Map tsMap = new HashMap<>(); tsUpdate.forEach((key, tsValue) -> tsMap.put(key, tsValue.toArray(new TsValue[tsValue.size()]))); entityData = new EntityData(entityId, null, tsMap); - wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData))); + wsService.sendWsMsg(sessionId, new EntityDataUpdate(cmdId, null, Collections.singletonList(entityData), maxEntitiesPerDataSubscription)); } } @@ -207,7 +203,7 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { ); } } - wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null)); + wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription)); subIdsToCancel.forEach(subId -> localSubscriptionService.cancelSubscription(getSessionId(), subId)); subsToAdd.forEach(localSubscriptionService::addSubscription); } diff --git a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java index 81a1ee1778..1467090874 100644 --- a/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/telemetry/cmd/v2/EntityDataUpdate.java @@ -18,6 +18,7 @@ package org.thingsboard.server.service.telemetry.cmd.v2; import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnoreProperties; import com.fasterxml.jackson.annotation.JsonProperty; +import lombok.Getter; import org.thingsboard.server.common.data.page.PageData; import org.thingsboard.server.common.data.query.EntityData; import org.thingsboard.server.service.telemetry.sub.SubscriptionErrorCode; @@ -27,8 +28,12 @@ import java.util.List; public class EntityDataUpdate extends DataUpdate { - public EntityDataUpdate(int cmdId, PageData data, List update) { + @Getter + private long allowedEntities; + + public EntityDataUpdate(int cmdId, PageData data, List update, long allowedEntities) { super(cmdId, data, update, SubscriptionErrorCode.NO_ERROR.getCode(), null); + this.allowedEntities = allowedEntities; } public EntityDataUpdate(int cmdId, int errorCode, String errorMsg) { diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 54a3dac870..9d5897b483 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -50,7 +50,8 @@ server: refresh_interval: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_INTERVAL_SEC:60}" refresh_pool_size: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_REFRESH_POOL_SIZE:1}" max_per_user: "${TB_SERVER_WS_DYNAMIC_PAGE_LINK_MAX_PER_USER:10}" - max_entities_per_alarm_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_ALARM_SUBSCRIPTION:1000}" + max_entities_per_data_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_DATA_SUBSCRIPTION:10000}" + max_entities_per_alarm_subscription: "${TB_SERVER_WS_MAX_ENTITIES_PER_ALARM_SUBSCRIPTION:10000}" rest: limits: tenant: