From 237568997491b967e4b9b73be8e0c1ba9761b3ca Mon Sep 17 00:00:00 2001 From: dashevchenko Date: Thu, 4 Jan 2024 18:55:00 +0200 Subject: [PATCH] refactoring: moved common method for retrieving key dictionary id to abstract class --- .../install/ThingsboardInstallService.java | 1 - .../install/SqlTsDatabaseUpgradeService.java | 2 - .../update/DefaultCacheCleanupService.java | 5 ++ ...paAbstractDaoListeningExecutorService.java | 55 +++++++++++++++++++ .../dao/sql/attributes/JpaAttributeDao.java | 47 ---------------- .../sqlts/BaseAbstractSqlTimeseriesDao.java | 54 ------------------ 6 files changed, 60 insertions(+), 104 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java index ef5f3edd58..730ec0ec97 100644 --- a/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java +++ b/application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java @@ -276,7 +276,6 @@ public class ThingsboardInstallService { } else { log.info("Skipping images migration. Run the upgrade with fromVersion as '3.6.2-images' to migrate"); } - //TODO DON'T FORGET to update switch statement in the CacheCleanupService if you need to clear the cache case "3.6.3": log.info("Upgrading ThingsBoard from version 3.6.3 to 3.7.0 ..."); databaseEntitiesUpgradeService.upgradeDatabase("3.6.3"); diff --git a/application/src/main/java/org/thingsboard/server/service/install/SqlTsDatabaseUpgradeService.java b/application/src/main/java/org/thingsboard/server/service/install/SqlTsDatabaseUpgradeService.java index b61da45631..6400463e11 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/SqlTsDatabaseUpgradeService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/SqlTsDatabaseUpgradeService.java @@ -213,8 +213,6 @@ public class SqlTsDatabaseUpgradeService extends AbstractSqlTsDatabaseUpgradeSer loadSql(conn, LOAD_DROP_PARTITIONS_FUNCTIONS_SQL, "2.4.3"); } break; - case "3.6.2": - break; default: throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); } diff --git a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java index 350005536a..e401ec2360 100644 --- a/application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java +++ b/application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java @@ -27,6 +27,7 @@ import org.springframework.stereotype.Service; import java.util.Objects; import java.util.Optional; +import static org.thingsboard.server.common.data.CacheConstants.ATTRIBUTES_CACHE; import static org.thingsboard.server.common.data.CacheConstants.RESOURCE_INFO_CACHE; import static org.thingsboard.server.common.data.CacheConstants.SECURITY_SETTINGS_CACHE; @@ -94,6 +95,10 @@ public class DefaultCacheCleanupService implements CacheCleanupService { clearCacheByName(SECURITY_SETTINGS_CACHE); clearCacheByName(RESOURCE_INFO_CACHE); break; + case "3.6.3": + log.info("Clearing cache to upgrade from version 3.6.3 to 3.7.0"); + clearCacheByName(ATTRIBUTES_CACHE); + break; default: //Do nothing, since cache cleanup is optional. } diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java index 5646aa1b6c..3a10b7f619 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java @@ -16,17 +16,29 @@ package org.thingsboard.server.dao.sql; import lombok.extern.slf4j.Slf4j; +import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.dao.DataIntegrityViolationException; import org.springframework.jdbc.core.JdbcTemplate; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; +import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; +import org.thingsboard.server.dao.sqlts.dictionary.KeyDictionaryRepository; import javax.sql.DataSource; import java.sql.SQLException; import java.sql.SQLWarning; import java.sql.Statement; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.locks.ReentrantLock; @Slf4j public abstract class JpaAbstractDaoListeningExecutorService { + private final ConcurrentMap keyDictionaryMap = new ConcurrentHashMap<>(); + protected static final ReentrantLock creationLock = new ReentrantLock(); + @Autowired protected JpaExecutorService service; @@ -36,6 +48,49 @@ public abstract class JpaAbstractDaoListeningExecutorService { @Autowired protected JdbcTemplate jdbcTemplate; + @Autowired + protected KeyDictionaryRepository keyDictionaryRepository; + + protected Integer getOrSaveKeyId(String strKey) { + Integer keyId = keyDictionaryMap.get(strKey); + if (keyId == null) { + Optional tsKvDictionaryOptional; + tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); + if (tsKvDictionaryOptional.isEmpty()) { + creationLock.lock(); + try { + keyId = keyDictionaryMap.get(strKey); + if (keyId != null) { + return keyId; + } + tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); + if (tsKvDictionaryOptional.isEmpty()) { + KeyDictionaryEntry keyDictionaryEntry = new KeyDictionaryEntry(); + keyDictionaryEntry.setKey(strKey); + try { + KeyDictionaryEntry saved = keyDictionaryRepository.save(keyDictionaryEntry); + keyDictionaryMap.put(saved.getKey(), saved.getKeyId()); + keyId = saved.getKeyId(); + } catch (DataIntegrityViolationException | ConstraintViolationException e) { + tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); + KeyDictionaryEntry dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get KeyDictionaryEntry entity from DB!")); + keyDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); + keyId = dictionary.getKeyId(); + } + } else { + keyId = tsKvDictionaryOptional.get().getKeyId(); + } + } finally { + creationLock.unlock(); + } + } else { + keyId = tsKvDictionaryOptional.get().getKeyId(); + keyDictionaryMap.put(strKey, keyId); + } + } + return keyId; + } + protected void printWarnings(Statement statement) throws SQLException { SQLWarning warnings = statement.getWarnings(); if (warnings != null) { diff --git a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java index 01149d3d49..54d0bf9495 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java @@ -22,10 +22,8 @@ import com.google.common.util.concurrent.MoreExecutors; import jakarta.annotation.PostConstruct; import jakarta.annotation.PreDestroy; import lombok.extern.slf4j.Slf4j; -import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Value; -import org.springframework.dao.DataIntegrityViolationException; import org.springframework.stereotype.Component; import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.id.DeviceProfileId; @@ -37,13 +35,11 @@ import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.attributes.AttributesDao; import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey; import org.thingsboard.server.dao.model.sql.AttributeKvEntity; -import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; -import org.thingsboard.server.dao.sqlts.dictionary.KeyDictionaryRepository; import org.thingsboard.server.dao.util.SqlDao; import java.util.ArrayList; @@ -51,9 +47,6 @@ import java.util.Collection; import java.util.Comparator; import java.util.List; import java.util.Optional; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.locks.ReentrantLock; import java.util.function.Function; import java.util.stream.Collectors; @@ -65,9 +58,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl @Autowired ScheduledLogExecutorComponent logExecutor; - @Autowired - private KeyDictionaryRepository keyDictionaryRepository; - @Autowired private AttributeKvRepository attributeKvRepository; @@ -92,9 +82,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl @Value("${sql.batch_sort:true}") private boolean batchSortEnabled; - private final ConcurrentMap attributeDictionaryMap = new ConcurrentHashMap<>(); - private static final ReentrantLock attributeCreationLock = new ReentrantLock(); - private TbSqlBlockingQueueWrapper queue; @PostConstruct @@ -214,40 +201,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl attributeKey); } - private Integer getOrSaveKeyId(String attributeKey) { - Integer keyId = attributeDictionaryMap.get(attributeKey); - if (keyId == null) { - Optional byIdOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(attributeKey)); - if (byIdOptional.isEmpty()) { - attributeCreationLock.lock(); - try { - byIdOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(attributeKey)); - if (byIdOptional.isEmpty()) { - KeyDictionaryEntry attributeKvDictionaryEntry = new KeyDictionaryEntry(); - attributeKvDictionaryEntry.setKey(attributeKey); - try { - KeyDictionaryEntry saved = keyDictionaryRepository.save(attributeKvDictionaryEntry); - attributeDictionaryMap.put(saved.getKey(), saved.getKeyId()); - keyId = saved.getKeyId(); - } catch (DataIntegrityViolationException | ConstraintViolationException e) { - byIdOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(attributeKey)); - KeyDictionaryEntry dictionary = byIdOptional.orElseThrow(() -> new RuntimeException("Failed to get AttributeKvDictionary entity from DB!")); - attributeDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); - keyId = dictionary.getKeyId(); - } - } else { - keyId = byIdOptional.get().getKeyId(); - } - } finally { - attributeCreationLock.unlock(); - } - } else { - keyId = byIdOptional.get().getKeyId(); - attributeDictionaryMap.put(attributeKey, keyId); - } - } - return keyId; - } private String getKey(Integer attributeKey) { Optional byKeyId = keyDictionaryRepository.findByKeyId(attributeKey); return byKeyId.map(KeyDictionaryEntry::getKey).orElse(null); diff --git a/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java b/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java index e018cc1ef9..56f21aee39 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java +++ b/dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java @@ -19,75 +19,21 @@ import com.google.common.base.Function; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; -import org.hibernate.exception.ConstraintViolationException; -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.dao.DataIntegrityViolationException; import org.thingsboard.server.common.data.kv.ReadTsKvQuery; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; -import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryEntry; -import org.thingsboard.server.dao.model.sqlts.dictionary.KeyDictionaryCompositeKey; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; -import org.thingsboard.server.dao.sqlts.dictionary.KeyDictionaryRepository; import jakarta.annotation.Nullable; import java.util.List; import java.util.Objects; import java.util.Optional; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; @Slf4j public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService { - private final ConcurrentMap tsKvDictionaryMap = new ConcurrentHashMap<>(); - protected static final ReentrantLock tsCreationLock = new ReentrantLock(); - @Autowired - protected KeyDictionaryRepository keyDictionaryRepository; - - protected Integer getOrSaveKeyId(String strKey) { - Integer keyId = tsKvDictionaryMap.get(strKey); - if (keyId == null) { - Optional tsKvDictionaryOptional; - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - tsCreationLock.lock(); - try { - keyId = tsKvDictionaryMap.get(strKey); - if (keyId != null) { - return keyId; - } - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - if (tsKvDictionaryOptional.isEmpty()) { - KeyDictionaryEntry keyDictionaryEntry = new KeyDictionaryEntry(); - keyDictionaryEntry.setKey(strKey); - try { - KeyDictionaryEntry saved = keyDictionaryRepository.save(keyDictionaryEntry); - tsKvDictionaryMap.put(saved.getKey(), saved.getKeyId()); - keyId = saved.getKeyId(); - } catch (DataIntegrityViolationException | ConstraintViolationException e) { - tsKvDictionaryOptional = keyDictionaryRepository.findById(new KeyDictionaryCompositeKey(strKey)); - KeyDictionaryEntry dictionary = tsKvDictionaryOptional.orElseThrow(() -> new RuntimeException("Failed to get TsKvDictionary entity from DB!")); - tsKvDictionaryMap.put(dictionary.getKey(), dictionary.getKeyId()); - keyId = dictionary.getKeyId(); - } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - } - } finally { - tsCreationLock.unlock(); - } - } else { - keyId = tsKvDictionaryOptional.get().getKeyId(); - tsKvDictionaryMap.put(strKey, keyId); - } - } - return keyId; - } - protected ListenableFuture getReadTsKvQueryResultFuture(ReadTsKvQuery query, ListenableFuture>> future) { return Futures.transform(future, new Function<>() { @Nullable