Browse Source

Optimization of Device Creation and Lookup Performance

pull/2205/head
Andrew Shvayka 7 years ago
parent
commit
ac67516dc2
  1. 48
      application/src/main/java/org/thingsboard/server/service/state/DefaultDeviceStateService.java
  2. 45
      dao/src/main/java/org/thingsboard/server/dao/asset/BaseAssetService.java
  3. 35
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceCredentialsServiceImpl.java
  4. 43
      dao/src/main/java/org/thingsboard/server/dao/device/DeviceServiceImpl.java
  5. 25
      dao/src/main/java/org/thingsboard/server/dao/entity/AbstractEntityService.java
  6. 12
      dao/src/main/resources/sql/schema-entities-idx.sql
  7. 9
      dao/src/main/resources/sql/schema-entities.sql

48
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<String> PERSISTENT_ATTRIBUTES = Arrays.asList(ACTIVITY_STATE, LAST_CONNECT_TIME, LAST_DISCONNECT_TIME, LAST_ACTIVITY_TIME, INACTIVITY_ALARM_TIME, INACTIVITY_TIMEOUT);
public static final List<String> 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<TenantId, Set<DeviceId>> tenantDevices = new ConcurrentHashMap<>();
private ConcurrentMap<DeviceId, DeviceStateData> deviceStates = new ConcurrentHashMap<>();
@ -250,20 +246,28 @@ public class DefaultDeviceStateService implements DeviceStateService {
}
private void initStateFromDB() {
List<Tenant> tenants = tenantService.findTenants(new TextPageLink(Integer.MAX_VALUE)).getData();
for (Tenant tenant : tenants) {
List<ListenableFuture<DeviceStateData>> fetchFutures = new ArrayList<>();
List<Device> 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<Tenant> tenants = tenantService.findTenants(new TextPageLink(Integer.MAX_VALUE)).getData();
for (Tenant tenant : tenants) {
List<ListenableFuture<DeviceStateData>> fetchFutures = new ArrayList<>();
TextPageLink pageLink = new TextPageLink(initFetchPackSize);
while (pageLink != null) {
TextPageData<Device> 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);
}
}

45
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

35
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!");
}
}
}

43
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

25
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<ConstraintViolationException> 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();
}
}
}

12
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);

9
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 (

Loading…
Cancel
Save