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 75577fc5a4..b1d47c9864 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.util.TbCoreComponent; import org.thingsboard.server.service.queue.TbClusterService; import org.thingsboard.server.service.telemetry.TelemetrySubscriptionService; +import org.thingsboard.server.utils.EventDeduplicationExecutor; import javax.annotation.Nullable; import javax.annotation.PostConstruct; @@ -126,13 +127,13 @@ public class DefaultDeviceStateService implements DeviceStateService { @Getter private int initFetchPackSize; - private volatile boolean clusterUpdatePending = false; - private ListeningScheduledExecutorService queueExecutor; private final ConcurrentMap> partitionedDevices = new ConcurrentHashMap<>(); private final ConcurrentMap deviceStates = new ConcurrentHashMap<>(); private final ConcurrentMap deviceLastReportedActivity = new ConcurrentHashMap<>(); private final ConcurrentMap deviceLastSavedActivity = new ConcurrentHashMap<>(); + private volatile EventDeduplicationExecutor> deduplicationExecutor; + public DefaultDeviceStateService(TenantService tenantService, DeviceService deviceService, AttributesService attributesService, TimeseriesService tsService, @@ -155,6 +156,7 @@ public class DefaultDeviceStateService implements DeviceStateService { // Should be always single threaded due to absence of locks. queueExecutor = MoreExecutors.listeningDecorator(Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("device-state"))); queueExecutor.scheduleAtFixedRate(this::updateState, new Random().nextInt(defaultStateCheckIntervalInSec), defaultStateCheckIntervalInSec, TimeUnit.SECONDS); + deduplicationExecutor = new EventDeduplicationExecutor<>(DefaultDeviceStateService.class.getSimpleName(), queueExecutor, this::initStateFromDB); } @PreDestroy @@ -292,25 +294,14 @@ public class DefaultDeviceStateService implements DeviceStateService { } } - volatile Set pendingPartitions; - @Override public void onApplicationEvent(PartitionChangeEvent partitionChangeEvent) { if (ServiceType.TB_CORE.equals(partitionChangeEvent.getServiceType())) { - synchronized (this) { - pendingPartitions = partitionChangeEvent.getPartitions(); - if (!clusterUpdatePending) { - clusterUpdatePending = true; - queueExecutor.submit(() -> { - clusterUpdatePending = false; - initStateFromDB(); - }); - } - } + deduplicationExecutor.submit(partitionChangeEvent.getPartitions()); } } - private void initStateFromDB() { + private void initStateFromDB(Set pendingPartitions) { try { log.info("CURRENT PARTITIONS: {}", partitionedDevices.keySet()); log.info("NEW PARTITIONS: {}", pendingPartitions); diff --git a/application/src/main/java/org/thingsboard/server/utils/EventDeduplicationExecutor.java b/application/src/main/java/org/thingsboard/server/utils/EventDeduplicationExecutor.java new file mode 100644 index 0000000000..1c67868740 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/utils/EventDeduplicationExecutor.java @@ -0,0 +1,85 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.utils; + +import lombok.extern.slf4j.Slf4j; + +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.function.Consumer; + +/** + * This class deduplicate executions of the specified function. + * Useful in cluster mode, when you get event about partition change multiple times. + * Assuming that the function execution is expensive, we should execute it immediately when first time event occurs and + * later, once the processing of first event is done, process last pending task. + * + * @param

parameters of the function + */ +@Slf4j +public class EventDeduplicationExecutor

{ + private final String name; + private final ExecutorService executor; + private final Consumer

function; + private P pendingTask; + private boolean busy; + + public EventDeduplicationExecutor(String name, ExecutorService executor, Consumer

function) { + this.name = name; + this.executor = executor; + this.function = function; + } + + public void submit(P params) { + log.info("[{}] Going to submit: {}", name, params); + synchronized (EventDeduplicationExecutor.this) { + if (!busy) { + busy = true; + pendingTask = null; + try { + log.info("[{}] Submitting task: {}", name, params); + executor.submit(() -> { + try { + log.info("[{}] Executing task: {}", name, params); + function.accept(params); + } catch (Throwable e) { + log.warn("Failed to process task with parameters: {}", params, e); + throw e; + } finally { + unlockAndProcessIfAny(); + } + }); + } catch (Throwable e) { + log.warn("Failed to submit task with parameters: {}", params, e); + unlockAndProcessIfAny(); + throw e; + } + } else { + log.info("[{}] Task is already in progress. {} pending task: {}", name, pendingTask == null ? "adding" : "updating", params); + pendingTask = params; + } + } + } + + private void unlockAndProcessIfAny() { + synchronized (EventDeduplicationExecutor.this) { + busy = false; + if (pendingTask != null) { + submit(pendingTask); + } + } + } +} diff --git a/application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java b/application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java new file mode 100644 index 0000000000..26718b6624 --- /dev/null +++ b/application/src/test/java/org/thingsboard/server/util/EventDeduplicationExecutorTest.java @@ -0,0 +1,119 @@ +/** + * Copyright © 2016-2021 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.util; + +import com.google.common.util.concurrent.MoreExecutors; +import lombok.extern.slf4j.Slf4j; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.mockito.Mockito; +import org.mockito.runners.MockitoJUnitRunner; +import org.thingsboard.server.utils.EventDeduplicationExecutor; + +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.function.Consumer; + +@Slf4j +@RunWith(MockitoJUnitRunner.class) +public class EventDeduplicationExecutorTest { + + @Test + public void testSimpleFlowSameThread() throws InterruptedException { + simpleFlow(MoreExecutors.newDirectExecutorService()); + } + + @Test + public void testPeriodicFlowSameThread() throws InterruptedException { + periodicFlow(MoreExecutors.newDirectExecutorService()); + } + + + @Test + public void testSimpleFlowSingleThread() throws InterruptedException { + simpleFlow(Executors.newFixedThreadPool(1)); + } + + @Test + public void testPeriodicFlowSingleThread() throws InterruptedException { + periodicFlow(Executors.newFixedThreadPool(1)); + } + + @Test + public void testSimpleFlowMultiThread() throws InterruptedException { + simpleFlow(Executors.newFixedThreadPool(3)); + } + + @Test + public void testPeriodicFlowMultiThread() throws InterruptedException { + periodicFlow(Executors.newFixedThreadPool(3)); + } + + private void simpleFlow(ExecutorService executorService) throws InterruptedException { + try { + Consumer function = Mockito.spy(StringConsumer.class); + EventDeduplicationExecutor executor = new EventDeduplicationExecutor<>(EventDeduplicationExecutorTest.class.getSimpleName(), executorService, function); + + String params1 = "params1"; + String params2 = "params2"; + String params3 = "params3"; + + executor.submit(params1); + executor.submit(params2); + executor.submit(params3); + Thread.sleep(500); + Mockito.verify(function).accept(params1); + Mockito.verify(function).accept(params3); + } finally { + executorService.shutdownNow(); + } + } + + private void periodicFlow(ExecutorService executorService) throws InterruptedException { + try { + Consumer function = Mockito.spy(StringConsumer.class); + EventDeduplicationExecutor executor = new EventDeduplicationExecutor<>(EventDeduplicationExecutorTest.class.getSimpleName(), executorService, function); + + String params1 = "params1"; + String params2 = "params2"; + String params3 = "params3"; + + executor.submit(params1); + Thread.sleep(500); + executor.submit(params2); + Thread.sleep(500); + executor.submit(params3); + Thread.sleep(500); + Mockito.verify(function).accept(params1); + Mockito.verify(function).accept(params2); + Mockito.verify(function).accept(params3); + } finally { + executorService.shutdownNow(); + } + } + + public static class StringConsumer implements Consumer { + @Override + public void accept(String s) { + try { + Thread.sleep(100); + } catch (InterruptedException e) { + throw new RuntimeException(e); + } + } + } + +}