Browse Source

code refactoring

pull/13025/head
dashevchenko 2 years ago
parent
commit
8441a6bca2
  1. 7
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbLocalSubscriptionService.java
  2. 65
      application/src/main/java/org/thingsboard/server/service/ws/DefaultWebSocketService.java
  3. 1
      application/src/main/java/org/thingsboard/server/service/ws/telemetry/sub/TelemetrySubscriptionUpdate.java
  4. 2
      common/data/src/main/java/org/thingsboard/server/common/data/DataConstants.java

7
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<TsKvEntry> updateData = null;
Map<String, Long> 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<>();
}

65
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<String, WsSessionMetaData> 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<String, Long> 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<List<TsKvEntry>> getSubscriptionCallback(final WebSocketSessionRef sessionRef, final TimeseriesSubscriptionCmd cmd,
final String sessionId, final EntityId entityId, final long queryTs, final long startTs, final List<String> 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 {

1
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;

2
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";
}

Loading…
Cancel
Save