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)); + } + } 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..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 @@ -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; + } + + protected 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(); + } + } + } 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/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(); + } + } + }