From ad590da509aba32ee8384c2257ec178711f8cdb8 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Thu, 4 Aug 2022 16:18:15 +0300 Subject: [PATCH] Improvement to Device State Service --- .../AbstractPartitionBasedService.java | 2 +- .../state/DefaultDeviceStateService.java | 43 +++++++++++++------ 2 files changed, 30 insertions(+), 15 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java b/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java index 56a78474a1..305696662b 100644 --- a/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java +++ b/application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java @@ -126,7 +126,7 @@ public abstract class AbstractPartitionBasedService extends } List> fetchTasks = partitionedFetchTasks.remove(partition); if (fetchTasks != null) { - fetchTasks.forEach(f -> f.cancel(true)); + fetchTasks.forEach(f -> f.cancel(false)); } partitionListChanged = true; } diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index 667323c487..cd9d02f3e4 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -307,8 +307,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService partition : Lists.partition(entry.getValue(), 1000)) { log.info("[{}] Submit task for device states: {}", entry.getKey(), partition.size()); + DevicePackFutureHolder devicePackFutureHolder = new DevicePackFutureHolder(); var devicePackFuture = deviceStateExecutor.submit(() -> { - List states; - if (persistToTelemetry && !dbTypeInfoComponent.isLatestTsDaoStoredToSql()) { - states = fetchDeviceStateDataUsingSeparateRequests(partition); - } else { - states = fetchDeviceStateDataUsingEntityDataQuery(partition); - } - for (var state : states) { - addDeviceUsingState(entry.getKey(), state); - checkAndUpdateState(state.getDeviceId(), state); + try { + List states; + if (persistToTelemetry && !dbTypeInfoComponent.isLatestTsDaoStoredToSql()) { + states = fetchDeviceStateDataUsingSeparateRequests(partition); + } else { + states = fetchDeviceStateDataUsingEntityDataQuery(partition); + } + if (devicePackFutureHolder.future != null && !devicePackFutureHolder.future.isCancelled()) { + for (var state : states) { + if (!addDeviceUsingState(entry.getKey(), state)) { + return; + } + checkAndUpdateState(state.getDeviceId(), state); + } + log.info("[{}] Initialized {} out of {} device states", entry.getKey().getPartition().orElse(0), counter.addAndGet(states.size()), entry.getValue().size()); + } + } catch (Throwable t) { + log.error("Unexpected exception while device pack fetching", t); + throw t; } - log.info("[{}] Initialized {} out of {} device states", entry.getKey().getPartition().orElse(0), counter.addAndGet(states.size()), entry.getValue().size()); }); + devicePackFutureHolder.future = devicePackFuture; result.computeIfAbsent(entry.getKey(), tmp -> new ArrayList<>()).add(devicePackFuture); } } return result; } + private static class DevicePackFutureHolder { + private volatile ListenableFuture future; + } + void checkAndUpdateState(@Nonnull DeviceId deviceId, @Nonnull DeviceStateData state) { if (state.getState().isActive()) { updateInactivityStateIfExpired(System.currentTimeMillis(), deviceId, state); @@ -391,14 +405,15 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIds = partitionedEntities.get(tpi); if (deviceIds != null) { deviceIds.add(state.getDeviceId()); deviceStates.putIfAbsent(state.getDeviceId(), state); + return true; } else { log.debug("[{}] Device belongs to external partition {}", state.getDeviceId(), tpi.getFullTopicName()); - throw new RuntimeException("Device belongs to external partition " + tpi.getFullTopicName() + "!"); + return false; } }