From ac67516dc209c33dddd3fc86c0371d638c117f9d Mon Sep 17 00:00:00 2001 From: Andrew Shvayka Date: Fri, 22 Nov 2019 15:47:54 +0200 Subject: [PATCH] Optimization of Device Creation and Lookup Performance --- .../state/DefaultDeviceStateService.java | 48 ++++++++++--------- .../server/dao/asset/BaseAssetService.java | 45 ++++++++++++----- .../device/DeviceCredentialsServiceImpl.java | 35 ++++++++++---- .../server/dao/device/DeviceServiceImpl.java | 43 ++++++++++++----- .../dao/entity/AbstractEntityService.java | 25 ++++++++++ .../resources/sql/schema-entities-idx.sql | 12 +++++ .../main/resources/sql/schema-entities.sql | 9 ++-- 7 files changed, 160 insertions(+), 57 deletions(-) diff --git a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java index aee2444b2e..5c83a5f89b 100644 --- a/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java +++ b/application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java @@ -40,6 +40,7 @@ import org.thingsboard.server.common.data.kv.AttributeKvEntry; import org.thingsboard.server.common.data.kv.BasicTsKvEntry; import org.thingsboard.server.common.data.kv.BooleanDataEntry; import org.thingsboard.server.common.data.kv.LongDataEntry; +import org.thingsboard.server.common.data.page.TextPageData; import org.thingsboard.server.common.data.page.TextPageLink; import org.thingsboard.server.common.msg.TbMsg; import org.thingsboard.server.common.msg.TbMsgDataType; @@ -93,7 +94,8 @@ public class DefaultDeviceStateService implements DeviceStateService { public static final String INACTIVITY_ALARM_TIME = "inactivityAlarmTime"; public static final String INACTIVITY_TIMEOUT = "inactivityTimeout"; - public static final List PERSISTENT_ATTRIBUTES = Arrays.asList(ACTIVITY_STATE, LAST_CONNECT_TIME, LAST_DISCONNECT_TIME, LAST_ACTIVITY_TIME, INACTIVITY_ALARM_TIME, INACTIVITY_TIMEOUT); + public static final List PERSISTENT_ATTRIBUTES = Arrays.asList(ACTIVITY_STATE, LAST_CONNECT_TIME, + LAST_DISCONNECT_TIME, LAST_ACTIVITY_TIME, INACTIVITY_ALARM_TIME, INACTIVITY_TIMEOUT); @Autowired private TenantService tenantService; @@ -129,17 +131,11 @@ public class DefaultDeviceStateService implements DeviceStateService { @Getter private boolean persistToTelemetry; -// TODO in v2.1 -// @Value("${state.defaultStatePersistenceIntervalInSec}") -// @Getter -// private long defaultStatePersistenceIntervalInSec; -// -// @Value("${state.defaultStatePersistencePack}") -// @Getter -// private long defaultStatePersistencePack; + @Value("${state.initFetchPackSize:1000}") + @Getter + private int initFetchPackSize; private ListeningScheduledExecutorService queueExecutor; - private ConcurrentMap> tenantDevices = new ConcurrentHashMap<>(); private ConcurrentMap deviceStates = new ConcurrentHashMap<>(); @@ -250,20 +246,28 @@ public class DefaultDeviceStateService implements DeviceStateService { } private void initStateFromDB() { - List tenants = tenantService.findTenants(new TextPageLink(Integer.MAX_VALUE)).getData(); - for (Tenant tenant : tenants) { - List> fetchFutures = new ArrayList<>(); - List devices = deviceService.findDevicesByTenantId(tenant.getId(), new TextPageLink(Integer.MAX_VALUE)).getData(); - for (Device device : devices) { - if (!routingService.resolveById(device.getId()).isPresent()) { - fetchFutures.add(fetchDeviceState(device)); + try { + List tenants = tenantService.findTenants(new TextPageLink(Integer.MAX_VALUE)).getData(); + for (Tenant tenant : tenants) { + List> fetchFutures = new ArrayList<>(); + TextPageLink pageLink = new TextPageLink(initFetchPackSize); + while (pageLink != null) { + TextPageData page = deviceService.findDevicesByTenantId(tenant.getId(), pageLink); + pageLink = page.getNextPageLink(); + for (Device device : page.getData()) { + if (!routingService.resolveById(device.getId()).isPresent()) { + fetchFutures.add(fetchDeviceState(device)); + } + } + try { + Futures.successfulAsList(fetchFutures).get().forEach(this::addDeviceUsingState); + } catch (InterruptedException | ExecutionException e) { + log.warn("Failed to init device state service from DB", e); + } } } - try { - Futures.successfulAsList(fetchFutures).get().forEach(this::addDeviceUsingState); - } catch (InterruptedException | ExecutionException e) { - log.warn("Failed to init device state service from DB", e); - } + } catch (Throwable t) { + log.warn("Failed to init device states from DB", t); } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java index cf679a103f..8cd3c476f8 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java @@ -20,6 +20,7 @@ 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.cache.Cache; import org.springframework.cache.CacheManager; @@ -28,6 +29,7 @@ import org.springframework.cache.annotation.Cacheable; import org.springframework.stereotype.Service; import org.springframework.util.StringUtils; import org.thingsboard.server.common.data.Customer; +import org.thingsboard.server.common.data.Device; import org.thingsboard.server.common.data.EntitySubtype; import org.thingsboard.server.common.data.EntityType; import org.thingsboard.server.common.data.EntityView; @@ -113,7 +115,22 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ public Asset saveAsset(Asset asset) { log.trace("Executing saveAsset [{}]", asset); assetValidator.validate(asset, Asset::getTenantId); - return assetDao.save(asset.getTenantId(), asset); + Asset savedAsset; + if (!sqlDatabaseUsed) { + savedAsset = assetDao.save(asset.getTenantId(), asset); + } else { + try { + savedAsset = assetDao.save(asset.getTenantId(), asset); + } catch (Exception t) { + ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); + if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("asset_name_unq_key")) { + throw new DataValidationException("Asset with such name already exists!"); + } else { + throw t; + } + } + } + return savedAsset; } @Override @@ -265,22 +282,26 @@ public class BaseAssetService extends AbstractEntityService implements AssetServ @Override protected void validateCreate(TenantId tenantId, Asset asset) { - assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent( - d -> { - throw new DataValidationException("Asset with such name already exists!"); - } - ); + if (!sqlDatabaseUsed) { + assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent( + d -> { + throw new DataValidationException("Asset with such name already exists!"); + } + ); + } } @Override protected void validateUpdate(TenantId tenantId, Asset asset) { - assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent( - d -> { - if (!d.getId().equals(asset.getId())) { - throw new DataValidationException("Asset with such name already exists!"); + if (!sqlDatabaseUsed) { + assetDao.findAssetsByTenantIdAndName(asset.getTenantId().getId(), asset.getName()).ifPresent( + d -> { + if (!d.getId().equals(asset.getId())) { + throw new DataValidationException("Asset with such name already exists!"); + } } - } - ); + ); + } } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java index 7c7910ba8f..c14d25b65b 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java @@ -17,6 +17,7 @@ package org.thingsboard.server.dao.device; import lombok.extern.slf4j.Slf4j; +import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.annotation.CacheEvict; import org.springframework.cache.annotation.Cacheable; @@ -29,6 +30,7 @@ import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.common.data.security.DeviceCredentials; import org.thingsboard.server.common.data.security.DeviceCredentialsType; import org.thingsboard.server.common.msg.EncryptionUtil; +import org.thingsboard.server.dao.entity.AbstractEntityService; import org.thingsboard.server.dao.exception.DataValidationException; import org.thingsboard.server.dao.service.DataValidator; @@ -38,7 +40,7 @@ import static org.thingsboard.server.dao.service.Validator.validateString; @Service @Slf4j -public class DeviceCredentialsServiceImpl implements DeviceCredentialsService { +public class DeviceCredentialsServiceImpl extends AbstractEntityService implements DeviceCredentialsService { @Autowired private DeviceCredentialsDao deviceCredentialsDao; @@ -78,7 +80,20 @@ public class DeviceCredentialsServiceImpl implements DeviceCredentialsService { } log.trace("Executing updateDeviceCredentials [{}]", deviceCredentials); credentialsValidator.validate(deviceCredentials, id -> tenantId); - return deviceCredentialsDao.save(tenantId, deviceCredentials); + if (!sqlDatabaseUsed) { + return deviceCredentialsDao.save(tenantId, deviceCredentials); + } else { + try { + return deviceCredentialsDao.save(tenantId, deviceCredentials); + } catch (Exception t) { + ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); + if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("device_credentials_id_unq_key")) { + throw new DataValidationException("Specified credentials are already registered!"); + } else { + throw t; + } + } + } } private void formatCertData(DeviceCredentials deviceCredentials) { @@ -100,9 +115,11 @@ public class DeviceCredentialsServiceImpl implements DeviceCredentialsService { @Override protected void validateCreate(TenantId tenantId, DeviceCredentials deviceCredentials) { - DeviceCredentials existingCredentialsEntity = deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()); - if (existingCredentialsEntity != null) { - throw new DataValidationException("Create of existent device credentials!"); + if (!sqlDatabaseUsed) { + DeviceCredentials existingCredentialsEntity = deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()); + if (existingCredentialsEntity != null) { + throw new DataValidationException("Create of existent device credentials!"); + } } } @@ -112,9 +129,11 @@ public class DeviceCredentialsServiceImpl implements DeviceCredentialsService { if (existingCredentials == null) { throw new DataValidationException("Unable to update non-existent device credentials!"); } - DeviceCredentials sameCredentialsId = deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()); - if (sameCredentialsId != null && !sameCredentialsId.getUuidId().equals(deviceCredentials.getUuidId())) { - throw new DataValidationException("Specified credentials are already registered!"); + if (!sqlDatabaseUsed) { + DeviceCredentials sameCredentialsId = deviceCredentialsDao.findByCredentialsId(tenantId, deviceCredentials.getCredentialsId()); + if (sameCredentialsId != null && !sameCredentialsId.getUuidId().equals(deviceCredentials.getUuidId())) { + throw new DataValidationException("Specified credentials are already registered!"); + } } } diff --git a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java index 872e03a387..a26e275a96 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java +++ b/dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java @@ -20,6 +20,7 @@ import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import lombok.extern.slf4j.Slf4j; import org.apache.commons.lang3.RandomStringUtils; +import org.hibernate.exception.ConstraintViolationException; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.cache.Cache; import org.springframework.cache.CacheManager; @@ -123,7 +124,21 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe public Device saveDevice(Device device) { log.trace("Executing saveDevice [{}]", device); deviceValidator.validate(device, Device::getTenantId); - Device savedDevice = deviceDao.save(device.getTenantId(), device); + Device savedDevice; + if (!sqlDatabaseUsed) { + savedDevice = deviceDao.save(device.getTenantId(), device); + } else { + try { + savedDevice = deviceDao.save(device.getTenantId(), device); + } catch (Exception t) { + ConstraintViolationException e = extractConstraintViolationException(t).orElse(null); + if (e != null && e.getConstraintName() != null && e.getConstraintName().equalsIgnoreCase("device_name_unq_key")) { + throw new DataValidationException("Device with such name already exists!"); + } else { + throw t; + } + } + } if (device.getId() == null) { DeviceCredentials deviceCredentials = new DeviceCredentials(); deviceCredentials.setDeviceId(new DeviceId(savedDevice.getUuidId())); @@ -296,22 +311,26 @@ public class DeviceServiceImpl extends AbstractEntityService implements DeviceSe @Override protected void validateCreate(TenantId tenantId, Device device) { - deviceDao.findDeviceByTenantIdAndName(device.getTenantId().getId(), device.getName()).ifPresent( - d -> { - throw new DataValidationException("Device with such name already exists!"); - } - ); + if (!sqlDatabaseUsed) { + deviceDao.findDeviceByTenantIdAndName(device.getTenantId().getId(), device.getName()).ifPresent( + d -> { + throw new DataValidationException("Device with such name already exists!"); + } + ); + } } @Override protected void validateUpdate(TenantId tenantId, Device device) { - deviceDao.findDeviceByTenantIdAndName(device.getTenantId().getId(), device.getName()).ifPresent( - d -> { - if (!d.getUuidId().equals(device.getUuidId())) { - throw new DataValidationException("Device with such name already exists!"); + if (!sqlDatabaseUsed) { + deviceDao.findDeviceByTenantIdAndName(device.getTenantId().getId(), device.getName()).ifPresent( + d -> { + if (!d.getUuidId().equals(device.getUuidId())) { + throw new DataValidationException("Device with such name already exists!"); + } } - } - ); + ); + } } @Override diff --git a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java index 5280637020..2fce1fd758 100644 --- a/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java +++ b/dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java @@ -16,20 +16,45 @@ package org.thingsboard.server.dao.entity; 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.thingsboard.server.common.data.id.EntityId; import org.thingsboard.server.common.data.id.TenantId; import org.thingsboard.server.dao.relation.RelationService; +import javax.annotation.PostConstruct; +import java.util.Optional; + @Slf4j public abstract class AbstractEntityService { @Autowired protected RelationService relationService; + @Value("${database.entities.type:sql}") + private String databaseType; + + protected boolean sqlDatabaseUsed; + + @PostConstruct + public void init() { + sqlDatabaseUsed = "sql".equalsIgnoreCase(databaseType); + } + protected void deleteEntityRelations(TenantId tenantId, EntityId entityId) { log.trace("Executing deleteEntityRelations [{}]", entityId); relationService.deleteEntityRelations(tenantId, entityId); } + protected Optional extractConstraintViolationException(Exception t) { + if (t instanceof ConstraintViolationException) { + return Optional.of ((ConstraintViolationException) t); + } else if (t.getCause() instanceof ConstraintViolationException) { + return Optional.of ((ConstraintViolationException) (t.getCause())); + } else { + return Optional.empty(); + } + } + } diff --git a/dao/src/main/resources/sql/schema-entities-idx.sql b/dao/src/main/resources/sql/schema-entities-idx.sql index 59785ae758..9809219ff9 100644 --- a/dao/src/main/resources/sql/schema-entities-idx.sql +++ b/dao/src/main/resources/sql/schema-entities-idx.sql @@ -21,3 +21,15 @@ CREATE INDEX IF NOT EXISTS idx_event_type_entity_id ON event(tenant_id, event_ty CREATE INDEX IF NOT EXISTS idx_relation_to_id ON relation(relation_type_group, to_type, to_id); CREATE INDEX IF NOT EXISTS idx_relation_from_id ON relation(relation_type_group, from_type, from_id); + +CREATE INDEX IF NOT EXISTS idx_device_customer_id ON device(tenant_id, customer_id); + +CREATE INDEX IF NOT EXISTS idx_device_customer_id_and_type ON device(tenant_id, customer_id, type); + +CREATE INDEX IF NOT EXISTS idx_device_type ON device(tenant_id, type); + +CREATE INDEX IF NOT EXISTS idx_asset_customer_id ON asset(tenant_id, customer_id); + +CREATE INDEX IF NOT EXISTS idx_asset_customer_id_and_type ON asset(tenant_id, customer_id, type); + +CREATE INDEX IF NOT EXISTS idx_asset_type ON asset(tenant_id, type); \ No newline at end of file diff --git a/dao/src/main/resources/sql/schema-entities.sql b/dao/src/main/resources/sql/schema-entities.sql index 089ed28afc..4b21d851a5 100644 --- a/dao/src/main/resources/sql/schema-entities.sql +++ b/dao/src/main/resources/sql/schema-entities.sql @@ -45,7 +45,8 @@ CREATE TABLE IF NOT EXISTS asset ( label varchar(255), search_text varchar(255), tenant_id varchar(31), - type varchar(255) + type varchar(255), + CONSTRAINT asset_name_unq_key UNIQUE (tenant_id, name) ); CREATE TABLE IF NOT EXISTS audit_log ( @@ -120,7 +121,8 @@ CREATE TABLE IF NOT EXISTS device ( name varchar(255), label varchar(255), search_text varchar(255), - tenant_id varchar(31) + tenant_id varchar(31), + CONSTRAINT device_name_unq_key UNIQUE (tenant_id, name) ); CREATE TABLE IF NOT EXISTS device_credentials ( @@ -128,7 +130,8 @@ CREATE TABLE IF NOT EXISTS device_credentials ( credentials_id varchar, credentials_type varchar(255), credentials_value varchar, - device_id varchar(31) + device_id varchar(31), + CONSTRAINT device_credentials_id_unq_key UNIQUE (credentials_id) ); CREATE TABLE IF NOT EXISTS event (