From 829433151207a3aa0cc6d4cd81238f92ce87fd62 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 21 Jul 2023 11:10:39 +0300 Subject: [PATCH] Configurable topic deletion delay --- .../service/entitiy/queue/DefaultTbQueueService.java | 12 +++++++----- application/src/main/resources/thingsboard.yml | 2 ++ 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java index 63e11aeb7a..abf7b694d9 100644 --- a/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java +++ b/application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.service.entitiy.queue; import lombok.AllArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import org.thingsboard.server.cluster.TbClusterService; import org.thingsboard.server.common.data.TenantProfile; @@ -44,13 +45,15 @@ import java.util.stream.Collectors; @TbCoreComponent @AllArgsConstructor public class DefaultTbQueueService extends AbstractTbEntityService implements TbQueueService { - private static final long DELETE_DELAY = 30; private final QueueService queueService; private final TbClusterService tbClusterService; private final TbQueueAdmin tbQueueAdmin; private final SchedulerComponent scheduler; + @Value("${queue.rule-engine.topic_deletion_delay:60}") + private int topicDeletionDelay; + @Override public Queue saveQueue(Queue queue) { boolean create = queue.getId() == null; @@ -119,10 +122,9 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb for (int i = currentPartitions; i < oldPartitions; i++) { String fullTopicName = new TopicPartitionInfo(queue.getTopic(), queue.getTenantId(), i, false).getFullTopicName(); log.info("Removed partition [{}]", fullTopicName); - tbQueueAdmin.deleteTopic( - fullTopicName); + tbQueueAdmin.deleteTopic(fullTopicName); } - }, DELETE_DELAY, TimeUnit.SECONDS); + }, topicDeletionDelay, TimeUnit.SECONDS); } } else if (!oldQueue.equals(queue)) { tbClusterService.onQueueChange(queue); @@ -144,7 +146,7 @@ public class DefaultTbQueueService extends AbstractTbEntityService implements Tb log.error("Failed to delete queue [{}]", fullTopicName); } } - }, DELETE_DELAY, TimeUnit.SECONDS); + }, topicDeletionDelay, TimeUnit.SECONDS); notificationEntityService.notifySendMsgToEdgeService(queue.getTenantId(), queue.getId(), EdgeEventActionType.DELETED); } diff --git a/application/src/main/resources/thingsboard.yml b/application/src/main/resources/thingsboard.yml index 589540096c..2c09221148 100644 --- a/application/src/main/resources/thingsboard.yml +++ b/application/src/main/resources/thingsboard.yml @@ -1232,6 +1232,8 @@ queue: failure-percentage: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_FAILURE_PERCENTAGE:0}" # Skip retry if failures or timeouts are less then X percentage of messages; pause-between-retries: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_RETRY_PAUSE:5}" # Time in seconds to wait in consumer thread before retries; max-pause-between-retries: "${TB_QUEUE_RE_SQ_PROCESSING_STRATEGY_MAX_RETRY_PAUSE:5}" # Max allowed time in seconds for pause between retries. + # Delay between Queue update/delete and actual topic deletion. The delay is for Rule Engines to have time to unsubscribe from the topics, and for other services to stop publishing + topic_deletion_delay: "${TB_QUEUE_RULE_ENGINE_TOPIC_DELETION_DELAY_SECS:60}" transport: # For high priority notifications that require minimum latency and processing time notifications_topic: "${TB_QUEUE_TRANSPORT_NOTIFICATIONS_TOPIC:tb_transport.notifications}"