Browse Source

refactoring: moved common method for retrieving key dictionary id to abstract class

pull/9850/head
dashevchenko 3 years ago
parent
commit
2375689974
  1. 1
      application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java
  2. 2
      application/src/main/java/org/thingsboard/server/service/install/SqlTsDatabaseUpgradeService.java
  3. 5
      application/src/main/java/org/thingsboard/server/service/install/update/DefaultCacheCleanupService.java
  4. 55
      dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java
  5. 47
      dao/src/main/java/org/thingsboard/server/dao/sql/attributes/JpaAttributeDao.java
  6. 54
      dao/src/main/java/org/thingsboard/server/dao/sqlts/BaseAbstractSqlTimeseriesDao.java

1
application/src/main/java/org/thingsboard/server/install/ThingsboardInstallService.java

@ -276,7 +276,6 @@ public class ThingsboardInstallService {
} else { } else {
log.info("Skipping images migration. Run the upgrade with fromVersion as '3.6.2-images' to migrate"); 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": case "3.6.3":
log.info("Upgrading ThingsBoard from version 3.6.3 to 3.7.0 ..."); log.info("Upgrading ThingsBoard from version 3.6.3 to 3.7.0 ...");
databaseEntitiesUpgradeService.upgradeDatabase("3.6.3"); databaseEntitiesUpgradeService.upgradeDatabase("3.6.3");

2
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"); loadSql(conn, LOAD_DROP_PARTITIONS_FUNCTIONS_SQL, "2.4.3");
} }
break; break;
case "3.6.2":
break;
default: default:
throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion); throw new RuntimeException("Unable to upgrade SQL database, unsupported fromVersion: " + fromVersion);
} }

5
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.Objects;
import java.util.Optional; 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.RESOURCE_INFO_CACHE;
import static org.thingsboard.server.common.data.CacheConstants.SECURITY_SETTINGS_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(SECURITY_SETTINGS_CACHE);
clearCacheByName(RESOURCE_INFO_CACHE); clearCacheByName(RESOURCE_INFO_CACHE);
break; 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: default:
//Do nothing, since cache cleanup is optional. //Do nothing, since cache cleanup is optional.
} }

55
dao/src/main/java/org/thingsboard/server/dao/sql/JpaAbstractDaoListeningExecutorService.java

