From 988dfce500add4fce5eb0c11c08f27ee89825226 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Fri, 29 Nov 2024 15:18:27 +0200 Subject: [PATCH] Refactor KafkaEdgeTopicsCleanUpService to properly clean up both deleted edges and expired topics --- .../edge/rpc/KafkaEdgeGrpcSession.java | 2 +- .../ttl/KafkaEdgeTopicsCleanUpService.java | 128 +++++++++++++----- .../server/queue/kafka/TbKafkaAdmin.java | 9 ++ 3 files changed, 105 insertions(+), 34 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java index 5754f7c8f3..9cb3f2dffc 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java @@ -138,7 +138,7 @@ public class KafkaEdgeGrpcSession extends EdgeGrpcSession { @Override public void cleanUp() { - String topic = topicService.getEdgeEventNotificationsTopic(tenantId, edge.getId()).getTopic(); + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edge.getId()).getTopic(); TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); kafkaAdmin.deleteTopic(topic); } diff --git a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java index 03b46921eb..1601310de8 100644 --- a/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java +++ b/application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java @@ -15,18 +15,16 @@ */ package org.thingsboard.server.service.ttl; -import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.condition.ConditionalOnExpression; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.thingsboard.server.common.data.AttributeScope; +import org.thingsboard.server.common.data.edge.Edge; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.kv.AttributeKvEntry; -import org.thingsboard.server.common.data.page.PageDataIterable; -import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.dao.attributes.AttributesService; import org.thingsboard.server.dao.edge.EdgeService; import org.thingsboard.server.dao.tenant.TenantService; @@ -38,7 +36,14 @@ import org.thingsboard.server.queue.kafka.TbKafkaTopicConfigs; import org.thingsboard.server.queue.util.TbCoreComponent; import java.time.Instant; +import java.util.ArrayList; import java.util.Date; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAST_CONNECT_TIME; @@ -46,57 +51,98 @@ import static org.thingsboard.server.service.state.DefaultDeviceStateService.LAS @Slf4j @Service @TbCoreComponent -@RequiredArgsConstructor @ConditionalOnExpression("'${queue.type:null}'=='kafka' && ${edges.enabled:true} && ${sql.ttl.edge_events.edge_events_ttl:0} > 0") -public class KafkaEdgeTopicsCleanUpService { +public class KafkaEdgeTopicsCleanUpService extends AbstractCleanUpService { - private final EdgeService edgeService; - private final TenantService tenantService; - private final AttributesService attributesService; + private static final String EDGE_EVENT_TOPIC_NAME = "tb_edge_event.notifications."; private final TopicService topicService; - private final PartitionService partitionService; + private final TenantService tenantService; + private final EdgeService edgeService; + private final AttributesService attributesService; private final TbKafkaSettings kafkaSettings; private final TbKafkaTopicConfigs kafkaTopicConfigs; + @Value("${queue.prefix:}") + private String prefix; + @Value("${sql.ttl.edge_events.edge_events_ttl:2628000}") private long ttlSeconds; + public KafkaEdgeTopicsCleanUpService(PartitionService partitionService, EdgeService edgeService, + TenantService tenantService, AttributesService attributesService, + TopicService topicService, TbKafkaSettings kafkaSettings, TbKafkaTopicConfigs kafkaTopicConfigs) { + super(partitionService); + this.topicService = topicService; + this.tenantService = tenantService; + this.edgeService = edgeService; + this.attributesService = attributesService; + this.kafkaSettings = kafkaSettings; + this.kafkaTopicConfigs = kafkaTopicConfigs; + } + @Scheduled(initialDelayString = "#{T(org.apache.commons.lang3.RandomUtils).nextLong(0, ${sql.ttl.edge_events.execution_interval_ms})}", fixedDelayString = "${sql.ttl.edge_events.execution_interval_ms}") public void cleanUp() { - PageDataIterable tenants = new PageDataIterable<>(tenantService::findTenantsIds, 10_000); - for (TenantId tenantId : tenants) { - try { - cleanUp(tenantId); - } catch (Exception e) { - log.warn("Failed to drop kafka topics for tenant {}", tenantId, e); - } + if (!isSystemTenantPartitionMine()) { + return; } - } - private void cleanUp(TenantId tenantId) throws Exception { - if (!partitionService.resolve(ServiceType.TB_CORE, tenantId, tenantId).isMyPartition()) { + TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); + Set topics = kafkaAdmin.getAllTopics(); + if (topics == null || topics.isEmpty()) { + log.warn("No topics found in Kafka. Skipping cleanup."); return; } - PageDataIterable edgeIds = new PageDataIterable<>(link -> edgeService.findEdgeIdsByTenantId(tenantId, link), 1024); + String edgeTopicPrefix = prefix.isBlank() ? EDGE_EVENT_TOPIC_NAME : prefix + "." + EDGE_EVENT_TOPIC_NAME; + List matchingTopics = topics.stream().filter(topic -> topic.startsWith(edgeTopicPrefix)).toList(); + if (matchingTopics.isEmpty()) { + log.info("No matching topics found with prefix [{}]. Skipping cleanup.", edgeTopicPrefix); + return; + } + + Map> tenantEdgeMap = extractTenantAndEdgeIds(matchingTopics, edgeTopicPrefix); + long currentTimeMillis = System.currentTimeMillis(); long ttlMillis = TimeUnit.SECONDS.toMillis(ttlSeconds); - for (EdgeId edgeId : edgeIds) { - TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); - attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, LAST_CONNECT_TIME).get() - .flatMap(AttributeKvEntry::getLongValue) - .filter(lastConnectTime -> isTopicExpired(lastConnectTime, ttlMillis, currentTimeMillis)) - .ifPresent(lastConnectTime -> { - String topic = topicService.getEdgeEventNotificationsTopic(tenantId, edgeId).getTopic(); - if (kafkaAdmin.isTopicEmpty(topic)) { - kafkaAdmin.deleteTopic(topic); - log.info("Removed outdated topic for tenant {} and edge with id {} older than {}", - tenantId, edgeId, Date.from(Instant.ofEpochMilli(currentTimeMillis - ttlMillis))); - } - }); + tenantEdgeMap.forEach((tenantId, edgeIds) -> processTenantCleanUp(kafkaAdmin, tenantId, edgeIds, ttlMillis, currentTimeMillis)); + } + + private void processTenantCleanUp(TbKafkaAdmin kafkaAdmin, TenantId tenantId, List edgeIds, long ttlMillis, long currentTimeMillis) { + boolean tenantExists = tenantService.tenantExists(tenantId); + if (tenantExists) { + for (EdgeId edgeId : edgeIds) { + try { + attributesService.find(tenantId, edgeId, AttributeScope.SERVER_SCOPE, LAST_CONNECT_TIME).get() + .flatMap(AttributeKvEntry::getLongValue) + .filter(lastConnectTime -> isTopicExpired(lastConnectTime, ttlMillis, currentTimeMillis)) + .ifPresentOrElse(lastConnectTime -> { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + if (kafkaAdmin.isTopicEmpty(topic)) { + kafkaAdmin.deleteTopic(topic); + log.info("[{}] Removed outdated topic {} for edge {} older than {}", + tenantId, topic, edgeId, Date.from(Instant.ofEpochMilli(currentTimeMillis - ttlMillis))); + } + }, () -> { + Edge edge = edgeService.findEdgeById(tenantId, edgeId); + if (edge == null) { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + kafkaAdmin.deleteTopic(topic); + log.info("[{}] Removed topic {} for deleted edge {}", tenantId, topic, edgeId); + } + }); + } catch (InterruptedException | ExecutionException e) { + log.error("[{}] Failed to delete topic", tenantId); + } + } + } else { + for (EdgeId edgeId : edgeIds) { + String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + kafkaAdmin.deleteTopic(topic); + } + log.info("[{}] Removed topics for not existing tenant and edges {}", tenantId, edgeIds); } } @@ -104,4 +150,20 @@ public class KafkaEdgeTopicsCleanUpService { return lastConnectTime + ttlMillis < currentTimeMillis; } + private Map> extractTenantAndEdgeIds(List topics, String prefix) { + Map> tenantEdgeMap = new HashMap<>(); + for (String topic : topics) { + try { + String remaining = topic.substring(prefix.length()); + String[] parts = remaining.split("\\."); + TenantId tenantId = new TenantId(UUID.fromString(parts[0])); + EdgeId edgeId = new EdgeId(UUID.fromString(parts[1])); + tenantEdgeMap.computeIfAbsent(tenantId, id -> new ArrayList<>()).add(edgeId); + } catch (Exception e) { + log.warn("Failed to extract TenantId and EdgeId from topic [{}]", topic, e); + } + } + return tenantEdgeMap; + } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java index 2d750ae38b..2ea11c7afa 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java @@ -121,6 +121,15 @@ public class TbKafkaAdmin implements TbQueueAdmin { return topics; } + public Set getAllTopics() { + try { + return settings.getAdminClient().listTopics().names().get(); + } catch (InterruptedException | ExecutionException e) { + log.error("Failed to get all topics.", e); + } + return null; + } + public CreateTopicsResult createTopic(NewTopic topic) { return settings.getAdminClient().createTopics(Collections.singletonList(topic)); }