From 76bed44705a0e40bb4bf96ae73978c6915643312 Mon Sep 17 00:00:00 2001 From: Andrii Landiak Date: Thu, 28 Nov 2024 15:16:20 +0200 Subject: [PATCH] Minor improvement --- .../service/edge/rpc/KafkaEdgeGrpcSession.java | 2 +- .../service/edge/rpc/PostgresEdgeGrpcSession.java | 1 - .../service/ttl/KafkaEdgeTopicsCleanUpService.java | 5 +---- .../server/queue/discovery/TopicService.java | 3 ++- .../server/queue/kafka/TbKafkaSettings.java | 12 ++++-------- 5 files changed, 8 insertions(+), 15 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 6aaa3d72a4..1b2faebbee 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 deleteTopic(EdgeId edgeId) { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + String topic = topicService.getEdgeEventNotificationsTopic(tenantId, edgeId).getTopic(); TbKafkaAdmin kafkaAdmin = new TbKafkaAdmin(kafkaSettings, kafkaTopicConfigs.getEdgeEventConfigs()); kafkaAdmin.deleteTopic(topic); } diff --git a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java index 12516568a8..3f80e625b5 100644 --- a/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java +++ b/application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java @@ -36,7 +36,6 @@ public class PostgresEdgeGrpcSession extends EdgeGrpcSession { BiConsumer sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService, int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) { super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession); - initInputStream(); } @Override 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 1bd0e30c77..29e3b971d4 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 @@ -21,7 +21,6 @@ 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.common.util.ThingsBoardThreadFactory; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; @@ -41,8 +40,6 @@ import org.thingsboard.server.service.state.DefaultDeviceStateService; import java.time.Instant; import java.util.Date; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @Slf4j @@ -92,7 +89,7 @@ public class KafkaEdgeTopicsCleanUpService { .flatMap(AttributeKvEntry::getLongValue) .filter(lastConnectTime -> isTopicExpired(lastConnectTime, ttlMillis, currentTimeMillis)) .ifPresent(lastConnectTime -> { - String topic = topicService.buildEdgeEventNotificationsTopicPartitionInfo(tenantId, edgeId).getTopic(); + 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 {}", diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java index 168e90cd1b..927c311a2d 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java @@ -17,6 +17,7 @@ package org.thingsboard.server.queue.discovery; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; +import org.springframework.util.ConcurrentReferenceHashMap; import org.thingsboard.server.common.data.id.EdgeId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; @@ -34,7 +35,7 @@ public class TopicService { private final ConcurrentMap tbCoreNotificationTopics = new ConcurrentHashMap<>(); private final ConcurrentMap tbRuleEngineNotificationTopics = new ConcurrentHashMap<>(); private final ConcurrentMap tbEdgeNotificationTopics = new ConcurrentHashMap<>(); - private final ConcurrentMap tbEdgeEventsNotificationTopics = new ConcurrentHashMap<>(); + private final ConcurrentReferenceHashMap tbEdgeEventsNotificationTopics = new ConcurrentReferenceHashMap<>(); /** * Each Service should start a consumer for messages that target individual service instance based on serviceId. diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java index 5159b6f769..f228bbae9a 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java @@ -50,7 +50,7 @@ import java.util.Properties; @Component public class TbKafkaSettings { - private static final List BASIC_TOPIC_PREFIXES = List.of("tb_edge_event.notifications"); + private static final List DYNAMIC_TOPICS = List.of("tb_edge_event.notifications"); @Value("${queue.kafka.bootstrap.servers}") private String servers; @@ -167,19 +167,15 @@ public class TbKafkaSettings { .getOrDefault(topic, Collections.emptyList()) .forEach(kv -> props.put(kv.getKey(), kv.getValue())); - applyBaseTopicProperties(props, topic); - - return props; - } - - private void applyBaseTopicProperties(Properties props, String topic) { if (topic != null) { - BASIC_TOPIC_PREFIXES.stream() + DYNAMIC_TOPICS.stream() .filter(topic::startsWith) .findFirst() .ifPresent(prefix -> consumerPropertiesPerTopic.getOrDefault(prefix, Collections.emptyList()) .forEach(kv -> props.put(kv.getKey(), kv.getValue()))); } + + return props; } public Properties toProducerProps() {