Browse Source

Merge pull request #11368 from thingsboard/feature/dedicated-datasource

Dedicated datasource for events and audit logs
pull/11431/head
Andrew Shvayka 2 years ago
committed by GitHub
parent
commit
b62f1871b5
No known key found for this signature in database GPG Key ID: B5690EEEBB952194
  1. 23
      application/src/main/resources/thingsboard.yml
  2. 28
      application/src/test/java/org/thingsboard/server/controller/AuditLogControllerTest.java
  3. 36
      application/src/test/java/org/thingsboard/server/controller/AuditLogControllerTest_DedicatedEventsDataSource.java
  4. 37
      dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java
  5. 8
      dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java
  6. 11
      dao/src/main/java/org/thingsboard/server/dao/config/DedicatedEventsDataSource.java
  7. 91
      dao/src/main/java/org/thingsboard/server/dao/config/DedicatedEventsJpaDaoConfig.java
  8. 26
      dao/src/main/java/org/thingsboard/server/dao/config/DefaultDataSource.java
  9. 15
      dao/src/main/java/org/thingsboard/server/dao/config/DefaultDedicatedJpaDaoConfig.java
  10. 121
      dao/src/main/java/org/thingsboard/server/dao/config/JpaDaoConfig.java
  11. 4
      dao/src/main/java/org/thingsboard/server/dao/config/SqlTsDaoConfig.java
  12. 4
      dao/src/main/java/org/thingsboard/server/dao/config/SqlTsLatestDaoConfig.java
  13. 4
      dao/src/main/java/org/thingsboard/server/dao/config/TimescaleDaoConfig.java
  14. 4
      dao/src/main/java/org/thingsboard/server/dao/config/TimescaleTsLatestDaoConfig.java
  15. 1
      dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java
  16. 1
      dao/src/main/java/org/thingsboard/server/dao/model/sql/AuditLogEntity.java
  17. 19
      dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java
  18. 87
      dao/src/main/java/org/thingsboard/server/dao/sql/audit/DedicatedJpaAuditLogDao.java
  19. 35
      dao/src/main/java/org/thingsboard/server/dao/sql/audit/JpaAuditLogDao.java
  20. 36
      dao/src/main/java/org/thingsboard/server/dao/sql/event/DedicatedEventInsertRepository.java
  21. 45
      dao/src/main/java/org/thingsboard/server/dao/sql/event/DedicatedJpaEventDao.java
  22. 13
      dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java
  23. 73
      dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java
  24. 101
      dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java
  25. 55
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/DedicatedEventsSqlPartitioningRepository.java
  26. 16
      dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java
  27. 6
      dao/src/test/java/org/thingsboard/server/dao/AbstractDaoServiceTest.java
  28. 7
      dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java
  29. 1
      dao/src/test/java/org/thingsboard/server/dao/PostgreSqlInitializer.java
  30. 28
      dao/src/test/java/org/thingsboard/server/dao/service/event/sql/EventServiceSqlTest_DedicatedEventsDataSource.java

23
application/src/main/resources/thingsboard.yml

@ -773,7 +773,28 @@ spring:
leakDetectionThreshold: "${SPRING_DATASOURCE_HIKARI_LEAK_DETECTION_THRESHOLD:0}"
# This property increases the number of connections in the pool as demand increases. At the same time, the property ensures that the pool doesn't grow to the point of exhausting a system's resources, which ultimately affects an application's performance and availability
maximumPoolSize: "${SPRING_DATASOURCE_MAXIMUM_POOL_SIZE:16}"
registerMbeans: "${SPRING_DATASOURCE_HIKARI_REGISTER_MBEANS:false}" # true - enable MBean to diagnose pools state via JMX
# Enable MBean to diagnose pools state via JMX
registerMbeans: "${SPRING_DATASOURCE_HIKARI_REGISTER_MBEANS:false}"
events:
# Enable dedicated datasource (a separate database) for events and audit logs.
# Before enabling this, make sure you have set up the following tables in the new DB:
# error_event, lc_event, rule_chain_debug_event, rule_node_debug_event, stats_event, audit_log
enabled: "${SPRING_DEDICATED_EVENTS_DATASOURCE_ENABLED:false}"
# Database driver for Spring JPA for events datasource
driverClassName: "${SPRING_EVENTS_DATASOURCE_DRIVER_CLASS_NAME:org.postgresql.Driver}"
# Database connection URL for events datasource
url: "${SPRING_EVENTS_DATASOURCE_URL:jdbc:postgresql://localhost:5432/thingsboard_events}"
# Database username for events datasource
username: "${SPRING_EVENTS_DATASOURCE_USERNAME:postgres}"
# Database user password for events datasource
password: "${SPRING_EVENTS_DATASOURCE_PASSWORD:postgres}"
hikari:
# This property controls the amount of time that a connection can be out of the pool before a message is logged indicating a possible connection leak for events datasource. A value of 0 means leak detection is disabled
leakDetectionThreshold: "${SPRING_EVENTS_DATASOURCE_HIKARI_LEAK_DETECTION_THRESHOLD:0}"
# This property increases the number of connections in the pool as demand increases for events datasource. At the same time, the property ensures that the pool doesn't grow to the point of exhausting a system's resources, which ultimately affects an application's performance and availability
maximumPoolSize: "${SPRING_EVENTS_DATASOURCE_MAXIMUM_POOL_SIZE:16}"
# Enable MBean to diagnose pools state via JMX for events datasource
registerMbeans: "${SPRING_EVENTS_DATASOURCE_HIKARI_REGISTER_MBEANS:false}"
# Audit log parameters
audit-log:

28
application/src/test/java/org/thingsboard/server/controller/AuditLogControllerTest.java

