diff --git a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java index 64139530b1..3f05a5b923 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/DefaultTbRuleEngineConsumerService.java @@ -126,6 +126,7 @@ public class DefaultTbRuleEngineConsumerService extends AbstractConsumerService< @Override protected void stopConsumers() { consumers.values().forEach(TbRuleEngineQueueConsumerManager::stop); + consumers.values().forEach(TbRuleEngineQueueConsumerManager::awaitStop); ctx.stop(); } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/QueueEvent.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/QueueEvent.java index a4a3bcc59a..86955be7e7 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/QueueEvent.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/QueueEvent.java @@ -19,6 +19,6 @@ import java.io.Serializable; public enum QueueEvent implements Serializable { - PARTITION_CHANGE, CONFIG_UPDATE, STOP, DELETE + PARTITION_CHANGE, CONFIG_UPDATE, DELETE } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java index a0a5579951..2e3a4676e8 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbQueueConsumerTask.java @@ -15,7 +15,8 @@ */ package org.thingsboard.server.service.queue.ruleengine; -import lombok.Data; +import lombok.Getter; +import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos; @@ -23,16 +24,26 @@ import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; import java.util.Set; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; -@Data +@RequiredArgsConstructor @Slf4j public class TbQueueConsumerTask { + @Getter private final Object key; + @Getter private final TbQueueConsumer> consumer; - private volatile Future task; + + private Future task; + private CountDownLatch completionLatch; + + public void setTask(Future task) { + this.task = task; + this.completionLatch = new CountDownLatch(1); + } public void subscribe(Set partitions) { log.trace("[{}] Subscribing to partitions: {}", key, partitions); @@ -42,23 +53,32 @@ public class TbQueueConsumerTask { public void initiateStop() { log.debug("[{}] Initiating stop", key); consumer.stop(); + if (isRunning()) { + task.cancel(true); + } } - public void awaitFinish() { + public void awaitCompletion() { log.trace("[{}] Awaiting finish", key); if (isRunning()) { try { - task.get(60, TimeUnit.SECONDS); - task = null; + if (!completionLatch.await(30, TimeUnit.SECONDS)) { + throw new IllegalStateException("timeout of 30 seconds expired"); + } } catch (Exception e) { log.warn("[{}] Failed to await for consumer to stop", key, e); } + task = null; } log.trace("[{}] Awaited finish", key); } public boolean isRunning() { - return task != null && !task.isDone(); + return task != null; + } + + public void finished() { + completionLatch.countDown(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineConsumerContext.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineConsumerContext.java index 325a685972..f4c8e4c2ae 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineConsumerContext.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineConsumerContext.java @@ -36,7 +36,6 @@ import javax.annotation.PostConstruct; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; @Component @TbRuleEngineComponent @@ -83,12 +82,7 @@ public class TbRuleEngineConsumerContext { public void stop() { scheduler.shutdownNow(); - consumersExecutor.shutdown(); - mgmtExecutor.shutdown(); - try { - mgmtExecutor.awaitTermination(15, TimeUnit.SECONDS); - } catch (InterruptedException e) { - log.warn("Failed to await mgmtExecutor termination"); - } + consumersExecutor.shutdownNow(); + mgmtExecutor.shutdownNow(); } } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java index dc3966d0d2..3a7b2f15d8 100644 --- a/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java +++ b/application/src/main/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManager.java @@ -85,7 +85,13 @@ public class TbRuleEngineQueueConsumerManager { } public void init(Queue queue) { - doInit(queue); + this.queue = queue; + if (queue.isConsumerPerPartition()) { + this.consumerWrapper = new ConsumerPerPartitionWrapper(); + } else { + this.consumerWrapper = new SingleConsumerWrapper(); + } + log.debug("[{}] Initialized consumer for queue: {}", queueKey, queue); } public void update(Queue queue) { @@ -96,10 +102,6 @@ public class TbRuleEngineQueueConsumerManager { addTask(new TbQueueConsumerManagerTask(QueueEvent.PARTITION_CHANGE, partitions)); } - public void stop() { - addTask(new TbQueueConsumerManagerTask(QueueEvent.STOP)); - } - public void delete() { addTask(new TbQueueConsumerManagerTask(QueueEvent.DELETE)); } @@ -135,20 +137,22 @@ public class TbRuleEngineQueueConsumerManager { newPartitions = task.getPartitions(); } else if (task.getEvent() == QueueEvent.CONFIG_UPDATE) { newConfiguration = task.getQueue(); - } else if (task.getEvent() == QueueEvent.STOP) { - doStop(); - return; } else if (task.getEvent() == QueueEvent.DELETE) { doDelete(); return; } } + if (stopped) { + return; + } if (newConfiguration != null) { doUpdate(newConfiguration); } if (newPartitions != null) { doUpdate(newPartitions); } + } catch (Exception e) { + log.error("[{}] Failed to process tasks", queueKey, e); } finally { lock.unlock(); } @@ -159,16 +163,6 @@ public class TbRuleEngineQueueConsumerManager { }); } - private void doInit(Queue queue) { - this.queue = queue; - if (queue.isConsumerPerPartition()) { - consumerWrapper = new ConsumerPerPartitionWrapper(); - } else { - consumerWrapper = new SingleConsumerWrapper(); - } - log.debug("[{}] Initialized consumer for queue: {}", queueKey, queue); - } - private void doUpdate(Queue newQueue) { log.info("[{}] Processing queue update: {}", queueKey, newQueue); var oldQueue = this.queue; @@ -179,12 +173,12 @@ public class TbRuleEngineQueueConsumerManager { } if (oldQueue == null) { - doInit(queue); + init(queue); } else if (newQueue.isConsumerPerPartition() != oldQueue.isConsumerPerPartition()) { consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::initiateStop); - consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitFinish); + consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitCompletion); - doInit(queue); + init(queue); if (partitions != null) { doUpdate(partitions); // even if partitions number was changed, there can be no partition change event } @@ -200,17 +194,21 @@ public class TbRuleEngineQueueConsumerManager { consumerWrapper.updatePartitions(partitions); } - private void doStop() { + public void stop() { + log.debug("[{}] Stopping consumers", queueKey); consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::initiateStop); - consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitFinish); - log.debug("[{}] Unsubscribed and stopped consumers", queueKey); stopped = true; } + public void awaitStop() { + consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitCompletion); + log.debug("[{}] Unsubscribed and stopped consumers", queueKey); + } + private void doDelete() { stopped = true; log.info("[{}] Handling queue deletion", queueKey); - consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitFinish); + consumerWrapper.getConsumers().forEach(TbQueueConsumerTask::awaitCompletion); List>> queueConsumers = consumerWrapper.getConsumers().stream() .map(TbQueueConsumerTask::getConsumer).collect(Collectors.toList()); @@ -240,6 +238,7 @@ public class TbRuleEngineQueueConsumerManager { Future consumerLoop = ctx.getConsumersExecutor().submit(() -> { ThingsBoardThreadFactory.updateCurrentThreadName(consumerTask.getKey().toString()); consumerLoop(consumerTask.getConsumer()); + consumerTask.finished(); }); consumerTask.setTask(consumerLoop); } @@ -253,7 +252,7 @@ public class TbRuleEngineQueueConsumerManager { } processMsgs(msgs, consumer, queue); } catch (Exception e) { - if (!stopped) { + if (!consumer.isStopped()) { log.warn("Failed to process messages from queue", e); try { Thread.sleep(ctx.getPollDuration()); @@ -279,7 +278,7 @@ public class TbRuleEngineQueueConsumerManager { TbMsgPackProcessingContext packCtx = new TbMsgPackProcessingContext(queue.getName(), submitStrategy, ackStrategy.isSkipTimeoutMsgs()); submitStrategy.submitAttempt((id, msg) -> submitMessage(packCtx, id, msg)); - final boolean timeout = !packCtx.await(queue.getPackProcessingTimeout(), TimeUnit.MILLISECONDS); + final boolean timeout = !awaitPackProcessing(packCtx, queue.getPackProcessingTimeout(), true); TbRuleEngineProcessingResult result = new TbRuleEngineProcessingResult(queue.getName(), timeout, packCtx); if (timeout) { @@ -307,6 +306,19 @@ public class TbRuleEngineQueueConsumerManager { } } + private boolean awaitPackProcessing(TbMsgPackProcessingContext packCtx, long processingTimeout, boolean ignoreInterrupt) throws InterruptedException { + try { + return packCtx.await(processingTimeout, TimeUnit.MILLISECONDS); + } catch (InterruptedException e) { + if (ignoreInterrupt) { + log.debug("Interrupt happened while waiting for pack processing, trying to await one more time"); + return awaitPackProcessing(packCtx, processingTimeout, false); + } else { + throw new RuntimeException("Failed to await pack processing due to thread interrupt", e); + } + } + } + private TbRuleEngineSubmitStrategy getSubmitStrategy(Queue queue) { return ctx.getSubmitStrategyFactory().newInstance(queue.getName(), queue.getSubmitStrategy()); } @@ -430,7 +442,7 @@ public class TbRuleEngineQueueConsumerManager { consumers.get(tpi).initiateStop(); }); removedPartitions.forEach((tpi) -> { - consumers.remove(tpi).awaitFinish(); + consumers.remove(tpi).awaitCompletion(); }); addedPartitions.forEach((tpi) -> { @@ -457,7 +469,7 @@ public class TbRuleEngineQueueConsumerManager { if (partitions.isEmpty()) { if (consumer != null && consumer.isRunning()) { consumer.initiateStop(); - consumer.awaitFinish(); + consumer.awaitCompletion(); } consumer = null; return; diff --git a/application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java b/application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java index 276e65ec39..20b765da6d 100644 --- a/application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java +++ b/application/src/test/java/org/thingsboard/server/service/queue/ruleengine/TbRuleEngineQueueConsumerManagerTest.java @@ -180,13 +180,16 @@ public class TbRuleEngineQueueConsumerManagerTest { @After public void afterEach() { consumerManager.stop(); + consumerManager.awaitStop(); ruleEngineConsumerContext.stop(); if (generateQueueMsgs) { - await().atMost(1, TimeUnit.SECONDS) - .until(() -> totalProcessedMsgs.get() == totalConsumedMsgs.get()); + await().atMost(2, TimeUnit.SECONDS) + .untilAsserted(() -> { + log.debug("totalConsumedMsgs = {}, totalProcessedMsgs = {}", totalConsumedMsgs.get(), totalProcessedMsgs.get()); + assertThat(totalProcessedMsgs.get()).isEqualTo(totalConsumedMsgs.get()); + }); } - log.debug("totalConsumedMsgs = {}, totalProcessedMsgs = {}", totalConsumedMsgs.get(), totalProcessedMsgs.get()); } @Test @@ -510,7 +513,6 @@ public class TbRuleEngineQueueConsumerManagerTest { verify(consumer).unsubscribe(); int movedMsgs = consumer.msgCount - msgCount; - assertThat(movedMsgs).isGreaterThan(10); verify(ruleEngineMsgProducer, atLeast(movedMsgs)).send(any(), any(), any()); verify(actorContext, never()).tell(any()); generateQueueMsgs = false;