Browse Source

Refactor KafkaEdgeTopicsCleanUpService to properly clean up both deleted edges and expired topics

pull/11924/head
Andrii Landiak 2 years ago
parent
commit
988dfce500
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java
  2. 128
      application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java
  3. 9
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java

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

128
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<TenantId> 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<String> topics = kafkaAdmin.getAllTopics();
if (topics == null || topics.isEmpty()) {
log.warn("No topics found in Kafka. Skipping cleanup.");
return;
}
PageDataIterable<EdgeId> edgeIds = new PageDataIterable<>(link -> edgeService.findEdgeIdsByTenantId(tenantId, link), 1024);
String edgeTopicPrefix = prefix.isBlank() ? EDGE_EVENT_TOPIC_NAME : prefix + "." + EDGE_EVENT_TOPIC_NAME;
List<String> 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<TenantId, List<EdgeId>> 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<EdgeId> 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<TenantId, List<EdgeId>> extractTenantAndEdgeIds(List<String> topics, String prefix) {
Map<TenantId, List<EdgeId>> 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;
}
}

9
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<String> 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));
}

Loading…
Cancel
Save