diff --git a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java index 627b7f8b63..d39c5fcb12 100644 --- a/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java +++ b/application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java @@ -343,10 +343,11 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer s -> { TbTimeSeriesSubscription sub = (TbTimeSeriesSubscription) s; List updateData = null; + Map keyStates = sub.getKeyStates(); if (sub.isAllKeys()) { if (sub.isLatestValues()) { for (TsKvEntry kv : data) { - if (!sub.getKeyStates().containsKey((kv.getKey())) || kv.getTs() > sub.getKeyStates().get(kv.getKey())) { + if (!keyStates.containsKey((kv.getKey())) || kv.getTs() > keyStates.get(kv.getKey())) { if (updateData == null) { updateData = new ArrayList<>(); } @@ -358,8 +359,8 @@ public class DefaultTbLocalSubscriptionService implements TbLocalSubscriptionSer } } else { for (TsKvEntry kv : data) { - if (sub.getKeyStates().containsKey((kv.getKey()))) { - if (!sub.isLatestValues() || kv.getTs() > sub.getKeyStates().get(kv.getKey())) { + if (keyStates.containsKey((kv.getKey()))) { + if (!sub.isLatestValues() || kv.getTs() > keyStates.get(kv.getKey())) { if (updateData == null) { updateData = new ArrayList<>(); } diff --git a/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java index c52a20754d..11d962eb7c 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java @@ -105,6 +105,8 @@ import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.stream.Collectors; +import static org.thingsboard.server.common.data.DataConstants.LATEST_TELEMETRY_SCOPE; + /** * Created by ashvayka on 27.03.18. */ @@ -123,7 +125,6 @@ public class DefaultWebSocketService implements WebSocketService { private static final String FAILED_TO_FETCH_DATA = "Failed to fetch data!"; private static final String FAILED_TO_FETCH_ATTRIBUTES = "Failed to fetch attributes!"; private static final String SESSION_META_DATA_NOT_FOUND = "Session meta-data not found!"; - private static final String LATEST_TELEMETRY_SCOPE = "LATEST_TELEMETRY"; private final ConcurrentMap wsSessionsMap = new ConcurrentHashMap<>(); @@ -668,25 +669,7 @@ public class DefaultWebSocketService implements WebSocketService { data.forEach(v -> subState.put(v.getKey(), v.getTs())); Lock subLock = new ReentrantLock(); - TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() - .serviceId(serviceId) - .sessionId(sessionId) - .subscriptionId(registerNewSessionSubId(sessionId, sessionRef, cmd.getCmdId())) - .tenantId(sessionRef.getSecurityCtx().getTenantId()) - .entityId(entityId) - .updateProcessor((subscription, update) -> { - subLock.lock(); - try { - sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); - } finally { - subLock.unlock(); - } - }) - .queryTs(queryTs) - .allKeys(true) - .keyStates(subState) - .latestValues(LATEST_TELEMETRY_SCOPE.equals(cmd.getScope())) - .build(); + TbTimeSeriesSubscription sub = createTbTimeSeriesSubscription(subState, subLock, sessionId, sessionRef, cmd, entityId, queryTs, true); subLock.lock(); try { @@ -714,6 +697,28 @@ public class DefaultWebSocketService implements WebSocketService { on(r -> Futures.addCallback(tsService.findAllLatest(sessionRef.getSecurityCtx().getTenantId(), entityId), callback, executor), callback::onFailure)); } + private TbTimeSeriesSubscription createTbTimeSeriesSubscription(Map subState, Lock subLock, String sessionId, WebSocketSessionRef sessionRef, TimeseriesSubscriptionCmd cmd, EntityId entityId, long queryTs, boolean allKeys) { + return TbTimeSeriesSubscription.builder() + .serviceId(serviceId) + .sessionId(sessionId) + .subscriptionId(registerNewSessionSubId(sessionId, sessionRef, cmd.getCmdId())) + .tenantId(sessionRef.getSecurityCtx().getTenantId()) + .entityId(entityId) + .updateProcessor((subscription, update) -> { + subLock.lock(); + try { + sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); + } finally { + subLock.unlock(); + } + }) + .queryTs(queryTs) + .allKeys(allKeys) + .keyStates(subState) + .latestValues(LATEST_TELEMETRY_SCOPE.equals(cmd.getScope())) + .build(); + } + private FutureCallback> getSubscriptionCallback(final WebSocketSessionRef sessionRef, final TimeseriesSubscriptionCmd cmd, final String sessionId, final EntityId entityId, final long queryTs, final long startTs, final List keys) { return new FutureCallback<>() { @@ -724,25 +729,7 @@ public class DefaultWebSocketService implements WebSocketService { data.forEach(v -> subState.put(v.getKey(), v.getTs())); Lock subLock = new ReentrantLock(); - TbTimeSeriesSubscription sub = TbTimeSeriesSubscription.builder() - .serviceId(serviceId) - .sessionId(sessionId) - .subscriptionId(registerNewSessionSubId(sessionId, sessionRef, cmd.getCmdId())) - .tenantId(sessionRef.getSecurityCtx().getTenantId()) - .entityId(entityId) - .updateProcessor((subscription, update) -> { - subLock.lock(); - try { - sendUpdate(subscription.getSessionId(), cmd.getCmdId(), update); - } finally { - subLock.unlock(); - } - }) - .queryTs(queryTs) - .allKeys(false) - .keyStates(subState) - .latestValues(LATEST_TELEMETRY_SCOPE.equals(cmd.getScope())) - .build(); + TbTimeSeriesSubscription sub = createTbTimeSeriesSubscription(subState, subLock, sessionId, sessionRef, cmd, entityId, queryTs, false); subLock.lock(); try { diff --git a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/sub/TelemetrySubscriptionUpdate.java b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/sub/TelemetrySubscriptionUpdate.java index 1a3046e301..b22b021a03 100644 --- a/application/src/main/java/org/thingsboard/server/service/ws/telemetry/sub/TelemetrySubscriptionUpdate.java +++ b/application/src/main/java/org/thingsboard/server/service/ws/telemetry/sub/TelemetrySubscriptionUpdate.java @@ -16,7 +16,6 @@ package org.thingsboard.server.service.ws.telemetry.sub; import lombok.AllArgsConstructor; -import net.minidev.json.annotate.JsonIgnore; import org.thingsboard.server.common.data.kv.TsKvEntry; import org.thingsboard.server.service.subscription.SubscriptionErrorCode; diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java index b53d6daec2..b2d9d59cca 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java @@ -148,4 +148,6 @@ public class DataConstants { public static final String CF_QUEUE_NAME = "CalculatedFields"; public static final String CF_STATES_QUEUE_NAME = "CalculatedFieldStates"; + public static final String LATEST_TELEMETRY_SCOPE = "LATEST_TELEMETRY"; + }