From b4184d014c51193a6cca70055bf9024666e488ff Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Tue, 16 Mar 2021 13:46:12 +0200 Subject: [PATCH] Improvements to startTime and endTime for subscriptions --- .../DefaultSubscriptionManagerService.java | 57 ++++++++++++------- ...efaultTbEntityDataSubscriptionService.java | 10 ++-- .../subscription/TbAbstractDataSubCtx.java | 25 +++++--- .../subscription/TbAlarmDataSubCtx.java | 6 +- .../subscription/TbEntityDataSubCtx.java | 34 ++++------- .../subscription/TbSubscriptionUtils.java | 2 + .../TbTimeseriesSubscription.java | 5 +- common/queue/src/main/proto/queue.proto | 1 + 8 files changed, 80 insertions(+), 60 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java index b73641e9ad..7a7e7de251 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultSubscriptionManagerService.java @@ -213,7 +213,7 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene }, s -> true, s -> { List subscriptionUpdate = null; for (TsKvEntry kv : ts) { - if (isInTimeRange(s, kv.getTs()) && (s.isAllKeys() || s.getKeyStates().containsKey((kv.getKey())))) { + if ((s.isAllKeys() || s.getKeyStates().containsKey((kv.getKey())))) { if (subscriptionUpdate == null) { subscriptionUpdate = new ArrayList<>(); } @@ -375,11 +375,6 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene } } - private boolean isInTimeRange(TbTimeseriesSubscription subscription, long kvTime) { - return (subscription.getStartTime() == 0 || subscription.getStartTime() <= kvTime) - && (subscription.getEndTime() == 0 || subscription.getEndTime() >= kvTime); - } - private void removeSubscriptionFromEntityMap(TbSubscription sub) { Set entitySubSet = subscriptionsByEntityId.get(sub.getEntityId()); if (entitySubSet != null) { @@ -429,18 +424,9 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene serviceId, subscription.getSessionId(), subscription.getSubscriptionId(), subscription.getEntityId()); long curTs = System.currentTimeMillis(); - List queries = new ArrayList<>(); - subscription.getKeyStates().forEach((key, value) -> { - if (curTs > value) { - long startTs = subscription.getStartTime() > 0 ? Math.max(subscription.getStartTime(), value + 1L) : (value + 1L); - long endTs = subscription.getEndTime() > 0 ? Math.min(subscription.getEndTime(), curTs) : curTs; - if (startTs > 1) { - queries.add(new BaseReadTsKvQuery(key, startTs, endTs, 0, 1000, Aggregation.NONE)); - } - } - }); - if (!queries.isEmpty()) { - DonAsynchron.withCallback(tsService.findAll(subscription.getTenantId(), subscription.getEntityId(), queries), + + if (subscription.isLatestValues()) { + DonAsynchron.withCallback(tsService.findLatest(subscription.getTenantId(), subscription.getEntityId(), subscription.getKeyStates().keySet()), missedUpdates -> { if (missedUpdates != null && !missedUpdates.isEmpty()) { TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); @@ -449,6 +435,26 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene }, e -> log.error("Failed to fetch missed updates.", e), tsCallBackExecutor); + } else { + List queries = new ArrayList<>(); + subscription.getKeyStates().forEach((key, value) -> { + if (curTs > value) { + long startTs = subscription.getStartTime() > 0 ? Math.max(subscription.getStartTime(), value + 1L) : (value + 1L); + long endTs = subscription.getEndTime() > 0 ? Math.min(subscription.getEndTime(), curTs) : curTs; + queries.add(new BaseReadTsKvQuery(key, startTs, endTs, 0, 1000, Aggregation.NONE)); + } + }); + if (!queries.isEmpty()) { + DonAsynchron.withCallback(tsService.findAll(subscription.getTenantId(), subscription.getEntityId(), queries), + missedUpdates -> { + if (missedUpdates != null && !missedUpdates.isEmpty()) { + TopicPartitionInfo tpi = partitionService.getNotificationsTopic(ServiceType.TB_CORE, subscription.getServiceId()); + toCoreNotificationsProducer.send(tpi, toProto(subscription, missedUpdates), null); + } + }, + e -> log.error("Failed to fetch missed updates.", e), + tsCallBackExecutor); + } } } @@ -470,12 +476,19 @@ public class DefaultSubscriptionManagerService extends TbApplicationEventListene data.forEach((key, value) -> { TbSubscriptionUpdateValueListProto.Builder dataBuilder = TbSubscriptionUpdateValueListProto.newBuilder(); dataBuilder.setKey(key); - value.forEach(v -> { + boolean hasData = false; + for (Object v : value) { Object[] array = (Object[]) v; dataBuilder.addTs((long) array[0]); - dataBuilder.addValue((String) array[1]); - }); - builder.addData(dataBuilder.build()); + String strVal = (String) array[1]; + if (strVal != null) { + hasData = true; + dataBuilder.addValue(strVal); + } + } + if (hasData) { + builder.addData(dataBuilder.build()); + } }); ToCoreNotificationMsg toCoreMsg = ToCoreNotificationMsg.newBuilder().setToLocalSubscriptionServiceMsg( 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 330b096b78..127de752d0 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 @@ -215,7 +215,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } else { historyFuture = Futures.immediateFuture(ctx); } - Futures.addCallback(historyFuture, new FutureCallback() { + Futures.addCallback(historyFuture, new FutureCallback<>() { @Override public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { if (cmd.getLatestCmd() != null) { @@ -278,7 +278,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc wsService.sendWsMsg(ctx.getSessionId(), update); } else { ctx.fetchAlarms(); - ctx.createSubscriptions(cmd.getQuery().getLatestValues(), true); + ctx.createLatestValuesSubscriptions(cmd.getQuery().getLatestValues()); if (adq.getPageLink().getTimeWindow() > 0) { TbAlarmDataSubCtx finalCtx = ctx; ScheduledFuture task = scheduler.scheduleWithFixedDelay( @@ -419,7 +419,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc } wsService.sendWsMsg(ctx.getSessionId(), update); if (subscribe) { - ctx.createSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), false); + ctx.createTimeseriesSubscriptions(keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key)).collect(Collectors.toList()), cmd.getStartTs(), cmd.getEndTs()); } ctx.getData().getData().forEach(ed -> ed.getTimeseries().clear()); return ctx; @@ -468,7 +468,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc update = new EntityDataUpdate(ctx.getCmdId(), null, ctx.getData().getData(), ctx.getMaxEntitiesPerDataSubscription()); } wsService.sendWsMsg(ctx.getSessionId(), update); - ctx.createSubscriptions(latestCmd.getKeys(), true); + ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); } @Override @@ -484,7 +484,7 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc wsService.sendWsMsg(ctx.getSessionId(), update); ctx.setInitialDataSent(true); } - ctx.createSubscriptions(latestCmd.getKeys(), true); + ctx.createLatestValuesSubscriptions(latestCmd.getKeys()); } } 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 08f7bfa1a6..e054c9436f 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 @@ -115,10 +115,18 @@ public abstract class TbAbstractDataSubCtx keys, boolean resultToLatestValues) { + public void createLatestValuesSubscriptions(List keys) { + createSubscriptions(keys, true, 0, 0); + } + + public void createTimeseriesSubscriptions(List keys, long startTs, long endTs) { + createSubscriptions(keys, false, startTs, endTs); + } + + private void createSubscriptions(List keys, boolean latestValues, long startTs, long endTs) { Map> keysByType = getEntityKeyByTypeMap(keys); for (EntityData entityData : data.getData()) { - List entitySubscriptions = addSubscriptions(entityData, keysByType, resultToLatestValues); + List entitySubscriptions = addSubscriptions(entityData, keysByType, latestValues, startTs, endTs); entitySubscriptions.forEach(localSubscriptionService::addSubscription); } } @@ -129,14 +137,14 @@ public abstract class TbAbstractDataSubCtx addSubscriptions(EntityData entityData, Map> keysByType, boolean resultToLatestValues) { + protected List addSubscriptions(EntityData entityData, Map> keysByType, boolean latestValues, long startTs, long endTs) { List subscriptionList = new ArrayList<>(); keysByType.forEach((keysType, keysList) -> { int subIdx = sessionRef.getSessionSubIdSeq().incrementAndGet(); subToEntityIdMap.put(subIdx, entityData.getEntityId()); switch (keysType) { case TIME_SERIES: - subscriptionList.add(createTsSub(entityData, subIdx, keysList, resultToLatestValues)); + subscriptionList.add(createTsSub(entityData, subIdx, keysList, latestValues, startTs, endTs)); break; case CLIENT_ATTRIBUTE: subscriptionList.add(createAttrSub(entityData, subIdx, keysType, TbAttributeSubscriptionScope.CLIENT_SCOPE, keysList)); @@ -171,9 +179,9 @@ public abstract class TbAbstractDataSubCtx subKeys, boolean resultToLatestValues) { + private TbSubscription createTsSub(EntityData entityData, int subIdx, List subKeys, boolean latestValues, long startTs, long endTs) { Map keyStates = buildKeyStats(entityData, EntityKeyType.TIME_SERIES, subKeys); - if (entityData.getTimeseries() != null) { + if (!latestValues && entityData.getTimeseries() != null) { entityData.getTimeseries().forEach((k, v) -> { long ts = Arrays.stream(v).map(TsValue::getTs).max(Long::compareTo).orElse(0L); log.trace("[{}][{}] Updating key: {} with ts: {}", serviceId, cmdId, k, ts); @@ -187,9 +195,12 @@ public abstract class TbAbstractDataSubCtx sendWsMsg(sessionId, subscriptionUpdate, EntityKeyType.TIME_SERIES, resultToLatestValues)) + .updateConsumer((sessionId, subscriptionUpdate) -> sendWsMsg(sessionId, subscriptionUpdate, EntityKeyType.TIME_SERIES, latestValues)) .allKeys(false) .keyStates(keyStates) + .latestValues(latestValues) + .startTime(startTs) + .endTime(endTs) .build(); } 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 8af9ba2330..893d7ed91c 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 @@ -132,8 +132,8 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { } @Override - public void createSubscriptions(List keys, boolean resultToLatestValues) { - super.createSubscriptions(keys, resultToLatestValues); + public void createLatestValuesSubscriptions(List keys) { + super.createLatestValuesSubscriptions(keys); createAlarmSubscriptions(); } @@ -282,7 +282,7 @@ public class TbAlarmDataSubCtx extends TbAbstractDataSubCtx { newSubsList.forEach( entity -> { log.trace("[{}][{}] Found new subscription for entity: {}", sessionRef.getSessionId(), cmdId, entity.getEntityId()); - subsToAdd.addAll(addSubscriptions(entity, keysByType, true)); + subsToAdd.addAll(addSubscriptions(entity, keysByType, true, 0, 0)); } ); } 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 fd09536c72..7325fe7878 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 @@ -48,9 +48,6 @@ import java.util.stream.Collectors; @Slf4j public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { - @Getter - @Setter - private TimeSeriesCmd tsCmd; @Getter @Setter private boolean initialDataSent; @@ -183,25 +180,18 @@ public class TbEntityDataSubCtx extends TbAbstractDataSubCtx { 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)); - } - ); + // NOTE: We ignore the TS subscriptions for new entities here, because widgets will re-init it's content and will create new subscriptions. + if (curTsCmd == null && latestValueCmd != null) { + List keys = latestValueCmd.getKeys(); + 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, 0, 0)); + } + ); + } } } wsService.sendWsMsg(sessionRef.getSessionId(), new EntityDataUpdate(cmdId, data, null, maxEntitiesPerDataSubscription)); diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java index c6da1aad29..d7780ac2b2 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbSubscriptionUtils.java @@ -83,6 +83,7 @@ public class TbSubscriptionUtils { TbSubscriptionKetStateProto.newBuilder().setKey(key).setTs(value).build())); tSubProto.setStartTime(tSub.getStartTime()); tSubProto.setEndTime(tSub.getEndTime()); + tSubProto.setLatestValues(tSub.isLatestValues()); msgBuilder.setTelemetrySub(tSubProto.build()); break; case ATTRIBUTES: @@ -146,6 +147,7 @@ public class TbSubscriptionUtils { telemetrySub.getKeyStatesList().forEach(ksProto -> keyStates.put(ksProto.getKey(), ksProto.getTs())); builder.startTime(telemetrySub.getStartTime()); builder.endTime(telemetrySub.getEndTime()); + builder.latestValues(telemetrySub.getLatestValues()); builder.keyStates(keyStates); return builder.build(); } diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java b/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java index b9f7a732cc..962889f741 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/TbTimeseriesSubscription.java @@ -34,16 +34,19 @@ public class TbTimeseriesSubscription extends TbSubscription updateConsumer, - boolean allKeys, Map keyStates, long startTime, long endTime) { + boolean allKeys, Map keyStates, long startTime, long endTime, boolean latestValues) { super(serviceId, sessionId, subscriptionId, tenantId, entityId, TbSubscriptionType.TIMESERIES, updateConsumer); this.allKeys = allKeys; this.keyStates = keyStates; this.startTime = startTime; this.endTime = endTime; + this.latestValues = latestValues; } @Override diff --git a/common/queue/src/main/proto/queue.proto b/common/queue/src/main/proto/queue.proto index 6864259617..b3e15016c5 100644 --- a/common/queue/src/main/proto/queue.proto +++ b/common/queue/src/main/proto/queue.proto @@ -337,6 +337,7 @@ message TbTimeSeriesSubscriptionProto { repeated TbSubscriptionKetStateProto keyStates = 3; int64 startTime = 4; int64 endTime = 5; + bool latestValues = 6; } message TbAttributeSubscriptionProto {