Browse Source

Minor improvement

pull/11924/head
Andrii Landiak 2 years ago
parent
commit
76bed44705
  1. 2
      application/src/main/java/org/thingsboard/server/service/edge/rpc/KafkaEdgeGrpcSession.java
  2. 1
      application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java
  3. 5
      application/src/main/java/org/thingsboard/server/service/ttl/KafkaEdgeTopicsCleanUpService.java
  4. 3
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/TopicService.java
  5. 12
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.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 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);
}

1
application/src/main/java/org/thingsboard/server/service/edge/rpc/PostgresEdgeGrpcSession.java

@ -36,7 +36,6 @@ public class PostgresEdgeGrpcSession extends EdgeGrpcSession {
BiConsumer<Edge, UUID> sessionCloseListener, ScheduledExecutorService sendDownlinkExecutorService,
int maxInboundMessageSize, int maxHighPriorityQueueSizePerSession) {
super(ctx, outputStream, sessionOpenListener, sessionCloseListener, sendDownlinkExecutorService, maxInboundMessageSize, maxHighPriorityQueueSizePerSession);
initInputStream();
}
@Override

5
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 {}",

3
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<String, TopicPartitionInfo> tbCoreNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TopicPartitionInfo> tbRuleEngineNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentMap<String, TopicPartitionInfo> tbEdgeNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentMap<EdgeId, TopicPartitionInfo> tbEdgeEventsNotificationTopics = new ConcurrentHashMap<>();
private final ConcurrentReferenceHashMap<EdgeId, TopicPartitionInfo> tbEdgeEventsNotificationTopics = new ConcurrentReferenceHashMap<>();
/**
* Each Service should start a consumer for messages that target individual service instance based on serviceId.

12
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<String> BASIC_TOPIC_PREFIXES = List.of("tb_edge_event.notifications");
private static final List<String> 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() {

Loading…
Cancel
Save