From 073ce69872d53143c57632c880d47de15430e5cc Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Tue, 15 Feb 2022 16:04:20 +0200 Subject: [PATCH 1/7] handling of PartitionChangeEvent is not more synchronous --- .../DefaultTbApiUsageStateService.java | 140 ++++++++++++++---- .../queue/discovery/HashPartitionService.java | 3 +- 2 files changed, 112 insertions(+), 31 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 6e8e8c20e4..db85f3c852 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -16,6 +16,10 @@ package org.thingsboard.server.service.apiusage; import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningScheduledExecutorService; +import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.checkerframework.checker.nullness.qual.Nullable; @@ -23,9 +27,9 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; -import org.thingsboard.common.util.ThingsBoardExecutors; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.rule.engine.api.MailService; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.ApiFeature; import org.thingsboard.server.common.data.ApiUsageRecordKey; import org.thingsboard.server.common.data.ApiUsageState; @@ -49,8 +53,8 @@ import org.thingsboard.server.common.data.tenant.profile.TenantProfileConfigurat import org.thingsboard.server.common.data.tenant.profile.TenantProfileData; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TbCallback; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.common.msg.tools.SchedulerUtils; -import org.thingsboard.server.dao.customer.CustomerService; import org.thingsboard.server.dao.tenant.TbTenantProfileCache; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; @@ -58,11 +62,11 @@ import org.thingsboard.server.dao.usagerecord.ApiUsageStateService; import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceMsg; import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; +import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.scheduler.SchedulerComponent; -import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.telemetry.InternalTelemetryService; import javax.annotation.PostConstruct; @@ -73,13 +77,15 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Queue; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -102,12 +108,12 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener otherUsageStates = new ConcurrentHashMap<>(); + private final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); + private final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); @Value("${usage.stats.report.enabled:true}") @@ -130,25 +138,29 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener> subscribeQueue = new ConcurrentLinkedQueue<>(); + public DefaultTbApiUsageStateService(TbClusterService clusterService, PartitionService partitionService, TenantService tenantService, - CustomerService customerService, TimeseriesService tsService, ApiUsageStateService apiUsageStateService, SchedulerComponent scheduler, TbTenantProfileCache tenantProfileCache, - MailService mailService) { + MailService mailService, + DbCallbackExecutorService dbExecutor) { this.clusterService = clusterService; this.partitionService = partitionService; this.tenantService = tenantService; - this.customerService = customerService; this.tsService = tsService; this.apiUsageStateService = apiUsageStateService; this.scheduler = scheduler; this.tenantProfileCache = tenantProfileCache; this.mailService = mailService; this.mailExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("api-usage-svc-mail")); + this.dbExecutor = dbExecutor; } @PostConstruct @@ -158,6 +170,7 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener { - return !partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition(); - }); - otherUsageStates.entrySet().removeIf(entry -> { - return partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition(); + subscribeQueue.add(partitionChangeEvent.getPartitions()); + tenantStateExecutor.submit(this::pollInitStateFromDB); + } + } + + void pollInitStateFromDB() { + final Set partitions = getLatestPartitionsFromQueue(); + if (partitions == null) { + log.info("Tenant state service. Nothing to do. partitions is null"); + return; + } + initStateFromDB(partitions); + } + + Set getLatestPartitionsFromQueue() { + log.debug("getLatestPartitionsFromQueue, queue size {}", subscribeQueue.size()); + Set partitions = null; + while (!subscribeQueue.isEmpty()) { + partitions = subscribeQueue.poll(); + log.debug("polled from the queue partitions {}", partitions); + } + log.debug("getLatestPartitionsFromQueue, partitions {}", partitions); + return partitions; + } + + private void initStateFromDB(Set partitions) { + try { + Set addedPartitions = new HashSet<>(partitions); + addedPartitions.removeAll(partitionedEntities.keySet()); + + Set removedPartitions = new HashSet<>(partitionedEntities.keySet()); + removedPartitions.removeAll(partitions); + + removedPartitions.forEach(partition -> { + Set entities = partitionedEntities.remove(partition); + entities.forEach(this::cleanUpEntitiesStateMap); }); - initStatesFromDataBase(); + + addedPartitions.forEach(tpi -> + partitionedEntities.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); + + otherUsageStates.entrySet().removeIf(entry -> + partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition()); + + initStatesFromDataBase(addedPartitions); + } catch (Throwable t) { + log.warn("Failed to init tenant states from DB", t); } } @@ -311,6 +364,18 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener entityIds = partitionedEntities.get(tpi); + if (entityIds != null) { + entityIds.add(entityId); + myUsageStates.put(entityId, state); + } else { + log.debug("[{}] belongs to external partition {}", entityId, tpi.getFullTopicName()); + throw new RuntimeException(entityId.getEntityType() + " belongs to external partition " + tpi.getFullTopicName() + "!"); + } + } + private void updateProfileThresholds(TenantId tenantId, ApiUsageStateId id, TenantProfileConfiguration oldData, TenantProfileConfiguration newData) { long ts = System.currentTimeMillis(); @@ -339,6 +404,10 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener result) { log.info("[{}] Detected update of the API state for {}: {}", state.getEntityId(), state.getEntityType(), result); apiUsageStateService.update(state.getApiUsageState()); @@ -473,7 +542,12 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener addedPartitions) { + if (addedPartitions.isEmpty()) { + return; + } + try { log.info("Initializing tenant states."); updateLock.lock(); try { - ExecutorService tmpInitExecutor = ThingsBoardExecutors.newWorkStealingPool(20, "init-tenant-states-from-db"); - try { - PageDataIterable tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); - List> futures = new ArrayList<>(); - for (Tenant tenant : tenantIterator) { - if (!myUsageStates.containsKey(tenant.getId()) && partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()).isMyPartition()) { + PageDataIterable tenantIterator = new PageDataIterable<>(tenantService::findTenants, 1024); + List> futures = new ArrayList<>(); + for (Tenant tenant : tenantIterator) { + TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), tenant.getId()); + if (addedPartitions.contains(tpi)) { + if (!myUsageStates.containsKey(tenant.getId()) && tpi.isMyPartition()) { log.debug("[{}] Initializing tenant state.", tenant.getId()); - futures.add(tmpInitExecutor.submit(() -> { + futures.add(dbExecutor.submit(() -> { try { updateTenantState((TenantApiUsageState) getOrFetchState(tenant.getId(), tenant.getId()), tenantProfileCache.get(tenant.getTenantProfileId())); log.debug("[{}] Initialized tenant state.", tenant.getId()); } catch (Exception e) { log.warn("[{}] Failed to initialize tenant API state", tenant.getId(), e); } + return null; })); } + } else { + log.debug("[{}][{}] Tenant doesn't belong to current partition. tpi [{}]", tenant.getName(), tenant.getId(), tpi); } - for (Future future : futures) { - future.get(); - } - } finally { - tmpInitExecutor.shutdownNow(); } + Futures.whenAllComplete(futures); } finally { updateLock.unlock(); } - log.info("Initialized tenant states."); + log.info("Initialized {} tenant states.", myUsageStates.size()); } catch (Exception e) { log.warn("Unknown failure", e); } @@ -524,5 +601,8 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener { if (!myPartitions.containsKey(serviceQueueKey)) { log.info("[{}] NO MORE PARTITIONS FOR CURRENT KEY", serviceQueueKey); @@ -174,7 +176,6 @@ public class HashPartitionService implements PartitionService { applicationEventPublisher.publishEvent(new PartitionChangeEvent(this, serviceQueueKey, tpiList)); } }); - tpiCache.clear(); if (currentOtherServices == null) { currentOtherServices = new ArrayList<>(otherServices); From b2e063ba1e0ca214e97b7286add62e0cd4620764 Mon Sep 17 00:00:00 2001 From: YevhenBondarenko Date: Wed, 16 Feb 2022 17:01:22 +0200 Subject: [PATCH 2/7] added tests for DefaultTbApiUsageStateService --- .../DefaultTbApiUsageStateService.java | 10 +- .../DefaultTbApiUsageStateServiceTest.java | 101 ++++++++++++++++++ 2 files changed, 106 insertions(+), 5 deletions(-) create mode 100644 application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index db85f3c852..c22ced2387 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -120,13 +120,13 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener myUsageStates = new ConcurrentHashMap<>(); + final Map myUsageStates = new ConcurrentHashMap<>(); // Entities that should be processed on other servers - private final Map otherUsageStates = new ConcurrentHashMap<>(); + final Map otherUsageStates = new ConcurrentHashMap<>(); - private final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); + final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); - private final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); + final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); @Value("${usage.stats.report.enabled:true}") private boolean enabled; @@ -489,7 +489,7 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener Date: Fri, 18 Feb 2022 12:44:22 +0200 Subject: [PATCH 3/7] Correct executor for the service --- .../apiusage/DefaultTbApiUsageStateService.java | 17 +++++++---------- 1 file changed, 7 insertions(+), 10 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index c22ced2387..3e60c4d49c 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -110,7 +110,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener> subscribeQueue = new ConcurrentLinkedQueue<>(); @@ -147,7 +146,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener partitions = getLatestPartitionsFromQueue(); if (partitions == null) { - log.info("Tenant state service. Nothing to do. partitions is null"); + log.info("Api Usage state service. Nothing to do. Partitions are empty"); return; } initStateFromDB(partitions); @@ -601,8 +598,8 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener Date: Fri, 18 Feb 2022 13:57:59 +0200 Subject: [PATCH 4/7] Issue 6056. Refactoring of duplicated code to abstract class --- .../DefaultTbApiUsageStateService.java | 87 +++------- .../AbstractPartitionBasedService.java | 148 ++++++++++++++++++ .../state/DefaultDeviceStateService.java | 134 ++++------------ .../DefaultTbApiUsageStateServiceTest.java | 6 +- 4 files changed, 199 insertions(+), 176 deletions(-) create mode 100644 application/src/main/java/org/thingsboard/server/service/partition/AbstractPartitionBasedService.java diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index 3e60c4d49c..f39477b543 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -67,6 +67,7 @@ import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.scheduler.SchedulerComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.InternalTelemetryService; import javax.annotation.PostConstruct; @@ -93,7 +94,7 @@ import java.util.stream.Collectors; @Slf4j @Service -public class DefaultTbApiUsageStateService extends TbApplicationEventListener implements TbApiUsageStateService { +public class DefaultTbApiUsageStateService extends AbstractPartitionBasedService implements TbApiUsageStateService { public static final String HOURLY = "Hourly"; public static final FutureCallback VOID_CALLBACK = new FutureCallback() { @@ -123,8 +124,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener otherUsageStates = new ConcurrentHashMap<>(); - final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); - final Set deletedEntities = Collections.newSetFromMap(new ConcurrentHashMap<>()); @Value("${usage.stats.report.enabled:true}") @@ -137,10 +136,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener> subscribeQueue = new ConcurrentLinkedQueue<>(); - public DefaultTbApiUsageStateService(TbClusterService clusterService, PartitionService partitionService, TenantService tenantService, @@ -162,7 +157,7 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener msg, TbCallback callback) { ToUsageStatsServiceMsg statsMsg = msg.getValue(); @@ -226,59 +226,6 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener partitions = getLatestPartitionsFromQueue(); - if (partitions == null) { - log.info("Api Usage state service. Nothing to do. Partitions are empty"); - return; - } - initStateFromDB(partitions); - } - - Set getLatestPartitionsFromQueue() { - log.debug("getLatestPartitionsFromQueue, queue size {}", subscribeQueue.size()); - Set partitions = null; - while (!subscribeQueue.isEmpty()) { - partitions = subscribeQueue.poll(); - log.debug("polled from the queue partitions {}", partitions); - } - log.debug("getLatestPartitionsFromQueue, partitions {}", partitions); - return partitions; - } - - private void initStateFromDB(Set partitions) { - try { - Set addedPartitions = new HashSet<>(partitions); - addedPartitions.removeAll(partitionedEntities.keySet()); - - Set removedPartitions = new HashSet<>(partitionedEntities.keySet()); - removedPartitions.removeAll(partitions); - - removedPartitions.forEach(partition -> { - Set entities = partitionedEntities.remove(partition); - entities.forEach(this::cleanUpEntitiesStateMap); - }); - - addedPartitions.forEach(tpi -> - partitionedEntities.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); - - otherUsageStates.entrySet().removeIf(entry -> - partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition()); - - initStatesFromDataBase(addedPartitions); - } catch (Throwable t) { - log.warn("Failed to init tenant states from DB", t); - } - } - @Override public ApiUsageState getApiUsageState(TenantId tenantId) { TenantApiUsageState tenantState = (TenantApiUsageState) myUsageStates.get(tenantId); @@ -401,7 +348,8 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener addedPartitions) { - if (addedPartitions.isEmpty()) { - return; - } + @Override + protected void onRepartitionEvent() { + otherUsageStates.entrySet().removeIf(entry -> + partitionService.resolve(ServiceType.TB_CORE, entry.getValue().getTenantId(), entry.getKey()).isMyPartition()); + } + @Override + protected void onAddedPartitions(Set addedPartitions) { try { log.info("Initializing tenant states."); updateLock.lock(); @@ -595,11 +546,9 @@ public class DefaultTbApiUsageStateService extends TbApplicationEventListener extends TbApplicationEventListener { + + protected final ConcurrentMap> partitionedEntities = new ConcurrentHashMap<>(); + final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); + + protected ListeningScheduledExecutorService scheduledExecutor; + + abstract protected String getSchedulerExecutorName(); + + abstract protected void onAddedPartitions(Set addedPartitions); + + abstract protected void cleanupEntityOnPartitionRemoval(T entityId); + + public Set getPartitionedEntities(TopicPartitionInfo tpi) { + return partitionedEntities.get(tpi); + } + + protected void init() { + // Should be always single threaded due to absence of locks. + scheduledExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-state-scheduled"))); + } + + protected ServiceType getServiceType() { + return ServiceType.TB_CORE; + } + + protected void stop() { + if (scheduledExecutor != null) { + scheduledExecutor.shutdownNow(); + } + } + + /** + * DiscoveryService will call this event from the single thread (one-by-one). + * Events order is guaranteed by DiscoveryService. + * The only concurrency is expected from the [main] thread on Application started. + * Async implementation. Locks is not allowed by design. + * Any locks or delays in this module will affect DiscoveryService and entire system + */ + @Override + protected void onTbApplicationEvent(PartitionChangeEvent partitionChangeEvent) { + if (getServiceType().equals(partitionChangeEvent.getServiceType())) { + log.debug("onTbApplicationEvent, processing event: {}", partitionChangeEvent); + subscribeQueue.add(partitionChangeEvent.getPartitions()); + scheduledExecutor.submit(this::pollInitStateFromDB); + } + } + + protected void pollInitStateFromDB() { + final Set partitions = getLatestPartitions(); + if (partitions == null) { + log.debug("Nothing to do. Partitions are empty."); + return; + } + initStateFromDB(partitions); + } + + private void initStateFromDB(Set partitions) { + try { + log.info("CURRENT PARTITIONS: {}", partitionedEntities.keySet()); + log.info("NEW PARTITIONS: {}", partitions); + + Set addedPartitions = new HashSet<>(partitions); + addedPartitions.removeAll(partitionedEntities.keySet()); + + log.info("ADDED PARTITIONS: {}", addedPartitions); + + Set removedPartitions = new HashSet<>(partitionedEntities.keySet()); + removedPartitions.removeAll(partitions); + + log.info("REMOVED PARTITIONS: {}", removedPartitions); + + // We no longer manage current partition of entities; + removedPartitions.forEach(partition -> { + Set entities = partitionedEntities.remove(partition); + entities.forEach(this::cleanupEntityOnPartitionRemoval); + }); + + onRepartitionEvent(); + + addedPartitions.forEach(tpi -> partitionedEntities.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); + + if (!addedPartitions.isEmpty()) { + onAddedPartitions(addedPartitions); + } + + scheduledExecutor.submit(() -> { + log.info("Managing following partitions:"); + partitionedEntities.forEach((tpi, entities) -> { + log.info("[{}]: {} entities", tpi.getFullTopicName(), entities.size()); + }); + }); + } catch (Throwable t) { + log.warn("Failed to init entities state from DB", t); + } + } + + protected void onRepartitionEvent() { + } + + private Set getLatestPartitions() { + log.debug("getLatestPartitionsFromQueue, queue size {}", subscribeQueue.size()); + Set partitions = null; + while (!subscribeQueue.isEmpty()) { + partitions = subscribeQueue.poll(); + log.debug("polled from the queue partitions {}", partitions); + } + log.debug("getLatestPartitionsFromQueue, partitions {}", partitions); + return partitions; + } + +} 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 8a24cdbf1a..1d6b0dcb4b 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 @@ -59,6 +59,7 @@ import org.thingsboard.server.queue.discovery.PartitionService; import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.cluster.TbClusterService; +import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; import javax.annotation.Nonnull; @@ -94,7 +95,7 @@ import static org.thingsboard.server.common.data.DataConstants.SERVER_SCOPE; @Service @TbCoreComponent @Slf4j -public class DefaultDeviceStateService extends TbApplicationEventListener implements DeviceStateService { +public class DefaultDeviceStateService extends AbstractPartitionBasedService implements DeviceStateService { public static final String ACTIVITY_STATE = "active"; public static final String LAST_CONNECT_TIME = "lastConnectTime"; @@ -131,12 +132,9 @@ public class DefaultDeviceStateService extends TbApplicationEventListener> partitionedDevices = new ConcurrentHashMap<>(); - final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); - final Queue> subscribeQueue = new ConcurrentLinkedQueue<>(); + final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, AttributesService attributesService, TimeseriesService tsService, @@ -156,21 +154,23 @@ public class DefaultDeviceStateService extends TbApplicationEventListener() { + Futures.addCallback(fetchDeviceState(device), new FutureCallback<>() { @Override public void onSuccess(@Nullable DeviceStateData state) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenantId, device.getId()); - if (partitionedDevices.containsKey(tpi)) { + if (partitionedEntities.containsKey(tpi)) { addDeviceUsingState(tpi, state); save(deviceId, ACTIVITY_STATE, false); callback.onSuccess(); @@ -286,94 +286,14 @@ public class DefaultDeviceStateService extends TbApplicationEventListener partitions = getLatestPartitionsFromQueue(); - if (partitions == null) { - log.info("Device state service. Nothing to do. partitions is null"); - return; - } - initStateFromDB(partitions); - } - - // TODO: move to utils - Set getLatestPartitionsFromQueue() { - log.debug("getLatestPartitionsFromQueue, queue size {}", subscribeQueue.size()); - Set partitions = null; - while (!subscribeQueue.isEmpty()) { - partitions = subscribeQueue.poll(); - log.debug("polled from the queue partitions {}", partitions); - } - log.debug("getLatestPartitionsFromQueue, partitions {}", partitions); - return partitions; - } - - private void initStateFromDB(Set partitions) { - try { - log.info("CURRENT PARTITIONS: {}", partitionedDevices.keySet()); - log.info("NEW PARTITIONS: {}", partitions); - - Set addedPartitions = new HashSet<>(partitions); - addedPartitions.removeAll(partitionedDevices.keySet()); - - log.info("ADDED PARTITIONS: {}", addedPartitions); - - Set removedPartitions = new HashSet<>(partitionedDevices.keySet()); - removedPartitions.removeAll(partitions); - - log.info("REMOVED PARTITIONS: {}", removedPartitions); - - // We no longer manage current partition of devices; - removedPartitions.forEach(partition -> { - Set devices = partitionedDevices.remove(partition); - devices.forEach(this::cleanUpDeviceStateMap); - }); - - addedPartitions.forEach(tpi -> partitionedDevices.computeIfAbsent(tpi, key -> ConcurrentHashMap.newKeySet())); - - initPartitions(addedPartitions); - - scheduledExecutor.submit(() -> { - log.info("Managing following partitions:"); - partitionedDevices.forEach((tpi, devices) -> { - log.info("[{}]: {} devices", tpi.getFullTopicName(), devices.size()); - }); - }); - } catch (Throwable t) { - log.warn("Failed to init device states from DB", t); - } - } - - //TODO 3.0: replace this dummy search with new functionality to search by partitions using SQL capabilities. - //Adding only entities that are in new partitions - boolean initPartitions(Set addedPartitions) { - if (addedPartitions.isEmpty()) { - return false; - } - + protected void onAddedPartitions(Set addedPartitions) { List tenants = tenantService.findTenants(new PageLink(Integer.MAX_VALUE)).getData(); for (Tenant tenant : tenants) { log.debug("Finding devices for tenant [{}]", tenant.getName()); final PageLink pageLink = new PageLink(initFetchPackSize); scheduledExecutor.submit(() -> processPageAndSubmitNextPage(addedPartitions, tenant, pageLink, scheduledExecutor)); - } - return true; } private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink, final ExecutorService executor) { @@ -435,7 +355,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIds = partitionedDevices.get(tpi); + Set deviceIds = partitionedEntities.get(tpi); if (deviceIds != null) { deviceIds.add(state.getDeviceId()); deviceStates.put(state.getDeviceId(), state); @@ -447,7 +367,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener { + partitionedEntities.forEach((tpi, deviceIds) -> { log.debug("Calculating state updates. tpi {} for {} devices", tpi.getFullTopicName(), deviceIds.size()); for (DeviceId deviceId : deviceIds) { updateInactivityStateIfExpired(ts, deviceId); @@ -473,7 +393,7 @@ public class DefaultDeviceStateService extends TbApplicationEventListener deviceIdSet = partitionedDevices.get(tpi); - deviceIdSet.remove(deviceId); + Set deviceIdSet = partitionedEntities.get(tpi); + if (deviceIdSet != null) { + deviceIdSet.remove(deviceId); + } } - private void cleanUpDeviceStateMap(DeviceId deviceId) { + @Override + protected void cleanupEntityOnPartitionRemoval(DeviceId deviceId) { + cleanupEntity(deviceId); + } + + private void cleanupEntity(DeviceId deviceId) { deviceStates.remove(deviceId); } + private ListenableFuture fetchDeviceState(Device device) { ListenableFuture future; if (persistToTelemetry) { diff --git a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java index b446e4adc0..06146d0982 100644 --- a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java @@ -60,8 +60,6 @@ public class DefaultTbApiUsageStateServiceTest { @Mock ApiUsageStateService apiUsageStateService; @Mock - SchedulerComponent scheduler; - @Mock TbTenantProfileCache tenantProfileCache; @Mock MailService mailService; @@ -74,7 +72,7 @@ public class DefaultTbApiUsageStateServiceTest { @Before public void setUp() { - service = spy(new DefaultTbApiUsageStateService(clusterService, partitionService, tenantService, tsService, apiUsageStateService, scheduler, tenantProfileCache, mailService, dbExecutor)); + service = spy(new DefaultTbApiUsageStateService(clusterService, partitionService, tenantService, tsService, apiUsageStateService, tenantProfileCache, mailService, dbExecutor)); } @Test @@ -94,7 +92,7 @@ public class DefaultTbApiUsageStateServiceTest { willReturn(tenantUsageStateMock).given(service).getOrFetchState(tenantId, tenantId); ApiUsageState tenantUsageState = service.getApiUsageState(tenantId); assertThat(tenantUsageState, is(tenantUsageStateMock.getApiUsageState())); - assertThat(true, is(service.partitionedEntities.get(tpi).contains(tenantId))); + assertThat(true, is(service.getPartitionedEntities(tpi).contains(tenantId))); Mockito.verify(service, times(1)).getOrFetchState(tenantId, tenantId); } From 45df14e39e2ce1af5d104f50db90afcf9c5e3373 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 18 Feb 2022 16:33:51 +0200 Subject: [PATCH 5/7] Device State Service improvements and race condition fix --- .../AbstractPartitionBasedService.java | 8 +-- .../state/DefaultDeviceStateService.java | 56 ++++++++++++------- 2 files changed, 40 insertions(+), 24 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 632fb9938b..18a0a1f220 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 @@ -120,11 +120,9 @@ public abstract class AbstractPartitionBasedService extends onAddedPartitions(addedPartitions); } - scheduledExecutor.submit(() -> { - log.info("Managing following partitions:"); - partitionedEntities.forEach((tpi, entities) -> { - log.info("[{}]: {} entities", tpi.getFullTopicName(), entities.size()); - }); + log.info("Managing following partitions:"); + partitionedEntities.forEach((tpi, entities) -> { + log.info("[{}]: {} entities", tpi.getFullTopicName(), entities.size()); }); } catch (Throwable t) { log.warn("Failed to init entities state from DB", t); 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 1d6b0dcb4b..10f1fcea85 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 @@ -175,6 +175,9 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService 0 && lastReportedActivity > stateData.getState().getLastActivityTime()) { updateActivityState(deviceId, stateData, lastReportedActivity); } - cleanDeviceStateIfBelongsExternalPartition(tenantId, deviceId); } void updateActivityState(DeviceId deviceId, DeviceStateData stateData, long lastReportedActivity) { @@ -214,12 +219,14 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService processPageAndSubmitNextPage(addedPartitions, tenant, pageLink, scheduledExecutor)); + processPageAndSubmitNextPage(addedPartitions, tenant, pageLink); } } - private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink, final ExecutorService executor) { + private void processPageAndSubmitNextPage(final Set addedPartitions, final Tenant tenant, final PageLink pageLink) { 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()) { TopicPartitionInfo tpi = partitionService.resolve(ServiceType.TB_CORE, tenant.getId(), device.getId()); - if (addedPartitions.contains(tpi)) { + if (addedPartitions.contains(tpi) && !deviceStates.containsKey(device.getId())) { log.debug("[{}][{}] Device belong to current partition. tpi [{}]. Fetching state from DB", device.getName(), device.getId(), tpi); - ListenableFuture future = Futures.transform(fetchDeviceState(device), new Function() { + ListenableFuture future = Futures.transform(fetchDeviceState(device), new Function<>() { @Nullable @Override public Void apply(@Nullable DeviceStateData state) { @@ -323,7 +333,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService>() { + Futures.addCallback(Futures.successfulAsList(fetchFutures), new FutureCallback<>() { @Override public void onSuccess(List result) { log.trace("[{}] Success init device state from DB for batch size {}", tenant.getId(), result.size()); @@ -339,7 +349,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService processPageAndSubmitNextPage(addedPartitions, tenant, nextPageLink, executor)); + processPageAndSubmitNextPage(addedPartitions, tenant, nextPageLink); } } @@ -358,7 +368,7 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService deviceIds = partitionedEntities.get(tpi); if (deviceIds != null) { deviceIds.add(state.getDeviceId()); - deviceStates.put(state.getDeviceId(), state); + deviceStates.putIfAbsent(state.getDeviceId(), state); } else { log.debug("[{}] Device belongs to external partition {}", state.getDeviceId(), tpi.getFullTopicName()); throw new RuntimeException("Device belongs to external partition " + tpi.getFullTopicName() + "!"); @@ -384,12 +394,18 @@ public class DefaultDeviceStateService extends AbstractPartitionBasedService Date: Fri, 18 Feb 2022 16:40:12 +0200 Subject: [PATCH 6/7] Optimize imports --- .../apiusage/DefaultTbApiUsageStateService.java | 8 -------- .../service/state/DefaultDeviceStateService.java | 11 ++--------- 2 files changed, 2 insertions(+), 17 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java index f39477b543..57b10395fb 100644 --- a/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateService.java @@ -18,8 +18,6 @@ package org.thingsboard.server.service.apiusage; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningScheduledExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.StringUtils; import org.checkerframework.checker.nullness.qual.Nullable; @@ -63,9 +61,6 @@ import org.thingsboard.server.gen.transport.TransportProtos.ToUsageStatsServiceM import org.thingsboard.server.gen.transport.TransportProtos.UsageStatsKVProto; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.TbApplicationEventListener; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; -import org.thingsboard.server.queue.scheduler.SchedulerComponent; import org.thingsboard.server.service.executors.DbCallbackExecutorService; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.InternalTelemetryService; @@ -78,12 +73,9 @@ import java.util.Collections; import java.util.HashSet; import java.util.List; import java.util.Map; -import java.util.Queue; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentLinkedQueue; -import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; 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 10f1fcea85..8fc84de9e5 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 @@ -20,15 +20,15 @@ import com.google.common.base.Function; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningScheduledExecutorService; -import com.google.common.util.concurrent.MoreExecutors; import lombok.Getter; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; +import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.DataConstants; import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.Tenant; @@ -52,13 +52,9 @@ import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.device.DeviceService; import org.thingsboard.server.dao.tenant.TenantService; import org.thingsboard.server.dao.timeseries.TimeseriesService; -import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.gen.transport.TransportProtos; -import org.thingsboard.server.queue.discovery.event.PartitionChangeEvent; import org.thingsboard.server.queue.discovery.PartitionService; -import org.thingsboard.server.queue.discovery.TbApplicationEventListener; import org.thingsboard.server.queue.util.TbCoreComponent; -import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.service.partition.AbstractPartitionBasedService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; @@ -69,14 +65,11 @@ import javax.annotation.PreDestroy; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; -import java.util.HashSet; import java.util.List; -import java.util.Queue; import java.util.Random; import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentLinkedQueue; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; From 9ca1055d54a794806a974b8ad25e3a7200c95026 Mon Sep 17 00:00:00 2001 From: Andrii Shvaika Date: Fri, 18 Feb 2022 17:04:52 +0200 Subject: [PATCH 7/7] Remove incomplete test --- .../apiusage/DefaultTbApiUsageStateServiceTest.java | 13 ------------- 1 file changed, 13 deletions(-) diff --git a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java index 06146d0982..fde595ea3b 100644 --- a/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/apiusage/DefaultTbApiUsageStateServiceTest.java @@ -83,17 +83,4 @@ public class DefaultTbApiUsageStateServiceTest { Mockito.verify(service, never()).getOrFetchState(tenantId, tenantId); } - @Test - public void givenTenantIdWithoutTenantStateInMap_whenGetState_thenGetOrFetchState() { - TopicPartitionInfo tpi = Mockito.mock(TopicPartitionInfo.class); - Mockito.when(tpi.isMyPartition()).thenReturn(true); - willReturn(tpi).given(partitionService).resolve(ServiceType.TB_CORE, tenantId, tenantId); - service.myUsageStates.clear(); - willReturn(tenantUsageStateMock).given(service).getOrFetchState(tenantId, tenantId); - ApiUsageState tenantUsageState = service.getApiUsageState(tenantId); - assertThat(tenantUsageState, is(tenantUsageStateMock.getApiUsageState())); - assertThat(true, is(service.getPartitionedEntities(tpi).contains(tenantId))); - Mockito.verify(service, times(1)).getOrFetchState(tenantId, tenantId); - } - } \ No newline at end of file