Browse Source

Merge pull request #10813 from thingsboard/feature/startup-performance

Startup performance improvements
pull/10841/head
Andrew Shvayka 2 years ago
committed by GitHub
parent
commit
4db74f5f81
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 15
      application/src/main/java/org/thingsboard/server/ThingsboardServerApplication.java
  2. 45
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaAdmin.java
  3. 12
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaConsumerStatsService.java
  4. 35
      common/queue/src/main/java/org/thingsboard/server/queue/kafka/TbKafkaSettings.java
  5. 3
      dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java
  6. 4
      dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java
  7. 3
      dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java
  8. 3
      dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java
  9. 3
      dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java
  10. 3
      dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java
  11. 15
      dao/src/test/java/org/thingsboard/server/dao/service/EntityDaoRegistryTest.java

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

45
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<String, String> topicConfigs;
private final Set<String> topics = ConcurrentHashMap.newKeySet();
private final int numPartitions;
private volatile Set<String> topics;
private final short replicationFactor;
public TbKafkaAdmin(TbKafkaSettings settings, Map<String, String> 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<String> topics = getTopics();
if (topics.contains(topic)) {
return;
}
@ -86,12 +80,13 @@ public class TbKafkaAdmin implements TbQueueAdmin {
@Override
public void deleteTopic(String topic) {
Set<String> 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<String> 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() {
}
}

12
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<String, byte[]> 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<TopicPartition, OffsetAndMetadata> groupOffsets = adminClient.listConsumerGroupOffsets(groupId).partitionsToOffsetAndMetadata()
Map<TopicPartition, OffsetAndMetadata> groupOffsets = kafkaSettings.getAdminClient().listConsumerGroupOffsets(groupId).partitionsToOffsetAndMetadata()
.get(statsConfig.getKafkaResponseTimeoutMs(), TimeUnit.MILLISECONDS);
Map<TopicPartition, Long> 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();
}

35
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<String, List<TbProperty>> 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();
}
}
}

3
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 {

4
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 {

3
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

3
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

3
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

3
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

15
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<JpaRepository<?, ?>> 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();
}
}
}

Loading…
Cancel
Save