@ -16,17 +16,29 @@
package org.thingsboard.server.dao.sql; package org.thingsboard.server.dao.sql;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.jdbc.core.JdbcTemplate; 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 javax.sql.DataSource;
import java.sql.SQLException; import java.sql.SQLException;
import java.sql.SQLWarning; import java.sql.SQLWarning;
import java.sql.Statement; 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 @Slf4j
public abstract class JpaAbstractDaoListeningExecutorService { public abstract class JpaAbstractDaoListeningExecutorService {
private final ConcurrentMap<String, Integer> keyDictionaryMap = new ConcurrentHashMap<>();
protected static final ReentrantLock creationLock = new ReentrantLock();
@Autowired @Autowired
protected JpaExecutorService service; protected JpaExecutorService service;
@ -36,6 +48,49 @@ public abstract class JpaAbstractDaoListeningExecutorService {
@Autowired @Autowired
protected JdbcTemplate jdbcTemplate; protected JdbcTemplate jdbcTemplate;
@Autowired
protected KeyDictionaryRepository keyDictionaryRepository;
protected Integer getOrSaveKeyId(String strKey) {
Integer keyId = keyDictionaryMap.get(strKey);
if (keyId == null) {
Optional<KeyDictionaryEntry> 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 { protected void printWarnings(Statement statement) throws SQLException {
SQLWarning warnings = statement.getWarnings(); SQLWarning warnings = statement.getWarnings();
if (warnings != null) { if (warnings != null) {

47
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.PostConstruct;
import jakarta.annotation.PreDestroy; import jakarta.annotation.PreDestroy;
import lombok.extern.slf4j.Slf4j; import lombok.extern.slf4j.Slf4j;
import org.hibernate.exception.ConstraintViolationException;
import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.annotation.Value;
import org.springframework.dao.DataIntegrityViolationException;
import org.springframework.stereotype.Component; import org.springframework.stereotype.Component;
import org.thingsboard.server.common.data.AttributeScope; import org.thingsboard.server.common.data.AttributeScope;
import org.thingsboard.server.common.data.id.DeviceProfileId; 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.attributes.AttributesDao;
import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey; import org.thingsboard.server.dao.model.sql.AttributeKvCompositeKey;
import org.thingsboard.server.dao.model.sql.AttributeKvEntity; 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.model.sqlts.dictionary.KeyDictionaryEntry;
import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService; import org.thingsboard.server.dao.sql.JpaAbstractDaoListeningExecutorService;
import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent; import org.thingsboard.server.dao.sql.ScheduledLogExecutorComponent;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueParams;
import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper; import org.thingsboard.server.dao.sql.TbSqlBlockingQueueWrapper;
import org.thingsboard.server.dao.sqlts.dictionary.KeyDictionaryRepository;
import org.thingsboard.server.dao.util.SqlDao; import org.thingsboard.server.dao.util.SqlDao;
import java.util.ArrayList; import java.util.ArrayList;
@ -51,9 +47,6 @@ import java.util.Collection;
import java.util.Comparator; import java.util.Comparator;
import java.util.List; import java.util.List;
import java.util.Optional; 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.function.Function;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -65,9 +58,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
@Autowired @Autowired
ScheduledLogExecutorComponent logExecutor; ScheduledLogExecutorComponent logExecutor;
@Autowired
private KeyDictionaryRepository keyDictionaryRepository;
@Autowired @Autowired
private AttributeKvRepository attributeKvRepository; private AttributeKvRepository attributeKvRepository;
@ -92,9 +82,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
@Value("${sql.batch_sort:true}") @Value("${sql.batch_sort:true}")
private boolean batchSortEnabled; private boolean batchSortEnabled;
private final ConcurrentMap<String, Integer> attributeDictionaryMap = new ConcurrentHashMap<>();
private static final ReentrantLock attributeCreationLock = new ReentrantLock();
private TbSqlBlockingQueueWrapper<AttributeKvEntity> queue; private TbSqlBlockingQueueWrapper<AttributeKvEntity> queue;
@PostConstruct @PostConstruct
@ -214,40 +201,6 @@ public class JpaAttributeDao extends JpaAbstractDaoListeningExecutorService impl
attributeKey); attributeKey);
} }
private Integer getOrSaveKeyId(String attributeKey) {
Integer keyId = attributeDictionaryMap.get(attributeKey);
if (keyId == null) {
Optional<KeyDictionaryEntry> 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) { private String getKey(Integer attributeKey) {
Optional<KeyDictionaryEntry> byKeyId = keyDictionaryRepository.findByKeyId(attributeKey); Optional<KeyDictionaryEntry> byKeyId = keyDictionaryRepository.findByKeyId(attributeKey);
return byKeyId.map(KeyDictionaryEntry::getKey).orElse(null); return byKeyId.map(KeyDictionaryEntry::getKey).orElse(null);

54
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.Futures;
import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListenableFuture;
import lombok.extern.slf4j.Slf4j; 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.ReadTsKvQuery;
import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult; import org.thingsboard.server.common.data.kv.ReadTsKvQueryResult;
import org.thingsboard.server.dao.DaoUtil; import org.thingsboard.server.dao.DaoUtil;
import org.thingsboard.server.dao.model.sql.AbstractTsKvEntity; 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.sql.JpaAbstractDaoListeningExecutorService;
import org.thingsboard.server.dao.sqlts.dictionary.KeyDictionaryRepository;
import jakarta.annotation.Nullable; import jakarta.annotation.Nullable;
import java.util.List; import java.util.List;
import java.util.Objects; import java.util.Objects;
import java.util.Optional; 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; import java.util.stream.Collectors;
@Slf4j @Slf4j
public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService { public abstract class BaseAbstractSqlTimeseriesDao extends JpaAbstractDaoListeningExecutorService {
private final ConcurrentMap<String, Integer> 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<KeyDictionaryEntry> 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<ReadTsKvQueryResult> getReadTsKvQueryResultFuture(ReadTsKvQuery query, ListenableFuture<List<Optional<? extends AbstractTsKvEntity>>> future) { protected ListenableFuture<ReadTsKvQueryResult> getReadTsKvQueryResultFuture(ReadTsKvQuery query, ListenableFuture<List<Optional<? extends AbstractTsKvEntity>>> future) {
return Futures.transform(future, new Function<>() { return Futures.transform(future, new Function<>() {
@Nullable @Nullable

Loading…
Cancel
Save