@ -17,6 +17,7 @@ package org.thingsboard.server.controller;
import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.fasterxml.jackson.core.type.TypeReference;
import lombok.Getter;
import lombok.extern.slf4j.Slf4j;
import org.junit.After;
import org.junit.Assert;
@ -64,6 +65,7 @@ public class AuditLogControllerTest extends AbstractControllerTest {
@Autowired
private AuditLogDao auditLogDao;
@Getter
@SpyBean
private SqlPartitioningRepository partitioningRepository;
@SpyBean
@ -183,12 +185,12 @@ public class AuditLogControllerTest extends AbstractControllerTest {
@Test
public void whenSavingNewAuditLog_thenCheckAndCreatePartitionIfNotExists() throws ParseException {
long entityTs = ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2024-01-01T01:43:11Z").getTime();
reset(partitioningRepository);
reset(getPartitioningRepository());
AuditLog auditLog = createAuditLog(ActionType.LOGIN, tenantAdminUserId, entityTs);
verify(partitioningRepository).createPartitionIfNotExists(eq("audit_log"), eq(auditLog.getCreatedTime()), eq(partitionDurationInMs));
verify(getPartitioningRepository()).createPartitionIfNotExists(eq("audit_log"), eq(auditLog.getCreatedTime()), eq(partitionDurationInMs));
List<Long> partitions = partitioningRepository.fetchPartitions("audit_log");
assertThat(partitions).contains(partitioningRepository.calculatePartitionStartTime(auditLog.getCreatedTime(), partitionDurationInMs));
List<Long> partitions = getPartitioningRepository().fetchPartitions("audit_log");
assertThat(partitions).contains(getPartitioningRepository().calculatePartitionStartTime(auditLog.getCreatedTime(), partitionDurationInMs));
}
@Test
@ -197,15 +199,15 @@ public class AuditLogControllerTest extends AbstractControllerTest {
final long oldAuditLogTs = ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2020-10-01T00:00:00Z").getTime();
final long currentTimeMillis = oldAuditLogTs + TimeUnit.SECONDS.toMillis(auditLogsTtlInSec) * 2;
final long partitionStartTs = partitioningRepository.calculatePartitionStartTime(oldAuditLogTs, partitionDurationInMs);
partitioningRepository.createPartitionIfNotExists("audit_log", oldAuditLogTs, partitionDurationInMs);
List<Long> partitions = partitioningRepository.fetchPartitions("audit_log");
final long partitionStartTs = getPartitioningRepository().calculatePartitionStartTime(oldAuditLogTs, partitionDurationInMs);
getPartitioningRepository().createPartitionIfNotExists("audit_log", oldAuditLogTs, partitionDurationInMs);
List<Long> partitions = getPartitioningRepository().fetchPartitions("audit_log");
assertThat(partitions).contains(partitionStartTs);
willReturn(currentTimeMillis).given(auditLogsCleanUpService).getCurrentTimeMillis();
auditLogsCleanUpService.cleanUp();
partitions = partitioningRepository.fetchPartitions("audit_log");
partitions = getPartitioningRepository().fetchPartitions("audit_log");
assertThat(partitions).as("partitions cleared").doesNotContain(partitionStartTs);
assertThat(partitions).as("only newer partitions left").allSatisfy(partitionsStart -> {
long partitionEndTs = partitionsStart + partitionDurationInMs;
@ -218,18 +220,18 @@ public class AuditLogControllerTest extends AbstractControllerTest {
// creating partition bigger than sql.audit_logs.partition_size
long entityTs = ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2022-04-29T07:43:11Z").getTime();
//the partition 7 days is overlapping default partition size 1 day, use in the far past to not affect other tests
partitioningRepository.createPartitionIfNotExists("audit_log", entityTs, TimeUnit.DAYS.toMillis(7));
List<Long> partitions = partitioningRepository.fetchPartitions("audit_log");
getPartitioningRepository().createPartitionIfNotExists("audit_log", entityTs, TimeUnit.DAYS.toMillis(7));
List<Long> partitions = getPartitioningRepository().fetchPartitions("audit_log");
log.warn("entityTs [{}], fetched partitions {}", entityTs, partitions);
assertThat(partitions).contains(ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2022-04-28T00:00:00Z").getTime());
partitioningRepository.cleanupPartitionsCache("audit_log", entityTs, 0);
getPartitioningRepository().cleanupPartitionsCache("audit_log", entityTs, 0);
assertDoesNotThrow(() -> {
// expecting partition overlap error on partition save
createAuditLog(ActionType.LOGIN, tenantAdminUserId, entityTs);
});
assertThat(partitioningRepository.fetchPartitions("audit_log"))
.contains(ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2022-04-28T00:00:00Z").getTime());;
assertThat(getPartitioningRepository().fetchPartitions("audit_log"))
.contains(ISO_8601_EXTENDED_DATETIME_TIME_ZONE_FORMAT.parse("2022-04-28T00:00:00Z").getTime());
}
private AuditLog createAuditLog(ActionType actionType, EntityId entityId, long entityTs) {

36
application/src/test/java/org/thingsboard/server/controller/AuditLogControllerTest_DedicatedEventsDataSource.java

@ -0,0 +1,36 @@
/**
* 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.controller;
import lombok.Getter;
import org.springframework.boot.test.mock.mockito.SpyBean;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.dao.service.DaoSqlTest;
import org.thingsboard.server.dao.sqlts.insert.sql.DedicatedEventsSqlPartitioningRepository;
@DaoSqlTest
@TestPropertySource(properties = {
"spring.datasource.events.enabled=true",
"spring.datasource.events.url=${spring.datasource.url}",
"spring.datasource.events.driverClassName=${spring.datasource.driverClassName}",
})
public class AuditLogControllerTest_DedicatedEventsDataSource extends AuditLogControllerTest {
@Getter
@SpyBean
private DedicatedEventsSqlPartitioningRepository partitioningRepository;
}

37
dao/src/main/java/org/thingsboard/server/dao/JpaDaoConfig.java

@ -1,37 +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;
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;
/**
* @author Valerii Sosliuk
*/
@Configuration
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sql", "org.thingsboard.server.dao.attributes", "org.thingsboard.server.dao.cache", "org.thingsboard.server.cache"})
@EnableJpaRepositories(value = "org.thingsboard.server.dao.sql", bootstrapMode = BootstrapMode.LAZY)
@EntityScan("org.thingsboard.server.dao.model.sql")
@EnableTransactionManagement
public class JpaDaoConfig {
}

8
dao/src/main/java/org/thingsboard/server/dao/audit/AuditLogServiceImpl.java

@ -408,8 +408,12 @@ public class AuditLogServiceImpl implements AuditLogService {
}
return executor.submit(() -> {
AuditLog auditLog = auditLogDao.save(tenantId, auditLogEntry);
auditLogSink.logAction(auditLog);
try {
AuditLog auditLog = auditLogDao.save(tenantId, auditLogEntry);
auditLogSink.logAction(auditLog);
} catch (Throwable e) {
log.error("[{}] Failed to save audit log: {}", tenantId, auditLogEntry, e);
}
return null;
});
}

11
dao/src/main/java/org/thingsboard/server/dao/sql/event/EventCleanupRepository.java → dao/src/main/java/org/thingsboard/server/dao/config/DedicatedEventsDataSource.java

@ -13,11 +13,14 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao.sql.event;
package org.thingsboard.server.dao.config;
public interface EventCleanupRepository {
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
void cleanupEvents(long eventExpTime, boolean debug);
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
void migrateEvents(long regularEventTs, long debugEventTs);
@Retention(RetentionPolicy.RUNTIME)
@ConditionalOnProperty(value = "spring.datasource.events.enabled", havingValue = "true")
public @interface DedicatedEventsDataSource {
}

91
dao/src/main/java/org/thingsboard/server/dao/config/DedicatedEventsJpaDaoConfig.java

@ -0,0 +1,91 @@
/**
* 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.config;
import com.zaxxer.hikari.HikariDataSource;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.jdbc.DataSourceProperties;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.boot.orm.jpa.EntityManagerFactoryBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
import org.springframework.data.repository.config.BootstrapMode;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.orm.jpa.JpaTransactionManager;
import org.springframework.orm.jpa.LocalContainerEntityManagerFactoryBean;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.model.sql.AuditLogEntity;
import org.thingsboard.server.dao.model.sql.ErrorEventEntity;
import org.thingsboard.server.dao.model.sql.LifecycleEventEntity;
import org.thingsboard.server.dao.model.sql.RuleChainDebugEventEntity;
import org.thingsboard.server.dao.model.sql.RuleNodeDebugEventEntity;
import org.thingsboard.server.dao.model.sql.StatisticsEventEntity;
import javax.sql.DataSource;
import java.util.Objects;
@DedicatedEventsDataSource
@Configuration
@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sql.event", "org.thingsboard.server.dao.sql.audit"},
bootstrapMode = BootstrapMode.LAZY,
entityManagerFactoryRef = "eventsEntityManagerFactory", transactionManagerRef = "eventsTransactionManager")
public class DedicatedEventsJpaDaoConfig {
public static final String EVENTS_PERSISTENCE_UNIT = "events";
public static final String EVENTS_DATA_SOURCE = EVENTS_PERSISTENCE_UNIT + "DataSource";
public static final String EVENTS_TRANSACTION_MANAGER = EVENTS_PERSISTENCE_UNIT + "TransactionManager";
public static final String EVENTS_TRANSACTION_TEMPLATE = EVENTS_PERSISTENCE_UNIT + "TransactionTemplate";
public static final String EVENTS_JDBC_TEMPLATE = EVENTS_PERSISTENCE_UNIT + "JdbcTemplate";
@Bean
@ConfigurationProperties("spring.datasource.events")
public DataSourceProperties eventsDataSourceProperties() {
return new DataSourceProperties();
}
@ConfigurationProperties(prefix = "spring.datasource.events.hikari")
@Bean(EVENTS_DATA_SOURCE)
public DataSource eventsDataSource(@Qualifier("eventsDataSourceProperties") DataSourceProperties eventsDataSourceProperties) {
return eventsDataSourceProperties.initializeDataSourceBuilder().type(HikariDataSource.class).build();
}
@Bean
public LocalContainerEntityManagerFactoryBean eventsEntityManagerFactory(@Qualifier(EVENTS_DATA_SOURCE) DataSource eventsDataSource,
EntityManagerFactoryBuilder builder) {
return builder
.dataSource(eventsDataSource)
.packages(LifecycleEventEntity.class, StatisticsEventEntity.class, ErrorEventEntity.class, RuleNodeDebugEventEntity.class, RuleChainDebugEventEntity.class, AuditLogEntity.class)
.persistenceUnit(EVENTS_PERSISTENCE_UNIT)
.build();
}
@Bean(EVENTS_TRANSACTION_MANAGER)
public JpaTransactionManager eventsTransactionManager(@Qualifier("eventsEntityManagerFactory") LocalContainerEntityManagerFactoryBean eventsEntityManagerFactory) {
return new JpaTransactionManager(Objects.requireNonNull(eventsEntityManagerFactory.getObject()));
}
@Bean(EVENTS_TRANSACTION_TEMPLATE)
public TransactionTemplate eventsTransactionTemplate(@Qualifier(EVENTS_TRANSACTION_MANAGER) JpaTransactionManager eventsTransactionManager) {
return new TransactionTemplate(eventsTransactionManager);
}
@Bean(EVENTS_JDBC_TEMPLATE)
public JdbcTemplate eventsJdbcTemplate(@Qualifier(EVENTS_DATA_SOURCE) DataSource eventsDataSource) {
return new JdbcTemplate(eventsDataSource);
}
}

26
dao/src/main/java/org/thingsboard/server/dao/config/DefaultDataSource.java

@ -0,0 +1,26 @@
/**
* 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.config;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
@Retention(RetentionPolicy.RUNTIME)
@ConditionalOnProperty(value = "spring.datasource.events.enabled", havingValue = "false", matchIfMissing = true)
public @interface DefaultDataSource {
}

15
dao/src/main/java/org/thingsboard/server/dao/SqlTimeseriesDaoConfig.java → dao/src/main/java/org/thingsboard/server/dao/config/DefaultDedicatedJpaDaoConfig.java

@ -13,22 +13,15 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
package org.thingsboard.server.dao.config;
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;
@DefaultDataSource
@Configuration
@TbAutoConfiguration
@ComponentScan({"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 {
@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sql.event", "org.thingsboard.server.dao.sql.audit"}, bootstrapMode = BootstrapMode.LAZY)
public class DefaultDedicatedJpaDaoConfig {
}

121
dao/src/main/java/org/thingsboard/server/dao/config/JpaDaoConfig.java

@ -0,0 +1,121 @@
/**
* 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.config;
import com.zaxxer.hikari.HikariDataSource;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.autoconfigure.jdbc.DataSourceProperties;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.boot.orm.jpa.EntityManagerFactoryBuilder;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.ComponentScan;
import org.springframework.context.annotation.ComponentScan.Filter;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.FilterType;
import org.springframework.context.annotation.Primary;
import org.springframework.data.jpa.repository.config.EnableJpaRepositories;
import org.springframework.data.repository.config.BootstrapMode;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.jdbc.core.namedparam.NamedParameterJdbcTemplate;
import org.springframework.orm.jpa.JpaTransactionManager;
import org.springframework.orm.jpa.LocalContainerEntityManagerFactoryBean;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.sql.audit.AuditLogRepository;
import org.thingsboard.server.dao.sql.event.EventRepository;
import org.thingsboard.server.dao.util.TbAutoConfiguration;
import javax.sql.DataSource;
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
@Configuration
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sql", "org.thingsboard.server.dao.attributes", "org.thingsboard.server.dao.sqlts.dictionary", "org.thingsboard.server.dao.cache", "org.thingsboard.server.cache"})
@EnableJpaRepositories(value = {"org.thingsboard.server.dao.sql", "org.thingsboard.server.dao.sqlts.dictionary"},
excludeFilters = @Filter(type = FilterType.ASSIGNABLE_TYPE, classes = {EventRepository.class, AuditLogRepository.class}),
bootstrapMode = BootstrapMode.LAZY)
public class JpaDaoConfig {
@Bean
@ConfigurationProperties("spring.datasource")
public DataSourceProperties dataSourceProperties() {
return new DataSourceProperties();
}
@Primary
@ConfigurationProperties(prefix = "spring.datasource.hikari")
@Bean
public DataSource dataSource(@Qualifier("dataSourceProperties") DataSourceProperties dataSourceProperties) {
return dataSourceProperties.initializeDataSourceBuilder().type(HikariDataSource.class).build();
}
@Primary
@Bean
public LocalContainerEntityManagerFactoryBean entityManagerFactory(@Qualifier("dataSource") DataSource dataSource,
EntityManagerFactoryBuilder builder,
@Autowired(required = false) SqlTsLatestDaoConfig tsLatestDaoConfig,
@Autowired(required = false) SqlTsDaoConfig tsDaoConfig,
@Autowired(required = false) TimescaleDaoConfig timescaleDaoConfig,
@Autowired(required = false) TimescaleTsLatestDaoConfig timescaleTsLatestDaoConfig) {
List<String> packages = new ArrayList<>();
packages.add("org.thingsboard.server.dao.model.sql");
packages.add("org.thingsboard.server.dao.model.sqlts.dictionary");
if (tsLatestDaoConfig != null) {
packages.add("org.thingsboard.server.dao.model.sqlts.latest");
}
if (tsDaoConfig != null) {
packages.add("org.thingsboard.server.dao.model.sqlts.ts");
}
if (timescaleDaoConfig != null) {
packages.add("org.thingsboard.server.dao.model.sqlts.timescale");
}
if (timescaleTsLatestDaoConfig != null) {
packages.add("org.thingsboard.server.dao.model.sqlts.latest");
}
return builder
.dataSource(dataSource)
.packages(packages.toArray(String[]::new))
.persistenceUnit("default")
.build();
}
@Primary
@Bean
public JpaTransactionManager transactionManager(@Qualifier("entityManagerFactory") LocalContainerEntityManagerFactoryBean entityManagerFactory) {
return new JpaTransactionManager(Objects.requireNonNull(entityManagerFactory.getObject()));
}
@Primary
@Bean
public TransactionTemplate transactionTemplate(@Qualifier("transactionManager") JpaTransactionManager transactionManager) {
return new TransactionTemplate(transactionManager);
}
@Primary
@Bean
public JdbcTemplate jdbcTemplate(@Qualifier("dataSource") DataSource dataSource) {
return new JdbcTemplate(dataSource);
}
@Primary
@Bean
public NamedParameterJdbcTemplate namedParameterJdbcTemplate(@Qualifier("dataSource") DataSource dataSource) {
return new NamedParameterJdbcTemplate(dataSource);
}
}

4
dao/src/main/java/org/thingsboard/server/dao/SqlTsDaoConfig.java → dao/src/main/java/org/thingsboard/server/dao/config/SqlTsDaoConfig.java

@ -13,9 +13,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
package org.thingsboard.server.dao.config;
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;
@ -28,7 +27,6 @@ import org.thingsboard.server.dao.util.TbAutoConfiguration;
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sqlts.sql", "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
public class SqlTsDaoConfig {

4
dao/src/main/java/org/thingsboard/server/dao/SqlTsLatestDaoConfig.java → dao/src/main/java/org/thingsboard/server/dao/config/SqlTsLatestDaoConfig.java

@ -13,9 +13,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
package org.thingsboard.server.dao.config;
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;
@ -28,7 +27,6 @@ import org.thingsboard.server.dao.util.TbAutoConfiguration;
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sqlts.sql"})
@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
public class SqlTsLatestDaoConfig {

4
dao/src/main/java/org/thingsboard/server/dao/TimescaleDaoConfig.java → dao/src/main/java/org/thingsboard/server/dao/config/TimescaleDaoConfig.java

@ -13,9 +13,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
package org.thingsboard.server.dao.config;
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;
@ -28,7 +27,6 @@ import org.thingsboard.server.dao.util.TimescaleDBTsDao;
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sqlts.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
public class TimescaleDaoConfig {

4
dao/src/main/java/org/thingsboard/server/dao/TimescaleTsLatestDaoConfig.java → dao/src/main/java/org/thingsboard/server/dao/config/TimescaleTsLatestDaoConfig.java

@ -13,9 +13,8 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package org.thingsboard.server.dao;
package org.thingsboard.server.dao.config;
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;
@ -28,7 +27,6 @@ import org.thingsboard.server.dao.util.TimescaleDBTsLatestDao;
@TbAutoConfiguration
@ComponentScan({"org.thingsboard.server.dao.sqlts.timescale"})
@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
public class TimescaleTsLatestDaoConfig {

1
dao/src/main/java/org/thingsboard/server/dao/event/EventDao.java

@ -101,5 +101,4 @@ public interface EventDao {
*/
void removeEvents(UUID tenantId, UUID entityId, EventFilter eventFilter, Long startTime, Long endTime);
void migrateEvents(long regularEventTs, long debugEventTs);
}

1
dao/src/main/java/org/thingsboard/server/dao/model/sql/AuditLogEntity.java

@ -149,4 +149,5 @@ public class AuditLogEntity extends BaseSqlEntity<AuditLog> implements BaseEntit
auditLog.setActionFailureDetails(this.actionFailureDetails);
return auditLog;
}
}

19
dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDao.java

@ -55,10 +55,6 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D>
@PersistenceContext
private EntityManager entityManager;
protected abstract Class<E> getEntityClass();
protected abstract JpaRepository<E, UUID> getRepository();
@Override
@Transactional
public D save(TenantId tenantId, D domain) {
@ -85,6 +81,7 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D>
}
protected E doSave(E entity, boolean isNew) {
EntityManager entityManager = getEntityManager();
if (isNew) {
if (entity instanceof HasVersion versionedEntity) {
versionedEntity.setVersion(1L);
@ -176,11 +173,23 @@ public abstract class JpaAbstractDao<E extends BaseEntity<D>, D>
}
query += " ORDER BY id LIMIT ?";
return jdbcTemplate.queryForList(query, UUID.class, params);
return getJdbcTemplate().queryForList(query, UUID.class, params);
}
protected String getTenantIdColumn() {
return ModelConstants.TENANT_ID_COLUMN;
}
protected EntityManager getEntityManager() {
return entityManager;
}
protected JdbcTemplate getJdbcTemplate() {
return jdbcTemplate;
}
protected abstract Class<E> getEntityClass();
protected abstract JpaRepository<E, UUID> getRepository();
}

87
dao/src/main/java/org/thingsboard/server/dao/sql/audit/DedicatedJpaAuditLogDao.java

@ -0,0 +1,87 @@
/**
* 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.sql.audit;
import jakarta.persistence.EntityManager;
import jakarta.persistence.PersistenceContext;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.common.data.audit.AuditLog;
import org.thingsboard.server.common.data.id.TenantId;
import org.thingsboard.server.dao.config.DedicatedEventsDataSource;
import org.thingsboard.server.dao.sqlts.insert.sql.DedicatedEventsSqlPartitioningRepository;
import org.thingsboard.server.dao.util.SqlDao;
import java.util.Collection;
import java.util.UUID;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_JDBC_TEMPLATE;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_PERSISTENCE_UNIT;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_TRANSACTION_MANAGER;
@DedicatedEventsDataSource
@Component
@SqlDao
public class DedicatedJpaAuditLogDao extends JpaAuditLogDao {
@Autowired
@Qualifier(EVENTS_JDBC_TEMPLATE)
private JdbcTemplate jdbcTemplate;
@PersistenceContext(unitName = EVENTS_PERSISTENCE_UNIT)
private EntityManager entityManager;
public DedicatedJpaAuditLogDao(AuditLogRepository auditLogRepository, DedicatedEventsSqlPartitioningRepository partitioningRepository) {
super(auditLogRepository, partitioningRepository);
}
@Transactional(transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public AuditLog save(TenantId tenantId, AuditLog domain) {
return super.save(tenantId, domain);
}
@Transactional(transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public AuditLog saveAndFlush(TenantId tenantId, AuditLog domain) {
return super.saveAndFlush(tenantId, domain);
}
@Transactional(transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public void removeById(TenantId tenantId, UUID id) {
super.removeById(tenantId, id);
}
@Transactional(transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public void removeAllByIds(Collection<UUID> ids) {
super.removeAllByIds(ids);
}
@Override
protected EntityManager getEntityManager() {
return entityManager;
}
@Override
protected JdbcTemplate getJdbcTemplate() {
return jdbcTemplate;
}
}

35
dao/src/main/java/org/thingsboard/server/dao/sql/audit/JpaAuditLogDao.java

@ -19,7 +19,6 @@ import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.audit.ActionType;
import org.thingsboard.server.common.data.audit.AuditLog;
@ -30,7 +29,7 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.audit.AuditLogDao;
import org.thingsboard.server.dao.model.ModelConstants;
import org.thingsboard.server.dao.config.DefaultDataSource;
import org.thingsboard.server.dao.model.sql.AuditLogEntity;
import org.thingsboard.server.dao.sql.JpaPartitionedAbstractDao;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
@ -40,6 +39,9 @@ import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import static org.thingsboard.server.dao.model.ModelConstants.AUDIT_LOG_TABLE_NAME;
@DefaultDataSource
@Component
@SqlDao
@RequiredArgsConstructor
@ -48,24 +50,9 @@ public class JpaAuditLogDao extends JpaPartitionedAbstractDao<AuditLogEntity, Au
private final AuditLogRepository auditLogRepository;
private final SqlPartitioningRepository partitioningRepository;
private final JdbcTemplate jdbcTemplate;
@Value("${sql.audit_logs.partition_size:168}")
private int partitionSizeInHours;
@Value("${sql.ttl.audit_logs.ttl:0}")
private long ttlInSec;
private static final String TABLE_NAME = ModelConstants.AUDIT_LOG_TABLE_NAME;
@Override
protected Class<AuditLogEntity> getEntityClass() {
return AuditLogEntity.class;
}
@Override
protected JpaRepository<AuditLogEntity, UUID> getRepository() {
return auditLogRepository;
}
@Override
public PageData<AuditLog> findAuditLogsByTenantIdAndEntityId(UUID tenantId, EntityId entityId, List<ActionType> actionTypes, TimePageLink pageLink) {
@ -124,12 +111,22 @@ public class JpaAuditLogDao extends JpaPartitionedAbstractDao<AuditLogEntity, Au
@Override
public void cleanUpAuditLogs(long expTime) {
partitioningRepository.dropPartitionsBefore(TABLE_NAME, expTime, TimeUnit.HOURS.toMillis(partitionSizeInHours));
partitioningRepository.dropPartitionsBefore(AUDIT_LOG_TABLE_NAME, expTime, TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
@Override
public void createPartition(AuditLogEntity entity) {
partitioningRepository.createPartitionIfNotExists(TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
partitioningRepository.createPartitionIfNotExists(AUDIT_LOG_TABLE_NAME, entity.getCreatedTime(), TimeUnit.HOURS.toMillis(partitionSizeInHours));
}
@Override
protected Class<AuditLogEntity> getEntityClass() {
return AuditLogEntity.class;
}
@Override
protected JpaRepository<AuditLogEntity, UUID> getRepository() {
return auditLogRepository;
}
}

36
dao/src/main/java/org/thingsboard/server/dao/sql/event/DedicatedEventInsertRepository.java

@ -0,0 +1,36 @@
/**
* 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.sql.event;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.support.TransactionTemplate;
import org.thingsboard.server.dao.config.DedicatedEventsDataSource;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_JDBC_TEMPLATE;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_TRANSACTION_TEMPLATE;
@DedicatedEventsDataSource
@Repository
public class DedicatedEventInsertRepository extends EventInsertRepository {
public DedicatedEventInsertRepository(@Qualifier(EVENTS_JDBC_TEMPLATE) JdbcTemplate jdbcTemplate,
@Qualifier(EVENTS_TRANSACTION_TEMPLATE) TransactionTemplate transactionTemplate) {
super(jdbcTemplate, transactionTemplate);
}
}

45
dao/src/main/java/org/thingsboard/server/dao/sql/event/DedicatedJpaEventDao.java

@ -0,0 +1,45 @@
/**
* 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.sql.event;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.config.DedicatedEventsDataSource;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
import org.thingsboard.server.dao.sqlts.insert.sql.DedicatedEventsSqlPartitioningRepository;
import org.thingsboard.server.dao.util.SqlDao;
@DedicatedEventsDataSource
@Component
@SqlDao
public class DedicatedJpaEventDao extends JpaBaseEventDao {
public DedicatedJpaEventDao(EventPartitionConfiguration partitionConfiguration,
DedicatedEventsSqlPartitioningRepository partitioningRepository,
LifecycleEventRepository lcEventRepository,
StatisticsEventRepository statsEventRepository,
ErrorEventRepository errorEventRepository,
DedicatedEventInsertRepository eventInsertRepository,
RuleNodeDebugEventRepository ruleNodeDebugEventRepository,
RuleChainDebugEventRepository ruleChainDebugEventRepository,
ScheduledLogExecutorComponent logExecutor,
StatsFactory statsFactory) {
super(partitionConfiguration, partitioningRepository, lcEventRepository, statsEventRepository,
errorEventRepository, eventInsertRepository, ruleNodeDebugEventRepository,
ruleChainDebugEventRepository, logExecutor, statsFactory);
}
}

13
dao/src/main/java/org/thingsboard/server/dao/sql/event/EventInsertRepository.java

@ -16,7 +16,7 @@
package org.thingsboard.server.dao.sql.event;
import jakarta.annotation.PostConstruct;
import org.springframework.beans.factory.annotation.Autowired;
import lombok.RequiredArgsConstructor;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.jdbc.core.BatchPreparedStatementSetter;
import org.springframework.jdbc.core.JdbcTemplate;
@ -32,6 +32,7 @@ import org.thingsboard.server.common.data.event.LifecycleEvent;
import org.thingsboard.server.common.data.event.RuleChainDebugEvent;
import org.thingsboard.server.common.data.event.RuleNodeDebugEvent;
import org.thingsboard.server.common.data.event.StatisticsEvent;
import org.thingsboard.server.dao.config.DefaultDataSource;
import org.thingsboard.server.dao.util.SqlDao;
import java.sql.PreparedStatement;
@ -44,9 +45,11 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
@DefaultDataSource
@Repository
@Transactional
@SqlDao
@RequiredArgsConstructor
public class EventInsertRepository {
private static final ThreadLocal<Pattern> PATTERN_THREAD_LOCAL = ThreadLocal.withInitial(() -> Pattern.compile(String.valueOf(Character.MIN_VALUE)));
@ -55,11 +58,8 @@ public class EventInsertRepository {
private final Map<EventType, String> insertStmtMap = new ConcurrentHashMap<>();
@Autowired
protected JdbcTemplate jdbcTemplate;
@Autowired
private TransactionTemplate transactionTemplate;
private final JdbcTemplate jdbcTemplate;
private final TransactionTemplate transactionTemplate;
@Value("${sql.remove_null_chars:true}")
private boolean removeNullChars;
@ -236,4 +236,5 @@ public class EventInsertRepository {
}
return strValue;
}
}

73
dao/src/main/java/org/thingsboard/server/dao/sql/event/JpaBaseEventDao.java

@ -19,8 +19,8 @@ import com.datastax.oss.driver.api.core.uuid.Uuids;
import com.google.common.util.concurrent.ListenableFuture;
import jakarta.annotation.PostConstruct;
import jakarta.annotation.PreDestroy;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.StringUtils;
@ -38,6 +38,7 @@ import org.thingsboard.server.common.data.page.PageData;
import org.thingsboard.server.common.data.page.TimePageLink;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.config.DefaultDataSource;
import org.thingsboard.server.dao.event.EventDao;
import org.thingsboard.server.dao.model.sql.EventEntity;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
@ -54,46 +55,23 @@ import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Function;
/**
* Created by Valerii Sosliuk on 5/3/2017.
*/
@Slf4j
@DefaultDataSource
@Component
@SqlDao
@RequiredArgsConstructor
@Slf4j
public class JpaBaseEventDao implements EventDao {
@Autowired
private EventPartitionConfiguration partitionConfiguration;
@Autowired
private SqlPartitioningRepository partitioningRepository;
@Autowired
private LifecycleEventRepository lcEventRepository;
@Autowired
private StatisticsEventRepository statsEventRepository;
@Autowired
private ErrorEventRepository errorEventRepository;
@Autowired
private EventInsertRepository eventInsertRepository;
@Autowired
private EventCleanupRepository eventCleanupRepository;
@Autowired
private RuleNodeDebugEventRepository ruleNodeDebugEventRepository;
@Autowired
private RuleChainDebugEventRepository ruleChainDebugEventRepository;
@Autowired
ScheduledLogExecutorComponent logExecutor;
@Autowired
private StatsFactory statsFactory;
private final EventPartitionConfiguration partitionConfiguration;
private final SqlPartitioningRepository partitioningRepository;
private final LifecycleEventRepository lcEventRepository;
private final StatisticsEventRepository statsEventRepository;
private final ErrorEventRepository errorEventRepository;
private final EventInsertRepository eventInsertRepository;
private final RuleNodeDebugEventRepository ruleNodeDebugEventRepository;
private final RuleChainDebugEventRepository ruleChainDebugEventRepository;
private final ScheduledLogExecutorComponent logExecutor;
private final StatsFactory statsFactory;
@Value("${sql.events.batch_size:10000}")
private int batchSize;
@ -223,11 +201,6 @@ public class JpaBaseEventDao implements EventDao {
}
}
@Override
public void migrateEvents(long regularEventTs, long debugEventTs) {
eventCleanupRepository.migrateEvents(regularEventTs, debugEventTs);
}
private PageData<? extends Event> findEventByFilter(UUID tenantId, UUID entityId, RuleChainDebugEventFilter eventFilter, TimePageLink pageLink) {
return DaoUtil.toPageData(
ruleChainDebugEventRepository.findEvents(
@ -402,7 +375,7 @@ public class JpaBaseEventDao implements EventDao {
if (regularEventExpTs > 0) {
log.info("Going to cleanup regular events with exp time: {}", regularEventExpTs);
if (cleanupDb) {
eventCleanupRepository.cleanupEvents(regularEventExpTs, false);
cleanupEvents(regularEventExpTs, false);
} else {
cleanupPartitionsCache(regularEventExpTs, false);
}
@ -410,13 +383,25 @@ public class JpaBaseEventDao implements EventDao {
if (debugEventExpTs > 0) {
log.info("Going to cleanup debug events with exp time: {}", debugEventExpTs);
if (cleanupDb) {
eventCleanupRepository.cleanupEvents(debugEventExpTs, true);
cleanupEvents(debugEventExpTs, true);
} else {
cleanupPartitionsCache(debugEventExpTs, true);
}
}
}
private void cleanupEvents(long eventExpTime, boolean debug) {
for (EventType eventType : EventType.values()) {
if (eventType.isDebug() == debug) {
cleanupPartitions(eventType, eventExpTime);
}
}
}
private void cleanupPartitions(EventType eventType, long eventExpTime) {
partitioningRepository.dropPartitionsBefore(eventType.getTable(), eventExpTime, partitionConfiguration.getPartitionSizeInMs(eventType));
}
private void cleanupPartitionsCache(long expTime, boolean isDebug) {
for (EventType eventType : EventType.values()) {
if (eventType.isDebug() == isDebug) {

101
dao/src/main/java/org/thingsboard/server/dao/sql/event/SqlEventCleanupRepository.java

@ -1,101 +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.sql.event;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.dao.DataAccessException;
import org.springframework.stereotype.Repository;
import org.thingsboard.server.common.data.event.EventType;
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService;
import org.thingsboard.server.dao.sqlts.insert.sql.SqlPartitioningRepository;
import java.util.concurrent.TimeUnit;
@Slf4j
@Repository
public class SqlEventCleanupRepository extends JpaAbstractDaoListeningExecutorService implements EventCleanupRepository {
@Autowired
private EventPartitionConfiguration partitionConfiguration;
@Autowired
private SqlPartitioningRepository partitioningRepository;
@Override
public void cleanupEvents(long eventExpTime, boolean debug) {
for (EventType eventType : EventType.values()) {
if (eventType.isDebug() == debug) {
cleanupEvents(eventType, eventExpTime);
}
}
}
@Override
public void migrateEvents(long regularEventTs, long debugEventTs) {
regularEventTs = Math.max(regularEventTs, 1480982400000L);
debugEventTs = Math.max(debugEventTs, 1480982400000L);
callMigrateFunctionByPartitions("regular", "migrate_regular_events", regularEventTs, partitionConfiguration.getRegularPartitionSizeInHours());
callMigrateFunctionByPartitions("debug", "migrate_debug_events", debugEventTs, partitionConfiguration.getDebugPartitionSizeInHours());
try {
jdbcTemplate.execute("DROP PROCEDURE IF EXISTS migrate_regular_events(bigint, bigint, int)");
jdbcTemplate.execute("DROP PROCEDURE IF EXISTS migrate_debug_events(bigint, bigint, int)");
jdbcTemplate.execute("DROP TABLE IF EXISTS event");
} catch (DataAccessException e) {
log.error("Error occurred during drop of the `events` table", e);
throw e;
}
}
private void callMigrateFunctionByPartitions(String logTag, String functionName, long startTs, int partitionSizeInHours) {
long currentTs = System.currentTimeMillis();
var regularPartitionStepInMs = TimeUnit.HOURS.toMillis(partitionSizeInHours);
long numberOfPartitions = (currentTs - startTs) / regularPartitionStepInMs;
if (numberOfPartitions > 1000) {
log.error("Please adjust your {} events partitioning configuration. " +
"Configuration with partition size of {} hours and corresponding TTL will use {} (>1000) partitions which is not recommended!",
logTag, partitionSizeInHours, numberOfPartitions);
throw new RuntimeException("Please adjust your " + logTag + " events partitioning configuration. " +
"Configuration with partition size of " + partitionSizeInHours + " hours and corresponding TTL will use " +
+numberOfPartitions + " (>1000) partitions which is not recommended!");
}
while (startTs < currentTs) {
var endTs = startTs + regularPartitionStepInMs;
log.info("Migrate {} events for time period: [{},{}]", logTag, startTs, endTs);
callMigrateFunction(functionName, startTs, startTs + regularPartitionStepInMs, partitionSizeInHours);
startTs = endTs;
}
log.info("Migrate {} events done.", logTag);
}
private void callMigrateFunction(String functionName, long startTs, long endTs, int partitionSizeInHours) {
try {
jdbcTemplate.update("CALL " + functionName + "(?, ?, ?)", startTs, endTs, partitionSizeInHours);
} catch (DataAccessException e) {
if (e.getMessage() == null || !e.getMessage().contains("relation \"event\" does not exist")) {
log.error("[{}] SQLException occurred during execution of {} with parameters {} and {}", functionName, startTs, partitionSizeInHours, e);
throw new RuntimeException(e);
}
}
}
private void cleanupEvents(EventType eventType, long eventExpTime) {
partitioningRepository.dropPartitionsBefore(eventType.getTable(), eventExpTime, partitionConfiguration.getPartitionSizeInMs(eventType));
}
}

55
dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/DedicatedEventsSqlPartitioningRepository.java

@ -0,0 +1,55 @@
/**
* 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.sqlts.insert.sql;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository;
import org.springframework.transaction.annotation.Propagation;
import org.springframework.transaction.annotation.Transactional;
import org.thingsboard.server.dao.config.DedicatedEventsDataSource;
import org.thingsboard.server.dao.timeseries.SqlPartition;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_JDBC_TEMPLATE;
import static org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig.EVENTS_TRANSACTION_MANAGER;
@DedicatedEventsDataSource
@Repository
public class DedicatedEventsSqlPartitioningRepository extends SqlPartitioningRepository {
@Autowired
@Qualifier(EVENTS_JDBC_TEMPLATE)
private JdbcTemplate jdbcTemplate;
@Transactional(propagation = Propagation.NOT_SUPPORTED, transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public void save(SqlPartition partition) {
super.save(partition);
}
@Transactional(propagation = Propagation.NOT_SUPPORTED, transactionManager = EVENTS_TRANSACTION_MANAGER)
@Override
public void createPartitionIfNotExists(String table, long entityTs, long partitionDurationMs) {
super.createPartitionIfNotExists(table, entityTs, partitionDurationMs);
}
@Override
protected JdbcTemplate getJdbcTemplate() {
return jdbcTemplate;
}
}

16
dao/src/main/java/org/thingsboard/server/dao/sqlts/insert/sql/SqlPartitioningRepository.java

@ -19,6 +19,7 @@ import lombok.extern.slf4j.Slf4j;
import org.apache.commons.lang3.StringUtils;
import org.apache.commons.lang3.exception.ExceptionUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.context.annotation.Primary;
import org.springframework.dao.DataAccessException;
import org.springframework.jdbc.core.JdbcTemplate;
import org.springframework.stereotype.Repository;
@ -32,6 +33,7 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.locks.ReentrantLock;
@Primary
@Repository
@Slf4j
public class SqlPartitioningRepository {
@ -49,7 +51,7 @@ public class SqlPartitioningRepository {
@Transactional(propagation = Propagation.NOT_SUPPORTED)
public void save(SqlPartition partition) {
jdbcTemplate.execute(partition.getQuery());
getJdbcTemplate().execute(partition.getQuery());
}
@Transactional(propagation = Propagation.NOT_SUPPORTED) // executing non-transactionally, so that parent transaction is not aborted on partition save error
@ -119,8 +121,8 @@ public class SqlPartitioningRepository {
String dropStmtStr = "DROP TABLE " + tablePartition;
try {
jdbcTemplate.execute(detachPsqlStmtStr);
jdbcTemplate.execute(dropStmtStr);
getJdbcTemplate().execute(detachPsqlStmtStr);
getJdbcTemplate().execute(dropStmtStr);
return true;
} catch (DataAccessException e) {
log.error("[{}] Error occurred trying to detach and drop the partition {} ", table, partitionTs, e);
@ -134,7 +136,7 @@ public class SqlPartitioningRepository {
public List<Long> fetchPartitions(String table) {
List<Long> partitions = new ArrayList<>();
List<String> partitionsTables = jdbcTemplate.queryForList(SELECT_PARTITIONS_STMT, String.class, table);
List<String> partitionsTables = getJdbcTemplate().queryForList(SELECT_PARTITIONS_STMT, String.class, table);
for (String partitionTableName : partitionsTables) {
String partitionTsStr = partitionTableName.substring(table.length() + 1);
try {
@ -153,7 +155,7 @@ public class SqlPartitioningRepository {
private synchronized int getCurrentServerVersion() {
if (currentServerVersion == null) {
try {
currentServerVersion = jdbcTemplate.queryForObject("SELECT current_setting('server_version_num')", Integer.class);
currentServerVersion = getJdbcTemplate().queryForObject("SELECT current_setting('server_version_num')", Integer.class);
} catch (Exception e) {
log.warn("Error occurred during fetch of the server version", e);
}
@ -164,4 +166,8 @@ public class SqlPartitioningRepository {
return currentServerVersion;
}
protected JdbcTemplate getJdbcTemplate() {
return jdbcTemplate;
}
}

6
dao/src/test/java/org/thingsboard/server/dao/AbstractDaoServiceTest.java

@ -25,10 +25,14 @@ import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.test.context.support.DependencyInjectionTestExecutionListener;
import org.springframework.test.context.support.DirtiesContextTestExecutionListener;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig;
import org.thingsboard.server.dao.config.JpaDaoConfig;
import org.thingsboard.server.dao.config.SqlTsDaoConfig;
import org.thingsboard.server.dao.config.SqlTsLatestDaoConfig;
import org.thingsboard.server.dao.service.DaoSqlTest;
@RunWith(SpringRunner.class)
@ContextConfiguration(classes = {JpaDaoConfig.class, SqlTsDaoConfig.class, SqlTsLatestDaoConfig.class, SqlTimeseriesDaoConfig.class})
@ContextConfiguration(classes = {JpaDaoConfig.class, SqlTsDaoConfig.class, SqlTsLatestDaoConfig.class, DedicatedEventsJpaDaoConfig.class})
@DaoSqlTest
@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
@TestExecutionListeners({

7
dao/src/test/java/org/thingsboard/server/dao/AbstractJpaDaoTest.java

@ -24,13 +24,18 @@ import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.test.context.support.DependencyInjectionTestExecutionListener;
import org.springframework.test.context.support.DirtiesContextTestExecutionListener;
import org.thingsboard.server.common.stats.StatsFactory;
import org.thingsboard.server.dao.config.DedicatedEventsJpaDaoConfig;
import org.thingsboard.server.dao.config.DefaultDedicatedJpaDaoConfig;
import org.thingsboard.server.dao.config.JpaDaoConfig;
import org.thingsboard.server.dao.config.SqlTsDaoConfig;
import org.thingsboard.server.dao.config.SqlTsLatestDaoConfig;
import org.thingsboard.server.dao.service.DaoSqlTest;
/**
* Created by Valerii Sosliuk on 4/22/2017.
*/
@RunWith(SpringRunner.class)
@ContextConfiguration(classes = {JpaDaoConfig.class, SqlTsDaoConfig.class, SqlTsLatestDaoConfig.class, SqlTimeseriesDaoConfig.class})
@ContextConfiguration(classes = {JpaDaoConfig.class, SqlTsDaoConfig.class, SqlTsLatestDaoConfig.class, DedicatedEventsJpaDaoConfig.class, DefaultDedicatedJpaDaoConfig.class})
@DaoSqlTest
@TestExecutionListeners({
DependencyInjectionTestExecutionListener.class,

1
dao/src/test/java/org/thingsboard/server/dao/PostgreSqlInitializer.java

@ -63,4 +63,5 @@ public class PostgreSqlInitializer {
throw new RuntimeException("Unable to clean up the Postgres database. Reason: " + e.getMessage(), e);
}
}
}

28
dao/src/test/java/org/thingsboard/server/dao/service/event/sql/EventServiceSqlTest_DedicatedEventsDataSource.java

@ -0,0 +1,28 @@
/**
* 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.event.sql;
import org.springframework.test.context.TestPropertySource;
import org.thingsboard.server.dao.service.DaoSqlTest;
@DaoSqlTest
@TestPropertySource(properties = {
"spring.datasource.events.enabled=true",
"spring.datasource.events.url=${spring.datasource.url}",
"spring.datasource.events.driverClassName=${spring.datasource.driverClassName}"
})
public class EventServiceSqlTest_DedicatedEventsDataSource extends EventServiceSqlTest {
}
Loading…
Cancel
Save