From 4daf53bf25af545789b83ffbe8c98f0aa03e911b Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 16 May 2024 11:52:41 +0300 Subject: [PATCH 1/5] Lazy Jpa bootstrap mode --- .../thingsboard/server/dao/JpaDaoConfig.java | 3 +- .../server/dao/SqlTimeseriesDaoConfig.java | 4 +- .../server/dao/SqlTsDaoConfig.java | 3 +- .../server/dao/SqlTsLatestDaoConfig.java | 3 +- .../server/dao/TimescaleDaoConfig.java | 3 +- .../dao/TimescaleTsLatestDaoConfig.java | 3 +- .../dao/service/JpaRepositoriesTest.java | 42 +++++++++++++++++++ 7 files changed, 54 insertions(+), 7 deletions(-) create mode 100644 dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java diff --git a/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java index 05becc8bde..4dd0633420 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.TbAutoConfiguration; @@ -28,7 +29,7 @@ import org.thingsboard.server.dao.util.TbAutoConfiguration; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sql", "org.thingsboard.server.dao.attributes", "org.thingsboard.server.dao.cache", "org.thingsboard.server.cache"}) -@EnableJpaRepositories("org.thingsboard.server.dao.sql") +@EnableJpaRepositories(value = "org.thingsboard.server.dao.sql", bootstrapMode = BootstrapMode.LAZY) @EntityScan("org.thingsboard.server.dao.model.sql") @EnableTransactionManagement public class JpaDaoConfig { diff --git a/dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java index 9cf9532aea..09553fb973 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java @@ -19,14 +19,14 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; -import org.thingsboard.server.dao.util.SqlTsOrTsLatestAnyDao; import org.thingsboard.server.dao.util.TbAutoConfiguration; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sqlts.dictionary"}) -@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.dictionary"}) +@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sqlts.dictionary"}, bootstrapMode = BootstrapMode.LAZY) @EntityScan({"org.thingsboard.server.dao.model.sqlts.dictionary"}) @EnableTransactionManagement public class SqlTimeseriesDaoConfig { diff --git a/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java index 6bdf7c2170..bbd31846db 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.SqlTsDao; import org.thingsboard.server.dao.util.TbAutoConfiguration; @@ -26,7 +27,7 @@ import org.thingsboard.server.dao.util.TbAutoConfiguration; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sqlts.sql", "org.thingsboard.server.dao.sqlts.insert.sql"}) -@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.ts", "org.thingsboard.server.dao.sqlts.insert.sql"}) +@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sqlts.ts", "org.thingsboard.server.dao.sqlts.insert.sql"}, bootstrapMode = BootstrapMode.LAZY) @EntityScan({"org.thingsboard.server.dao.model.sqlts.ts"}) @EnableTransactionManagement @SqlTsDao diff --git a/dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java index 1793c4ed76..6fcf21a0a8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.SqlTsLatestDao; import org.thingsboard.server.dao.util.TbAutoConfiguration; @@ -26,7 +27,7 @@ import org.thingsboard.server.dao.util.TbAutoConfiguration; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sqlts.sql"}) -@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.insert.latest.sql", "org.thingsboard.server.dao.sqlts.latest"}) +@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sqlts.insert.latest.sql", "org.thingsboard.server.dao.sqlts.latest"}, bootstrapMode = BootstrapMode.LAZY) @EntityScan({"org.thingsboard.server.dao.model.sqlts.latest"}) @EnableTransactionManagement @SqlTsLatestDao diff --git a/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java index 16617c12a8..c134b88897 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.TbAutoConfiguration; import org.thingsboard.server.dao.util.TimescaleDBTsDao; @@ -26,7 +27,7 @@ import org.thingsboard.server.dao.util.TimescaleDBTsDao; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sqlts.timescale"}) -@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.timescale", "org.thingsboard.server.dao.sqlts.insert.timescale"}) +@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sqlts.timescale", "org.thingsboard.server.dao.sqlts.insert.timescale"}, bootstrapMode = BootstrapMode.LAZY) @EntityScan({"org.thingsboard.server.dao.model.sqlts.timescale"}) @EnableTransactionManagement @TimescaleDBTsDao diff --git a/dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java b/dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java index f0e6fc362b..f0e74f1c3d 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java +++ b/dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java @@ -19,6 +19,7 @@ import org.springframework.boot.autoconfigure.domain.EntityScan; import org.springframework.context.annotation.ComponentScan; import org.springframework.context.annotation.Configuration; import org.springframework.data.jpa.repository.config.EnableJpaRepositories; +import org.springframework.data.repository.config.BootstrapMode; import org.springframework.transaction.annotation.EnableTransactionManagement; import org.thingsboard.server.dao.util.TbAutoConfiguration; import org.thingsboard.server.dao.util.TimescaleDBTsLatestDao; @@ -26,7 +27,7 @@ import org.thingsboard.server.dao.util.TimescaleDBTsLatestDao; @Configuration @TbAutoConfiguration @ComponentScan({"org.thingsboard.server.dao.sqlts.timescale"}) -@EnableJpaRepositories({"org.thingsboard.server.dao.sqlts.insert.latest.sql", "org.thingsboard.server.dao.sqlts.latest"}) +@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sqlts.insert.latest.sql", "org.thingsboard.server.dao.sqlts.latest"}, bootstrapMode = BootstrapMode.LAZY) @EntityScan({"org.thingsboard.server.dao.model.sqlts.latest"}) @EnableTransactionManagement @TimescaleDBTsLatestDao diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java new file mode 100644 index 0000000000..a822401aeb --- /dev/null +++ b/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java @@ -0,0 +1,42 @@ +/** + * Copyright © 2016-2024 The Thingsboard Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.thingsboard.server.dao.service; + +import lombok.extern.slf4j.Slf4j; +import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.jpa.repository.JpaRepository; + +import java.util.List; + +@Slf4j +@DaoSqlTest +public class JpaRepositoriesTest extends AbstractServiceTest { + + @Autowired + List> repositories; + + /* + * Verifying that all the repositories are successfully bootstrapped, when using Lazy Jpa bootstrap mode + * */ + @Test + public void testRepositories() { + for (var repository : repositories) { + repository.count(); + } + } + +} From 45672d99885220f69a65b61cd7e67589a6798db5 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 16 May 2024 11:57:28 +0300 Subject: [PATCH 2/5] Move repositories test to EntityDaoRegistryTest --- .../dao/service/EntityDaoRegistryTest.java | 15 +++++++ .../dao/service/JpaRepositoriesTest.java | 42 ------------------- 2 files changed, 15 insertions(+), 42 deletions(-) delete mode 100644 dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/EntityDaoRegistryTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/EntityDaoRegistryTest.java index 8524efc4a0..30110a2741 100644 --- a/dao/src/test/java/org/thingsboard/server/dao/service/EntityDaoRegistryTest.java +++ b/dao/src/test/java/org/thingsboard/server/dao/service/EntityDaoRegistryTest.java @@ -19,11 +19,13 @@ import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.exception.ExceptionUtils; import org.junit.Test; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.jpa.repository.JpaRepository; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.Dao; import org.thingsboard.server.dao.entity.EntityDaoRegistry; +import java.util.List; import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; @@ -37,6 +39,9 @@ public class EntityDaoRegistryTest extends AbstractServiceTest { @Autowired EntityDaoRegistry entityDaoRegistry; + @Autowired + List> repositories; + @Test public void givenAllEntityTypes_whenGetDao_thenAllPresent() { for (EntityType entityType : EntityType.values()) { @@ -73,4 +78,14 @@ public class EntityDaoRegistryTest extends AbstractServiceTest { } } + /* + * Verifying that all the repositories are successfully bootstrapped, when using Lazy Jpa bootstrap mode + * */ + @Test + public void testJpaRepositories() { + for (var repository : repositories) { + repository.count(); + } + } + } diff --git a/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java b/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java deleted file mode 100644 index a822401aeb..0000000000 --- a/dao/src/test/java/org/thingsboard/server/dao/service/JpaRepositoriesTest.java +++ /dev/null @@ -1,42 +0,0 @@ -/** - * Copyright © 2016-2024 The Thingsboard Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.thingsboard.server.dao.service; - -import lombok.extern.slf4j.Slf4j; -import org.junit.Test; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.data.jpa.repository.JpaRepository; - -import java.util.List; - -@Slf4j -@DaoSqlTest -public class JpaRepositoriesTest extends AbstractServiceTest { - - @Autowired - List> repositories; - - /* - * Verifying that all the repositories are successfully bootstrapped, when using Lazy Jpa bootstrap mode - * */ - @Test - public void testRepositories() { - for (var repository : repositories) { - repository.count(); - } - } - -} From e8d14f37cb012243ebe5033dde274674c7d33435 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 16 May 2024 13:24:49 +0300 Subject: [PATCH 3/5] Improvements for Kafka admin creation --- .../server/queue/kafka/TbKafkaAdmin.java | 45 +++++++++++-------- .../kafka/TbKafkaConsumerStatsService.java | 12 ++--- .../server/queue/kafka/TbKafkaSettings.java | 35 ++++++++++++--- 3 files changed, 58 insertions(+), 34 deletions(-) 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 bf5555817d..a5a0427606 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 @@ -16,7 +16,6 @@ package org.thingsboard.server.queue.kafka; import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.CreateTopicsResult; import org.apache.kafka.clients.admin.NewTopic; import org.apache.kafka.common.errors.TopicExistsException; @@ -35,23 +34,17 @@ import java.util.concurrent.ExecutionException; @Slf4j public class TbKafkaAdmin implements TbQueueAdmin { - private final AdminClient client; + private final TbKafkaSettings settings; private final Map topicConfigs; - private final Set topics = ConcurrentHashMap.newKeySet(); private final int numPartitions; + private volatile Set topics; private final short replicationFactor; public TbKafkaAdmin(TbKafkaSettings settings, Map topicConfigs) { - client = AdminClient.create(settings.toAdminProps()); + this.settings = settings; this.topicConfigs = topicConfigs; - try { - topics.addAll(client.listTopics().names().get()); - } catch (InterruptedException | ExecutionException e) { - log.error("Failed to get all topics.", e); - } - String numPartitionsStr = topicConfigs.get(TbKafkaTopicConfigs.NUM_PARTITIONS_SETTING); if (numPartitionsStr != null) { numPartitions = Integer.parseInt(numPartitionsStr); @@ -64,6 +57,7 @@ public class TbKafkaAdmin implements TbQueueAdmin { @Override public void createTopicIfNotExists(String topic, String properties) { + Set topics = getTopics(); if (topics.contains(topic)) { return; } @@ -86,12 +80,13 @@ public class TbKafkaAdmin implements TbQueueAdmin { @Override public void deleteTopic(String topic) { + Set topics = getTopics(); if (topics.contains(topic)) { - client.deleteTopics(Collections.singletonList(topic)); + settings.getAdminClient().deleteTopics(Collections.singletonList(topic)); } else { try { - if (client.listTopics().names().get().contains(topic)) { - client.deleteTopics(Collections.singletonList(topic)); + if (settings.getAdminClient().listTopics().names().get().contains(topic)) { + settings.getAdminClient().deleteTopics(Collections.singletonList(topic)); } else { log.warn("Kafka topic [{}] does not exist.", topic); } @@ -101,14 +96,28 @@ public class TbKafkaAdmin implements TbQueueAdmin { } } - @Override - public void destroy() { - if (client != null) { - client.close(); + private Set getTopics() { + if (topics == null) { + synchronized (this) { + if (topics == null) { + topics = ConcurrentHashMap.newKeySet(); + try { + topics.addAll(settings.getAdminClient().listTopics().names().get()); + } catch (InterruptedException | ExecutionException e) { + log.error("Failed to get all topics.", e); + } + } + } } + return topics; } public CreateTopicsResult createTopic(NewTopic topic) { - return client.createTopics(Collections.singletonList(topic)); + return settings.getAdminClient().createTopics(Collections.singletonList(topic)); + } + + @Override + public void destroy() { } + } diff --git a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerStatsService.java b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerStatsService.java index dbc1b7b7c4..56440cf4a5 100644 --- a/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerStatsService.java +++ b/common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerStatsService.java @@ -15,11 +15,12 @@ */ package org.thingsboard.server.queue.kafka; +import jakarta.annotation.PostConstruct; +import jakarta.annotation.PreDestroy; import lombok.Builder; import lombok.Data; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; @@ -35,8 +36,6 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.msg.queue.ServiceType; import org.thingsboard.server.queue.discovery.PartitionService; -import jakarta.annotation.PostConstruct; -import jakarta.annotation.PreDestroy; import java.time.Duration; import java.util.ArrayList; import java.util.List; @@ -62,7 +61,6 @@ public class TbKafkaConsumerStatsService { @Autowired private PartitionService partitionService; - private AdminClient adminClient; private Consumer consumer; private ScheduledExecutorService statsPrintScheduler; @@ -71,7 +69,6 @@ public class TbKafkaConsumerStatsService { if (!statsConfig.getEnabled()) { return; } - this.adminClient = AdminClient.create(kafkaSettings.toAdminProps()); this.statsPrintScheduler = Executors.newSingleThreadScheduledExecutor(ThingsBoardThreadFactory.forName("kafka-consumer-stats")); Properties consumerProps = kafkaSettings.toConsumerProps(null); @@ -90,7 +87,7 @@ public class TbKafkaConsumerStatsService { } for (String groupId : monitoredGroups) { try { - Map groupOffsets = adminClient.listConsumerGroupOffsets(groupId).partitionsToOffsetAndMetadata() + Map groupOffsets = kafkaSettings.getAdminClient().listConsumerGroupOffsets(groupId).partitionsToOffsetAndMetadata() .get(statsConfig.getKafkaResponseTimeoutMs(), TimeUnit.MILLISECONDS); Map endOffsets = consumer.endOffsets(groupOffsets.keySet(), timeoutDuration); @@ -157,9 +154,6 @@ public class TbKafkaConsumerStatsService { if (statsPrintScheduler != null) { statsPrintScheduler.shutdownNow(); } - if (adminClient != null) { - adminClient.close(); - } if (consumer != null) { consumer.close(); } 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 96fa7afb7f..ebbc6fac7e 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 @@ -19,6 +19,7 @@ import lombok.Getter; import lombok.Setter; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.clients.admin.AdminClient; import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; @@ -34,6 +35,7 @@ import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.TbProperty; import org.thingsboard.server.queue.util.PropertyUtils; +import javax.annotation.PreDestroy; import java.util.Collections; import java.util.List; import java.util.Map; @@ -143,13 +145,7 @@ public class TbKafkaSettings { @Setter private Map> consumerPropertiesPerTopic = Collections.emptyMap(); - public Properties toAdminProps() { - Properties props = toProps(); - props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, servers); - props.put(AdminClientConfig.RETRIES_CONFIG, retries); - - return props; - } + private volatile AdminClient adminClient; public Properties toConsumerProps(String topic) { Properties props = toProps(); @@ -221,4 +217,29 @@ public class TbKafkaSettings { } } + public AdminClient getAdminClient() { + if (adminClient == null) { + synchronized (this) { + if (adminClient == null) { + adminClient = AdminClient.create(toAdminProps()); + } + } + } + return adminClient; + } + + private Properties toAdminProps() { + Properties props = toProps(); + props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, servers); + props.put(AdminClientConfig.RETRIES_CONFIG, retries); + return props; + } + + @PreDestroy + private void destroy() { + if (adminClient != null) { + adminClient.close(); + } + } + } From 1a264dafa70a31210795a9b45a9be151034e2002 Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Thu, 16 May 2024 14:34:50 +0300 Subject: [PATCH 4/5] Fix TbKafkaSettingsTest --- .../org/thingsboard/server/queue/kafka/TbKafkaSettings.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 ebbc6fac7e..760487c61e 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 @@ -228,7 +228,7 @@ public class TbKafkaSettings { return adminClient; } - private Properties toAdminProps() { + protected Properties toAdminProps() { Properties props = toProps(); props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, servers); props.put(AdminClientConfig.RETRIES_CONFIG, retries); From 9a5e3c81fc0ea5f7b9cfa34f23ffa098917eeceb Mon Sep 17 00:00:00 2001 From: ViacheslavKlimov Date: Fri, 17 May 2024 13:11:33 +0300 Subject: [PATCH 5/5] Log startup time --- .../server/ThingsboardServerApplication.java | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java b/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java index 60f69cd8a7..6ec00b5661 100644 --- a/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java +++ b/application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java @@ -15,24 +15,32 @@ */ package org.thingsboard.server; +import lombok.extern.slf4j.Slf4j; import org.springframework.boot.SpringApplication; import org.springframework.boot.SpringBootConfiguration; import org.springframework.context.annotation.ComponentScan; +import org.springframework.core.Ordered; import org.springframework.scheduling.annotation.EnableAsync; import org.springframework.scheduling.annotation.EnableScheduling; +import org.thingsboard.server.queue.util.AfterStartUp; import java.util.Arrays; +import java.util.concurrent.TimeUnit; @SpringBootConfiguration @EnableAsync @EnableScheduling @ComponentScan({"org.thingsboard.server", "org.thingsboard.script"}) +@Slf4j public class ThingsboardServerApplication { private static final String SPRING_CONFIG_NAME_KEY = "--spring.config.name"; private static final String DEFAULT_SPRING_CONFIG_PARAM = SPRING_CONFIG_NAME_KEY + "=" + "thingsboard"; + private static long startTs; + public static void main(String[] args) { + startTs = System.currentTimeMillis(); SpringApplication.run(ThingsboardServerApplication.class, updateArguments(args)); } @@ -45,4 +53,11 @@ public class ThingsboardServerApplication { } return args; } + + @AfterStartUp(order = Ordered.LOWEST_PRECEDENCE) + public void afterStartUp() { + long startupTimeMs = System.currentTimeMillis() - startTs; + log.info("Started ThingsBoard in {} seconds", TimeUnit.MILLISECONDS.toSeconds(startupTimeMs)); + } + }