Browse Source

Fix for multiple commands in same subscription

pull/7288/head
Andrii Shvaika 4 years ago
parent
commit
6879021027
  1. 124
      application/src/main/java/org/thingsboard/server/service/subscription/DefaultTbEntityDataSubscriptionService.java

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

@ -29,7 +29,6 @@ import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.web.socket.CloseStatus; import org.springframework.web.socket.CloseStatus;
import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.common.util.ThingsBoardThreadFactory;
import org.thingsboard.server.common.data.kv.Aggregation;
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;
@ -70,7 +69,6 @@ import javax.annotation.PreDestroy;
import java.util.ArrayList; import java.util.ArrayList;
import java.util.Arrays; import java.util.Arrays;
import java.util.HashMap; import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashSet; import java.util.LinkedHashSet;
import java.util.List; import java.util.List;
import java.util.Map; import java.util.Map;
@ -178,6 +176,8 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
ctx = createSubCtx(session, cmd); ctx = createSubCtx(session, cmd);
} }
ctx.setCurrentCmd(cmd); ctx.setCurrentCmd(cmd);
// Fetch entity list using entity data query
if (cmd.getQuery() != null) { if (cmd.getQuery() != null) {
if (ctx.getQuery() == null) { if (ctx.getQuery() == null) {
log.debug("[{}][{}] Initializing data using query: {}", session.getSessionId(), cmd.getCmdId(), cmd.getQuery()); log.debug("[{}][{}] Initializing data using query: {}", session.getSessionId(), cmd.getCmdId(), cmd.getQuery());
@ -209,54 +209,54 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
finalCtx.setRefreshTask(task); finalCtx.setRefreshTask(task);
} }
} }
if (cmd.getAggHistoryCmd() != null) {
handleAggHistoryCmd(ctx, cmd.getAggHistoryCmd(), cmd.getLatestCmd());
} else if (cmd.getAggTsCmd() != null) {
handleAggTsCmd(ctx, cmd.getAggTsCmd(), cmd.getLatestCmd());
} else if (cmd.hasRegularCmds()) {
handleRegularCommands(ctx, cmd);
} else {
checkAndSendInitialData(ctx);
}
}
private void handleRegularCommands(TbEntityDataSubCtx ctx, EntityDataCmd cmd) { try {
ListenableFuture<TbEntityDataSubCtx> historyFuture; List<ListenableFuture<?>> cmdFutures = new ArrayList<>();
if (cmd.getHistoryCmd() != null) { if (cmd.getAggHistoryCmd() != null) {
log.trace("[{}][{}] Going to process history command: {}", ctx.getSessionId(), cmd.getCmdId(), cmd.getHistoryCmd()); cmdFutures.add(handleAggHistoryCmd(ctx, cmd.getAggHistoryCmd()));
try {
historyFuture = handleHistoryCmd(ctx, cmd.getHistoryCmd());
} catch (RuntimeException e) {
handleWsCmdRuntimeException(ctx.getSessionId(), e, cmd);
return;
} }
} else { if (cmd.getAggTsCmd() != null) {
historyFuture = Futures.immediateFuture(ctx); cmdFutures.add(handleAggTsCmd(ctx, cmd.getAggTsCmd()));
} }
Futures.addCallback(historyFuture, new FutureCallback<>() { if (cmd.getHistoryCmd() != null) {
@Override cmdFutures.add(handleHistoryCmd(ctx, cmd.getHistoryCmd()));
public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) { }
try { if (cmdFutures.isEmpty()) {
if (cmd.getLatestCmd() != null || cmd.getTsCmd() != null) { handleRegularCommands(ctx, cmd);
if (cmd.getLatestCmd() != null) { } else {
handleLatestCmd(theCtx, cmd.getLatestCmd()); TbEntityDataSubCtx finalCtx = ctx;
} Futures.addCallback(Futures.allAsList(cmdFutures), new FutureCallback<>() {
if (cmd.getTsCmd() != null) { @Override
handleTimeSeriesCmd(theCtx, cmd.getTsCmd()); public void onSuccess(@Nullable List<Object> result) {
} handleRegularCommands(finalCtx, cmd);
} else {
checkAndSendInitialData(theCtx);
} }
} catch (RuntimeException e) {
handleWsCmdRuntimeException(theCtx.getSessionId(), e, cmd); @Override
} public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to process command", finalCtx.getSessionId(), finalCtx.getCmdId());
}
}, wsCallBackExecutor);
} }
} catch (RuntimeException e) {
handleWsCmdRuntimeException(ctx.getSessionId(), e, cmd);
}
}
@Override private void handleRegularCommands(TbEntityDataSubCtx ctx, EntityDataCmd cmd) {
public void onFailure(Throwable t) { try {
log.warn("[{}][{}] Failed to process command", ctx.getSessionId(), cmd.getCmdId()); if (cmd.getLatestCmd() != null || cmd.getTsCmd() != null) {
if (cmd.getLatestCmd() != null) {
handleLatestCmd(ctx, cmd.getLatestCmd());
}
if (cmd.getTsCmd() != null) {
handleTimeSeriesCmd(ctx, cmd.getTsCmd());
}
} else {
checkAndSendInitialData(ctx);
} }
}, wsCallBackExecutor); } catch (RuntimeException e) {
handleWsCmdRuntimeException(ctx.getSessionId(), e, cmd);
}
} }
private void checkAndSendInitialData(@Nullable TbEntityDataSubCtx theCtx) { private void checkAndSendInitialData(@Nullable TbEntityDataSubCtx theCtx) {
@ -267,30 +267,30 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
} }
} }
private void handleAggHistoryCmd(TbEntityDataSubCtx ctx, AggHistoryCmd cmd, LatestValueCmd latestCmd) { private ListenableFuture<TbEntityDataSubCtx> handleAggHistoryCmd(TbEntityDataSubCtx ctx, AggHistoryCmd cmd) {
var keys = cmd.getKeys(); var keys = cmd.getKeys();
long interval = cmd.getEndTs() - cmd.getStartTs(); long interval = cmd.getEndTs() - cmd.getStartTs();
List<ReadTsKvQuery> queries = keys.stream().map(key -> new BaseReadTsKvQuery( List<ReadTsKvQuery> queries = keys.stream().map(key -> new BaseReadTsKvQuery(
key.getKey(), cmd.getStartTs(), cmd.getEndTs(), interval, 1, key.getAgg() key.getKey(), cmd.getStartTs(), cmd.getEndTs(), interval, 1, key.getAgg()
)).distinct().collect(Collectors.toList()); )).distinct().collect(Collectors.toList());
handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getEndTs(), latestCmd, false); return handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), cmd.getEndTs(), false);
} }
private void handleAggTsCmd(TbEntityDataSubCtx ctx, AggTimeSeriesCmd cmd, LatestValueCmd latestCmd) { private ListenableFuture<TbEntityDataSubCtx> handleAggTsCmd(TbEntityDataSubCtx ctx, AggTimeSeriesCmd cmd) {
long endTs = cmd.getStartTs() + cmd.getTimeWindow(); long endTs = cmd.getStartTs() + cmd.getTimeWindow();
List<ReadTsKvQuery> queries = cmd.getKeys().stream() List<ReadTsKvQuery> queries = cmd.getKeys().stream()
.map(key -> new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), endTs, cmd.getTimeWindow(), 1, key.getAgg())) .map(key -> new BaseReadTsKvQuery(key.getKey(), cmd.getStartTs(), endTs, cmd.getTimeWindow(), 1, key.getAgg()))
.distinct().collect(Collectors.toList()); .distinct().collect(Collectors.toList());
handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), endTs, latestCmd, true); return handleAggCmd(ctx, cmd.getKeys(), queries, cmd.getStartTs(), endTs, true);
} }
private void handleAggCmd(TbEntityDataSubCtx ctx, List<AggKey> keys, List<ReadTsKvQuery> queries, private ListenableFuture<TbEntityDataSubCtx> handleAggCmd(TbEntityDataSubCtx ctx, List<AggKey> keys, List<ReadTsKvQuery> queries,
long startTs, long endTs, LatestValueCmd latestCmd, boolean subscribe) { long startTs, long endTs, boolean subscribe) {
Map<EntityData, ListenableFuture<List<ReadTsKvQueryResult>>> fetchResultMap = new HashMap<>(); Map<EntityData, ListenableFuture<List<ReadTsKvQueryResult>>> fetchResultMap = new HashMap<>();
List<EntityData> entityDataList = ctx.getData().getData(); List<EntityData> entityDataList = ctx.getData().getData();
entityDataList.forEach(entityData -> fetchResultMap.put(entityData, entityDataList.forEach(entityData -> fetchResultMap.put(entityData,
tsService.findAllByQueries(ctx.getTenantId(), entityData.getEntityId(), queries))); tsService.findAllByQueries(ctx.getTenantId(), entityData.getEntityId(), queries)));
var mainFuture = Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> { return Futures.transform(Futures.allAsList(fetchResultMap.values()), f -> {
// Map that holds last ts for each key for each entity. // Map that holds last ts for each key for each entity.
Map<EntityData, Map<String, Long>> lastTsEntityMap = new HashMap<>(); Map<EntityData, Map<String, Long>> lastTsEntityMap = new HashMap<>();
fetchResultMap.forEach((entityData, future) -> { fetchResultMap.forEach((entityData, future) -> {
@ -333,28 +333,6 @@ public class DefaultTbEntityDataSubscriptionService implements TbEntityDataSubsc
} }
return ctx; return ctx;
}, wsCallBackExecutor); }, wsCallBackExecutor);
Futures.addCallback(mainFuture, new FutureCallback<>() {
@Override
public void onSuccess(@Nullable TbEntityDataSubCtx theCtx) {
if (latestCmd != null) {
if (subscribe) {
var commandEntityKeys = new HashSet<>(latestCmd.getKeys());
var alreadySubscribedKeys = keys.stream().map(key -> new EntityKey(EntityKeyType.TIME_SERIES, key.getKey())).collect(Collectors.toList());
if (commandEntityKeys.removeAll(alreadySubscribedKeys)) {
latestCmd.setKeys(new ArrayList<>(commandEntityKeys));
}
}
handleLatestCmd(ctx, latestCmd);
}
}
@Override
public void onFailure(Throwable t) {
log.warn("[{}][{}] Failed to process command", ctx.getSessionId(), ctx.getCmdId());
}
}, wsCallBackExecutor);
} }
private void handleWsCmdRuntimeException(String sessionId, RuntimeException e, EntityDataCmd cmd) { private void handleWsCmdRuntimeException(String sessionId, RuntimeException e, EntityDataCmd cmd) {

Loading…
Cancel
Save