Browse Source

Configurable topic deletion delay

pull/8988/head
ViacheslavKlimov 3 years ago
parent
commit
8294331512
  1. 12
      application/src/main/java/org/thingsboard/server/service/entitiy/queue/DefaultTbQueueService.java
  2. 2
      application/src/main/resources/thingsboard.yml

12
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);
}

2
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}"

Loading…
Cancel
Save