diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java index 7410ec2903..ba762facb1 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperReprocessingService.java @@ -17,23 +17,22 @@ package org.thingsboard.server.service.housekeeper; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.exception.ExceptionUtils; -import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.StringUtils; -import org.thingsboard.server.common.data.id.TenantId; -import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos.HousekeeperTaskProto; import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; +import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.TbQueueProducer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; -import org.thingsboard.server.queue.discovery.PartitionService; +import org.thingsboard.server.queue.housekeeper.HousekeeperConfig; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; +import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; +import org.thingsboard.server.service.queue.consumer.QueueConsumerManager; -import javax.annotation.PostConstruct; import javax.annotation.PreDestroy; import java.util.LinkedHashSet; import java.util.List; @@ -41,8 +40,6 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; @TbCoreComponent @@ -50,160 +47,74 @@ import java.util.concurrent.atomic.AtomicInteger; @Slf4j public class HousekeeperReprocessingService { + private final HousekeeperConfig config; private final HousekeeperService housekeeperService; - private final PartitionService partitionService; - private final TbCoreQueueFactory queueFactory; - private TbQueueProducer> producer; - private TopicPartitionInfo submitTpi; - - @Value("${queue.core.housekeeper.reprocessing-start-delay-sec:300}") - private int startDelay; - @Value("${queue.core.housekeeper.task-reprocessing-delay-sec:3600}") - private int reprocessingDelay; - @Value("${queue.core.housekeeper.max-reprocessing-attempts:10}") - private int maxReprocessingAttempts; - @Value("${queue.core.housekeeper.poll-interval-ms:500}") - private int pollInterval; + private final QueueConsumerManager> consumer; + private final TbQueueProducer> producer; + private final TopicPartitionInfo submitTpi; private final ExecutorService consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-reprocessing-consumer")); - private final ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("housekeeper-reprocessing-scheduler")); protected AtomicInteger cycle = new AtomicInteger(); - private boolean stopped; - public HousekeeperReprocessingService(@Lazy HousekeeperService housekeeperService, - PartitionService partitionService, + public HousekeeperReprocessingService(HousekeeperConfig config, + @Lazy HousekeeperService housekeeperService, TbCoreQueueFactory queueFactory) { + this.config = config; this.housekeeperService = housekeeperService; - this.partitionService = partitionService; - this.queueFactory = queueFactory; - } - - @PostConstruct - private void init() { - producer = queueFactory.createHousekeeperReprocessingMsgProducer(); - submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); - - scheduler.scheduleWithFixedDelay(() -> { - try { - cycle.incrementAndGet(); - startReprocessing(); - } catch (Throwable e) { - log.error("Unexpected error during reprocessing", e); - } - }, startDelay, reprocessingDelay, TimeUnit.SECONDS); + this.consumer = QueueConsumerManager.>builder() + .name("Housekeeper reprocessing") + .msgPackProcessor(this::processMsgs) + .pollInterval(config.getPollInterval()) + .consumerCreator(queueFactory::createHousekeeperReprocessingMsgConsumer) + .consumerExecutor(consumerExecutor) + .build(); + this.producer = queueFactory.createHousekeeperReprocessingMsgProducer(); + this.submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); } - public void startReprocessing() { - if (!partitionService.isMyPartition(ServiceType.TB_CORE, TenantId.SYS_TENANT_ID, TenantId.SYS_TENANT_ID)) { - return; - } - - var consumer = queueFactory.createHousekeeperReprocessingMsgConsumer(); - consumer.subscribe(); - consumerExecutor.submit(() -> { - log.info("Starting Housekeeper tasks reprocessing"); - long startTs = System.currentTimeMillis(); - while (!stopped) { - try { - List> msgs = consumer.poll(pollInterval); - if (msgs.isEmpty() || msgs.stream().anyMatch(msg -> msg.getValue().getTask().getTs() >= startTs)) { - // it's not time yet to process the message - if (!consumer.isCommitSupported()) { - // resubmitting consumed messages if committing is not supported (for in-memory queue) - for (TbProtoQueueMsg msg : msgs) { - submit(msg.getKey(), msg.getValue()); - } - } - break; - } - - for (TbProtoQueueMsg msg : msgs) { - log.trace("Reprocessing task: {}", msg); - try { - reprocessTask(msg.getValue()); - } catch (InterruptedException e) { - return; - } catch (Throwable e) { - log.error("Unexpected error during message reprocessing [{}]", msg, e); - submitForReprocessing(msg.getValue(), e); - } - } - consumer.commit(); - } catch (Throwable t) { - if (!consumer.isStopped()) { - log.warn("Failed to process messages from queue", t); - try { - Thread.sleep(pollInterval); - } catch (InterruptedException interruptedException) { - log.trace("Failed to wait until the server has capacity to handle new requests", interruptedException); - } - } - } - } - consumer.unsubscribe(); - log.info("Stopped Housekeeper tasks reprocessing"); - }); + @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) + public void afterStartUp() { + consumer.subscribe(); // Kafka topic for tasks reprocessing has only 1 partition, so only one TB Core will reprocess tasks + consumer.launch(); } - private void reprocessTask(ToHousekeeperServiceMsg msg) throws Exception { - int attempt = msg.getTask().getAttempt(); - if (attempt > maxReprocessingAttempts) { - if (cycle.get() == 1) { // only reprocessing tasks with exceeded failures on first cycle (after start-up) - log.info("Trying to reprocess task with {} failed attempts: {}", attempt, msg); - } else { - // resubmitting msg to be processed on the next service start - msg = msg.toBuilder() - .setTask(msg.getTask().toBuilder() - .setTs(getReprocessingTs()) - .build()) - .build(); - submit(UUID.randomUUID(), msg); + private void processMsgs(List> msgs, TbQueueConsumer> consumer) throws Exception { + for (TbProtoQueueMsg msg : msgs) { + log.trace("Reprocessing task: {}", msg); + try { + housekeeperService.processTask(msg.getValue()); + } catch (InterruptedException e) { return; + } catch (Throwable e) { + log.error("Unexpected error during message reprocessing [{}]", msg, e); + submitForReprocessing(msg.getValue(), e); } } + consumer.commit(); - housekeeperService.processTask(msg); + Thread.sleep(config.getTaskReprocessingDelay()); } public void submitForReprocessing(ToHousekeeperServiceMsg msg, Throwable error) { HousekeeperTaskProto task = msg.getTask(); - - int attempt = task.getAttempt() + 1; Set errors = new LinkedHashSet<>(task.getErrorsList()); errors.add(StringUtils.truncate(ExceptionUtils.getStackTrace(error), 1024)); msg = msg.toBuilder() .setTask(task.toBuilder() - .setAttempt(attempt) + .setAttempt(task.getAttempt() + 1) .clearErrors().addAllErrors(errors) - .setTs(getReprocessingTs()) .build()) .build(); log.trace("Submitting for reprocessing: {}", msg); - submit(UUID.randomUUID(), msg); // reprocessing topic has single partition, so we don't care about the msg key - - if (task.getAttempt() >= maxReprocessingAttempts) { - log.warn("Failed to process task in {} attempts: {}", task.getAttempt(), msg); - } - } - - private void submit(UUID key, ToHousekeeperServiceMsg msg) { - producer.send(submitTpi, new TbProtoQueueMsg<>(key, msg), null); - } - - private long getReprocessingTs() { - return System.currentTimeMillis() + TimeUnit.SECONDS.toMillis((long) (reprocessingDelay * 0.8)); // *0.8 so that msgs submitted just after finishing reprocessing are processed on the next cycle + producer.send(submitTpi, new TbProtoQueueMsg<>(UUID.randomUUID(), msg), null); } @PreDestroy - private void stop() throws Exception { - stopped = true; - scheduler.shutdownNow(); - consumerExecutor.shutdown(); - if (!consumerExecutor.awaitTermination(10, TimeUnit.SECONDS)) { - consumerExecutor.shutdownNow(); - } + private void stop() { + consumer.stop(); + consumerExecutor.shutdownNow(); log.info("Stopped Housekeeper reprocessing service"); } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java index 1a8d07c9da..89953676c3 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/HousekeeperService.java @@ -16,7 +16,6 @@ package org.thingsboard.server.service.housekeeper; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; @@ -29,11 +28,13 @@ import org.thingsboard.server.gen.transport.TransportProtos.HousekeeperTaskProto import org.thingsboard.server.gen.transport.TransportProtos.ToHousekeeperServiceMsg; import org.thingsboard.server.queue.TbQueueConsumer; import org.thingsboard.server.queue.common.TbProtoQueueMsg; +import org.thingsboard.server.queue.housekeeper.HousekeeperConfig; import org.thingsboard.server.queue.provider.TbCoreQueueFactory; import org.thingsboard.server.queue.util.AfterStartUp; import org.thingsboard.server.queue.util.TbCoreComponent; import org.thingsboard.server.service.housekeeper.processor.HousekeeperTaskProcessor; import org.thingsboard.server.service.housekeeper.stats.HousekeeperStatsService; +import org.thingsboard.server.service.queue.consumer.QueueConsumerManager; import javax.annotation.PreDestroy; import java.util.List; @@ -54,90 +55,79 @@ public class HousekeeperService { private final Map> taskProcessors; + private final HousekeeperConfig config; private final HousekeeperReprocessingService reprocessingService; private final Optional statsService; private final NotificationRuleProcessor notificationRuleProcessor; - private final TbQueueConsumer> consumer; + private final QueueConsumerManager> consumer; private final ExecutorService consumerExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-consumer")); - private final ExecutorService executor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-task-processor")); + private final ExecutorService taskExecutor = Executors.newSingleThreadExecutor(ThingsBoardThreadFactory.forName("housekeeper-task-processor")); - @Value("${queue.core.housekeeper.task-processing-timeout-ms:120000}") - private int taskProcessingTimeout; - @Value("${queue.core.housekeeper.poll-interval-ms:500}") - private int pollInterval; - - private boolean stopped; - - public HousekeeperService(HousekeeperReprocessingService reprocessingService, + public HousekeeperService(HousekeeperConfig config, + HousekeeperReprocessingService reprocessingService, TbCoreQueueFactory queueFactory, Optional statsService, NotificationRuleProcessor notificationRuleProcessor, @Lazy List> taskProcessors) { + this.config = config; this.reprocessingService = reprocessingService; this.statsService = statsService; this.notificationRuleProcessor = notificationRuleProcessor; - this.consumer = queueFactory.createHousekeeperMsgConsumer(); + this.consumer = QueueConsumerManager.>builder() + .name("Housekeeper") + .msgPackProcessor(this::processMsgs) + .pollInterval(config.getPollInterval()) + .consumerCreator(queueFactory::createHousekeeperMsgConsumer) + .consumerExecutor(consumerExecutor) + .build(); this.taskProcessors = taskProcessors.stream().collect(Collectors.toMap(HousekeeperTaskProcessor::getTaskType, p -> p)); } @AfterStartUp(order = AfterStartUp.REGULAR_SERVICE) public void afterStartUp() { consumer.subscribe(); - consumerExecutor.submit(() -> { - while (!stopped && !consumer.isStopped()) { - try { - List> msgs = consumer.poll(pollInterval); - if (msgs.isEmpty()) { - continue; - } + consumer.launch(); + } - for (TbProtoQueueMsg msg : msgs) { - log.trace("Processing task: {}", msg); - try { - processTask(msg.getValue()); - } catch (InterruptedException e) { - return; - } catch (Throwable e) { - log.error("Unexpected error during message processing [{}]", msg, e); - reprocessingService.submitForReprocessing(msg.getValue(), e); - } - } - consumer.commit(); - } catch (Throwable t) { - if (!consumer.isStopped()) { - log.warn("Failed to process messages from queue", t); - try { - Thread.sleep(pollInterval); - } catch (InterruptedException interruptedException) { - log.trace("Failed to wait until the server has capacity to handle new requests", interruptedException); - } - } - } + private void processMsgs(List> msgs, TbQueueConsumer> consumer) { + for (TbProtoQueueMsg msg : msgs) { + log.trace("Processing task: {} attempt={}", msg, msg.getValue().getTask().getAttempt()); + try { + processTask(msg.getValue()); + } catch (InterruptedException e) { + return; + } catch (Throwable e) { + log.error("Unexpected error during message processing [{}]", msg, e); + reprocessingService.submitForReprocessing(msg.getValue(), e); } - }); - log.info("Started Housekeeper service"); + } + consumer.commit(); } @SuppressWarnings("unchecked") protected void processTask(ToHousekeeperServiceMsg msg) throws Exception { HousekeeperTask task = JacksonUtil.fromString(msg.getTask().getValue(), HousekeeperTask.class); - HousekeeperTaskProcessor taskProcessor = (HousekeeperTaskProcessor) taskProcessors.get(task.getTaskType()); + HousekeeperTaskType taskType = task.getTaskType(); + if (config.getDisabledTaskTypes().contains(taskType)) { + log.debug("Task type {} is disabled, ignoring {}", taskType, task); + return; + } + HousekeeperTaskProcessor taskProcessor = (HousekeeperTaskProcessor) taskProcessors.get(taskType); if (taskProcessor == null) { - throw new IllegalArgumentException("Unsupported task type " + task.getTaskType()); + throw new IllegalArgumentException("Unsupported task type " + taskType); } if (log.isDebugEnabled()) { log.debug("[{}] {} {}", task.getTenantId(), isNew(msg.getTask()) ? "Processing" : "Reprocessing", task.getDescription()); } try { - Future future = executor.submit(() -> { + Future future = taskExecutor.submit(() -> { taskProcessor.process((T) task); return null; }); - future.get(taskProcessingTimeout, TimeUnit.MILLISECONDS); - - statsService.ifPresent(statsService -> statsService.reportProcessed(task.getTaskType(), msg)); + future.get(config.getTaskProcessingTimeout(), TimeUnit.MILLISECONDS); + statsService.ifPresent(statsService -> statsService.reportProcessed(taskType, msg)); } catch (InterruptedException e) { throw e; } catch (Throwable e) { @@ -145,19 +135,23 @@ public class HousekeeperService { if (e instanceof ExecutionException) { error = e.getCause(); } else if (e instanceof TimeoutException) { - error = new TimeoutException("Timeout after " + taskProcessingTimeout + " seconds"); + error = new TimeoutException("Timeout after " + config.getTaskProcessingTimeout() + " seconds"); } log.error("[{}][{}][{}] {} task processing failed, submitting for reprocessing (attempt {}): {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), - task.getTaskType(), msg.getTask().getAttempt(), task, error); - reprocessingService.submitForReprocessing(msg, error); - - statsService.ifPresent(statsService -> statsService.reportFailure(task.getTaskType(), msg)); - notificationRuleProcessor.process(TaskProcessingFailureTrigger.builder() - .task(task) - .error(error) - .attempt(msg.getTask().getAttempt()) - .build()); + taskType, msg.getTask().getAttempt(), task, error); + + if (msg.getTask().getAttempt() < config.getMaxReprocessingAttempts()) { + reprocessingService.submitForReprocessing(msg, error); + } else { + log.error("Failed to process task in {} attempts: {}", msg.getTask().getAttempt(), msg); + notificationRuleProcessor.process(TaskProcessingFailureTrigger.builder() + .task(task) + .error(error) + .attempt(msg.getTask().getAttempt()) + .build()); + } + statsService.ifPresent(statsService -> statsService.reportFailure(taskType, msg)); } } @@ -167,12 +161,8 @@ public class HousekeeperService { @PreDestroy private void stop() throws Exception { - stopped = true; - consumer.unsubscribe(); - consumerExecutor.shutdown(); - if (!consumerExecutor.awaitTermination(10, TimeUnit.SECONDS)) { - consumerExecutor.shutdownNow(); - } + consumer.stop(); + consumerExecutor.shutdownNow(); log.info("Stopped Housekeeper service"); } diff --git a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java index 400bc36479..4944cf9f43 100644 --- a/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java +++ b/application/src/main/java/org/thingsboard/server/service/housekeeper/stats/HousekeeperStatsService.java @@ -105,7 +105,7 @@ public class HousekeeperStatsService { } private StatsCounter register(String statsName, StatsFactory statsFactory) { - StatsCounter counter = statsFactory.createStatsCounter(StatsType.HOUSEKEEPER.getName(), statsName, Map.of("taskType", taskType.name())); + StatsCounter counter = statsFactory.createStatsCounter(StatsType.HOUSEKEEPER.getName(), statsName, "taskType", taskType.name()); counters.add(counter); return counter; } diff --git a/application/src/main/java/org/thingsboard/server/service/queue/consumer/QueueConsumerManager.java b/application/src/main/java/org/thingsboard/server/service/queue/consumer/QueueConsumerManager.java new file mode 100644 index 0000000000..e7b3a2f2b7 --- /dev/null +++ b/application/src/main/java/org/thingsboard/server/service/queue/consumer/QueueConsumerManager.java @@ -0,0 +1,109 @@ +/** + * Copyright © 2016-2024 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.service.queue.consumer; + +import lombok.Builder; +import lombok.Getter; +import lombok.extern.slf4j.Slf4j; +import org.thingsboard.common.util.ThingsBoardThreadFactory; +import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; +import org.thingsboard.server.queue.TbQueueConsumer; +import org.thingsboard.server.queue.TbQueueMsg; + +import java.util.List; +import java.util.Set; +import java.util.concurrent.ExecutorService; +import java.util.function.Supplier; + +@Slf4j +public class QueueConsumerManager { + + private final String name; + private final MsgPackProcessor msgPackProcessor; + private final long pollInterval; + private final ExecutorService consumerExecutor; + private final String threadPrefix; + + @Getter + private final TbQueueConsumer consumer; + private volatile boolean stopped; + + @Builder + public QueueConsumerManager(String name, MsgPackProcessor msgPackProcessor, + long pollInterval, Supplier> consumerCreator, + ExecutorService consumerExecutor, String threadPrefix) { + this.name = name; + this.pollInterval = pollInterval; + this.msgPackProcessor = msgPackProcessor; + this.consumerExecutor = consumerExecutor; + this.threadPrefix = threadPrefix; + this.consumer = consumerCreator.get(); + } + + public void subscribe() { + consumer.subscribe(); + } + + public void subscribe(Set partitions) { + consumer.subscribe(partitions); + } + + public void launch() { + log.info("[{}] Launching consumer", name); + consumerExecutor.submit(() -> { + if (threadPrefix != null) { + ThingsBoardThreadFactory.addThreadNamePrefix(threadPrefix); + } + try { + consumerLoop(consumer); + } catch (Throwable e) { + log.error("Failure in consumer loop", e); + } + log.info("[{}] Consumer stopped", name); + }); + } + + private void consumerLoop(TbQueueConsumer consumer) { + while (!stopped && !consumer.isStopped()) { + try { + List msgs = consumer.poll(pollInterval); + if (msgs.isEmpty()) { + continue; + } + msgPackProcessor.process(msgs, consumer); + } catch (Exception e) { + if (!consumer.isStopped()) { + log.warn("Failed to process messages from queue", e); + try { + Thread.sleep(pollInterval); + } catch (InterruptedException interruptedException) { + log.trace("Failed to wait until the server has capacity to handle new requests", interruptedException); + } + } + } + } + } + + public void stop() { + log.debug("[{}] Stopping consumer", name); + stopped = true; + consumer.unsubscribe(); + } + + public interface MsgPackProcessor { + void process(List msgs, TbQueueConsumer consumer) throws Exception; + } +} \ No newline at end of file diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 938ef7986e..cbc4e79dd3 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1638,11 +1638,13 @@ queue: poll-interval-ms: "${TB_HOUSEKEEPER_POLL_INTERVAL_MS:500}" # Timeout in milliseconds for task processing. Tasks that fail to finish on time will be submitted for reprocessing task-processing-timeout-ms: "${TB_HOUSEKEEPER_TASK_PROCESSING_TIMEOUT_MS:120000}" - # Delay in seconds for starting tasks reprocessing after start-up - reprocessing-start-delay-sec: "${TB_HOUSEKEEPER_REPROCESSING_START_DELAY_SEC:300}" - # Period in seconds of tasks reprocessing - task-reprocessing-delay-sec: "${TB_HOUSEKEEPER_TASK_REPROCESSING_DELAY_SEC:3600}" - # Maximum amount of task reprocessing attempts. After exceeding, the task will be ignored until the next service start-up + # Comma-separated list of task types that shouldn't be processed. Available task types: + # DELETE_ATTRIBUTES, DELETE_TELEMETRY (both DELETE_LATEST_TS and DELETE_TS_HISTORY will be disabled), + # DELETE_LATEST_TS, DELETE_TS_HISTORY, DELETE_EVENTS, DELETE_ALARMS, UNASSIGN_ALARMS + disabled-task-types: "${TB_HOUSEKEEPER_DISABLED_TASK_TYPES:}" + # Delay in milliseconds between tasks reprocessing + task-reprocessing-delay-ms: "${TB_HOUSEKEEPER_TASK_REPROCESSING_DELAY_MS:5000}" + # Maximum amount of task reprocessing attempts. After exceeding, the task will be dropped max-reprocessing-attempts: "${TB_HOUSEKEEPER_MAX_REPROCESSING_ATTEMPTS:10}" stats: # Enable/disable statistics for Housekeeper diff --git a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java index 44891c2b7f..82c488f0ab 100644 --- a/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java +++ b/application/src/test/java/org/thingsboard/server/service/housekeeper/HousekeeperServiceTest.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.housekeeper; import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.node.TextNode; +import net.bytebuddy.implementation.bytecode.Throw; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -104,8 +105,7 @@ import static org.springframework.test.web.servlet.result.MockMvcResultMatchers. @DaoSqlTest @TestPropertySource(properties = { "transport.http.enabled=true", - "queue.core.housekeeper.reprocessing-start-delay-sec=1", - "queue.core.housekeeper.task-reprocessing-delay-sec=2", + "queue.core.housekeeper.task-reprocessing-delay-ms=2000", "queue.core.housekeeper.poll-interval-ms=1000", "queue.core.housekeeper.max-reprocessing-attempts=5" }) @@ -258,7 +258,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest { } @Test - public void whenTaskProcessingFails_thenReprocessUntilSuccessful() throws Exception { + public void whenTaskProcessingFails_thenReprocess() throws Exception { TimeoutException error = new TimeoutException("Test timeout"); doThrow(error).when(tsHistoryDeletionTaskProcessor).process(any()); @@ -267,8 +267,8 @@ public class HousekeeperServiceTest extends AbstractControllerTest { doDelete("/api/device/" + device.getId()).andExpect(status().isOk()); - int attempts = 3; - await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> { + int attempts = 2; + await().atMost(30, TimeUnit.SECONDS).pollInterval(1, TimeUnit.SECONDS).untilAsserted(() -> { verifyTaskProcessing(device.getId(), HousekeeperTaskType.DELETE_TS_HISTORY, 0); for (int i = 1; i <= attempts; i++) { int attempt = i; @@ -279,7 +279,7 @@ public class HousekeeperServiceTest extends AbstractControllerTest { }); assertThat(getTimeseriesHistory(device.getId())).isNotEmpty(); - doCallRealMethod().when(tsHistoryDeletionTaskProcessor).process(any()); // fixing the code + doCallRealMethod().when(tsHistoryDeletionTaskProcessor).process(any()); await().atMost(30, TimeUnit.SECONDS).untilAsserted(() -> { assertThat(getTimeseriesHistory(device.getId())).isEmpty(); }); @@ -314,7 +314,12 @@ public class HousekeeperServiceTest extends AbstractControllerTest { } private void verifyTaskProcessing(EntityId entityId, HousekeeperTaskType taskType, int expectedAttempt) throws Exception { - verify(housekeeperService).processTask(argThat(getTaskMatcher(entityId, taskType, task -> task.getAttempt() == expectedAttempt))); + try { + verify(housekeeperService).processTask(argThat(getTaskMatcher(entityId, taskType, task -> task.getAttempt() == expectedAttempt))); + } catch (Throwable e) { + e.printStackTrace(); + throw e; + } } private ArgumentMatcher getTaskMatcher(EntityId entityId, HousekeeperTaskType taskType, diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java index af4e746987..5911418dc3 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/AlarmsUnassignHousekeeperTask.java @@ -19,9 +19,11 @@ import lombok.AccessLevel; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.ToString; import org.thingsboard.server.common.data.User; @Data +@ToString(callSuper = true) @EqualsAndHashCode(callSuper = true) @NoArgsConstructor(access = AccessLevel.PROTECTED) public class AlarmsUnassignHousekeeperTask extends HousekeeperTask { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java index 515703991c..4156200758 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/EntitiesDeletionHousekeeperTask.java @@ -20,10 +20,12 @@ import lombok.AccessLevel; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.ToString; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.TenantId; @Data +@ToString(callSuper = true) @EqualsAndHashCode(callSuper = true) @NoArgsConstructor(access = AccessLevel.PROTECTED) public class EntitiesDeletionHousekeeperTask extends HousekeeperTask { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java index 59e815bbfa..cef43756e0 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/LatestTsDeletionHousekeeperTask.java @@ -19,10 +19,12 @@ import lombok.AccessLevel; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.ToString; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @Data +@ToString(callSuper = true) @EqualsAndHashCode(callSuper = true) @NoArgsConstructor(access = AccessLevel.PROTECTED) public class LatestTsDeletionHousekeeperTask extends HousekeeperTask { diff --git a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java index 1656626e6b..80b8edc345 100644 --- a/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java +++ b/common/data/src/main/java/org/thingsboard/server/common/data/housekeeper/TsHistoryDeletionHousekeeperTask.java @@ -19,10 +19,12 @@ import lombok.AccessLevel; import lombok.Data; import lombok.EqualsAndHashCode; import lombok.NoArgsConstructor; +import lombok.ToString; import org.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; @Data +@ToString(callSuper = true) @EqualsAndHashCode(callSuper = true) @NoArgsConstructor(access = AccessLevel.PROTECTED) public class TsHistoryDeletionHousekeeperTask extends HousekeeperTask { diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java index e9737bf90c..3b06f0cf3e 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/DefaultHousekeeperClient.java @@ -19,6 +19,7 @@ import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; import org.thingsboard.common.util.JacksonUtil; import org.thingsboard.server.common.data.housekeeper.HousekeeperTask; +import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; import org.thingsboard.server.common.msg.housekeeper.HousekeeperClient; import org.thingsboard.server.common.msg.queue.TopicPartitionInfo; import org.thingsboard.server.gen.transport.TransportProtos; @@ -33,11 +34,14 @@ import org.thingsboard.server.queue.provider.TbQueueProducerProvider; @Slf4j public class DefaultHousekeeperClient implements HousekeeperClient { + private final HousekeeperConfig config; private final TbQueueProducer> producer; private final TopicPartitionInfo submitTpi; private final TbQueueCallback submitCallback; - public DefaultHousekeeperClient(TbQueueProducerProvider producerProvider) { + public DefaultHousekeeperClient(HousekeeperConfig config, + TbQueueProducerProvider producerProvider) { + this.config = config; this.producer = producerProvider.getHousekeeperMsgProducer(); this.submitTpi = TopicPartitionInfo.builder().topic(producer.getDefaultTopic()).build(); this.submitCallback = new TbQueueCallback() { @@ -55,7 +59,13 @@ public class DefaultHousekeeperClient implements HousekeeperClient { @Override public void submitTask(HousekeeperTask task) { - log.debug("[{}][{}][{}] Submitting task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), task.getTaskType()); + HousekeeperTaskType taskType = task.getTaskType(); + if (config.getDisabledTaskTypes().contains(taskType)) { + log.trace("Task type {} is disabled, ignoring {}", taskType, task); + return; + } + + log.debug("[{}][{}][{}] Submitting task: {}", task.getTenantId(), task.getEntityId().getEntityType(), task.getEntityId(), taskType); /* * using msg key as entity id so that msgs related to certain entity are pushed to same partition, * e.g. on tenant deletion (entity id is tenant id), we need to clean up tenant entities in certain order diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/HousekeeperConfig.java b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/HousekeeperConfig.java new file mode 100644 index 0000000000..b03861a7ee --- /dev/null +++ b/common/queue/src/main/java/org/thingsboard/server/queue/housekeeper/HousekeeperConfig.java @@ -0,0 +1,40 @@ +/** + * Copyright © 2016-2024 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.queue.housekeeper; + +import lombok.Getter; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.thingsboard.server.common.data.housekeeper.HousekeeperTaskType; + +import java.util.Set; + +@Component +@Getter +public class HousekeeperConfig { + + @Value("${queue.core.housekeeper.disabled-task-types:}") + private Set disabledTaskTypes; + @Value("${queue.core.housekeeper.task-processing-timeout-ms:120000}") + private int taskProcessingTimeout; + @Value("${queue.core.housekeeper.poll-interval-ms:500}") + private int pollInterval; + @Value("${queue.core.housekeeper.task-reprocessing-delay-ms:5000}") + private int taskReprocessingDelay; + @Value("${queue.core.housekeeper.max-reprocessing-attempts:10}") + private int maxReprocessingAttempts; + +} diff --git a/common/util/src/main/java/org/thingsboard/common/util/ThingsBoardThreadFactory.java b/common/util/src/main/java/org/thingsboard/common/util/ThingsBoardThreadFactory.java index de097cfd6c..4d0876bf8e 100644 --- a/common/util/src/main/java/org/thingsboard/common/util/ThingsBoardThreadFactory.java +++ b/common/util/src/main/java/org/thingsboard/common/util/ThingsBoardThreadFactory.java @@ -51,6 +51,11 @@ public class ThingsBoardThreadFactory implements ThreadFactory { Thread.currentThread().setName(name); } + public static void addThreadNamePrefix(String prefix) { + String name = Thread.currentThread().getName(); + name = prefix + "-" + name; + Thread.currentThread().setName(name); + } @Override public Thread newThread(Runnable r) {