From fa279c5b13b3421f254fd29f85096bd889e7c9ea Mon Sep 17 00:00:00 2001 From: Sergey Matvienko Date: Mon, 17 May 2021 14:45:31 +0300 Subject: [PATCH] default state service: refactored events order. long submit task split to a smaller pieces. throw exception at addDeviceUsingState if Device belongs to external partition --- .../state/DefaultDeviceStateService.java | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) 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 286260cf8e..ac18ec71d1 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 @@ -386,13 +386,14 @@ public class DefaultDeviceStateService extends TbApplicationEventListener submitNextPage(addedPartitions, tenant, pageLink, scheduledExecutor)); + scheduledExecutor.submit(() -> processPageAndSubmitNextPage(addedPartitions, tenant, pageLink, scheduledExecutor)); } return true; } - private void submitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink, final ExecutorService executor) { + private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink, final ExecutorService executor) { + log.trace("[{}] Process page {} from {}", tenant, pageLink.getPage(), pageLink.getPageSize()); List> fetchFutures = new ArrayList<>(); PageData page = deviceService.findDevicesByTenantId(tenant.getId(), pageLink); for (Device device : page.getData()) { @@ -433,7 +434,8 @@ public class DefaultDeviceStateService extends TbApplicationEventListener submitNextPage(addedPartitions, tenant, nextPageLink, executor)); + log.trace("[{}] Submit next page {} from {}", tenant, nextPageLink.getPage(), nextPageLink.getPageSize()); + executor.submit(() -> processPageAndSubmitNextPage(addedPartitions, tenant, nextPageLink, executor)); } } @@ -449,8 +451,13 @@ public class DefaultDeviceStateService extends TbApplicationEventListener ConcurrentHashMap.newKeySet()).add(state.getDeviceId()); - deviceStates.put(state.getDeviceId(), state); + Set deviceIds = partitionedDevices.get(tpi); + if (deviceIds != null) { + deviceIds.add(state.getDeviceId()); + deviceStates.put(state.getDeviceId(), state); + } else { + new RuntimeException("Device belongs to external partition " + tpi.getFullTopicName() + "!"); + } } void updateInactivityStateIfExpired() {